Files

328 lines
9.2 KiB
Go

package basemap
import (
"bytes"
"context"
"database/sql"
"encoding/json"
"errors"
"fmt"
"net/http"
"strings"
"map-asset-gateway/api-go/internal/uid"
)
func (s *Store) CreateTargetSystem(ctx context.Context, input CreateTargetSystemInput) (TargetSystem, error) {
code := normalizeBasemapCode(input.Code)
if code == "" {
return TargetSystem{}, errors.New("target system code is required")
}
name := strings.TrimSpace(input.Name)
if name == "" {
name = code
}
callbackURL := strings.TrimSpace(input.CallbackURL)
if callbackURL == "" {
return TargetSystem{}, errors.New("callback url is required")
}
method := strings.ToUpper(strings.TrimSpace(input.CallbackMethod))
if method == "" {
method = http.MethodPost
}
if input.CallbackHeaders == nil {
input.CallbackHeaders = map[string]string{}
}
now := nowUTC()
id := uid.Deterministic("target-system", code)
_, err := s.db.ExecContext(ctx, `
INSERT INTO target_systems (
id, code, name, callback_url, callback_method, callback_headers_json, enabled, created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, 1, ?, ?)
ON CONFLICT(code) DO UPDATE SET
name = excluded.name,
callback_url = excluded.callback_url,
callback_method = excluded.callback_method,
callback_headers_json = excluded.callback_headers_json,
enabled = 1,
updated_at = excluded.updated_at
`, id, code, name, callbackURL, method, writeJSON(input.CallbackHeaders), toRFC3339(now), toRFC3339(now))
if err != nil {
return TargetSystem{}, err
}
return s.GetTargetSystemByCode(ctx, code)
}
func (s *Store) GetTargetSystemByCode(ctx context.Context, code string) (TargetSystem, error) {
row := s.db.QueryRowContext(ctx, `
SELECT id, code, name, callback_url, callback_method, callback_headers_json, enabled, created_at, updated_at
FROM target_systems
WHERE code = ?
`, normalizeBasemapCode(code))
item, err := scanTargetSystem(row)
if err != nil {
if errors.Is(err, sql.ErrNoRows) {
return TargetSystem{}, fmt.Errorf("target system %q not found", code)
}
return TargetSystem{}, err
}
return item, nil
}
func (s *Store) ListTargetSystems(ctx context.Context) ([]TargetSystem, error) {
rows, err := s.db.QueryContext(ctx, `
SELECT id, code, name, callback_url, callback_method, callback_headers_json, enabled, created_at, updated_at
FROM target_systems
ORDER BY code
`)
if err != nil {
return nil, err
}
defer rows.Close()
var items []TargetSystem
for rows.Next() {
item, err := scanTargetSystem(rows)
if err != nil {
return nil, err
}
items = append(items, item)
}
return items, rows.Err()
}
func (s *Store) PushBasemapVersion(ctx context.Context, targetCode, basemapCode, version string) (PushRecord, error) {
target, err := s.GetTargetSystemByCode(ctx, targetCode)
if err != nil {
return PushRecord{}, err
}
if !target.Enabled {
return PushRecord{}, fmt.Errorf("target system %q is disabled", target.Code)
}
basemapCode = normalizeBasemapCode(basemapCode)
version = strings.TrimSpace(version)
row := s.db.QueryRowContext(ctx, `
SELECT
v.id,
v.basemap_id,
b.code,
v.version,
v.status,
v.is_default,
v.manifest_path,
v.tile_root_path,
v.url_template,
v.tile_format,
v.tile_scheme,
v.min_zoom,
v.max_zoom,
v.bbox_json,
v.attribution,
v.metadata_json,
v.created_at,
v.updated_at
FROM basemap_versions v
JOIN basemaps b ON b.id = v.basemap_id
WHERE b.code = ? AND v.version = ?
`, basemapCode, version)
versionRow, err := scanBasemapVersionRow(row)
if err != nil {
if errors.Is(err, sql.ErrNoRows) {
return PushRecord{}, fmt.Errorf("basemap version %s/%s not found", basemapCode, version)
}
return PushRecord{}, err
}
versionItem := decodeBasemapVersion(versionRow)
basemapItem, err := s.GetBasemapByCode(ctx, basemapCode)
if err != nil {
return PushRecord{}, err
}
payload := map[string]any{
"event": "basemap.published",
"service": strings.TrimSpace(s.cfg.APIBaseURL),
"basemap": basemapItem,
"version": versionItem,
"pushed_at": toRFC3339(nowUTC()),
"targetCode": target.Code,
}
body, err := json.Marshal(payload)
if err != nil {
return PushRecord{}, err
}
request, err := http.NewRequestWithContext(ctx, target.CallbackMethod, target.CallbackURL, bytes.NewReader(body))
if err != nil {
return PushRecord{}, err
}
request.Header.Set("Content-Type", "application/json")
for key, value := range target.CallbackHeaders {
request.Header.Set(key, value)
}
response, err := s.httpClient.Do(request)
record := PushRecord{
ID: uid.New(),
TargetSystemID: target.ID,
TargetSystemCode: target.Code,
BasemapVersionID: versionItem.ID,
BasemapCode: basemapItem.Code,
BasemapVersion: versionItem.Version,
Request: payload,
PushedAt: nowUTC(),
}
if err != nil {
record.Status = "failed"
record.ErrorMessage = err.Error()
} else {
defer response.Body.Close()
record.ResponseBody = limitedReadAll(response.Body, 256*1024)
status := response.StatusCode
record.ResponseStatus = &status
if status >= 200 && status < 300 {
record.Status = "succeeded"
} else {
record.Status = "failed"
record.ErrorMessage = fmt.Sprintf("target returned status %d", status)
}
}
finishedAt := nowUTC()
record.FinishedAt = &finishedAt
_, insertErr := s.db.ExecContext(ctx, `
INSERT INTO push_records (
id, target_system_id, basemap_version_id, status, request_json, response_status, response_body, error_message, pushed_at, finished_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`, record.ID, record.TargetSystemID, record.BasemapVersionID, record.Status, writeJSON(record.Request), record.ResponseStatus, record.ResponseBody, record.ErrorMessage, toRFC3339(record.PushedAt), nullableTime(record.FinishedAt))
if insertErr != nil {
return PushRecord{}, insertErr
}
if err != nil {
return record, err
}
if record.Status != "succeeded" {
return record, errors.New(record.ErrorMessage)
}
return record, nil
}
func (s *Store) ListVectorPushRecords(ctx context.Context, limit int) ([]VectorPushRecord, error) {
if limit <= 0 {
limit = 20
}
rows, err := s.db.QueryContext(ctx, `
SELECT
p.id,
p.target_system_id,
t.code,
p.vector_asset_id,
v.code,
p.status,
p.request_json,
p.response_status,
p.response_body,
p.error_message,
p.pushed_at,
p.finished_at
FROM vector_push_records p
JOIN target_systems t ON t.id = p.target_system_id
JOIN vector_assets v ON v.id = p.vector_asset_id
ORDER BY p.pushed_at DESC
LIMIT ?
`, limit)
if err != nil {
return nil, err
}
defer rows.Close()
var items []VectorPushRecord
for rows.Next() {
item, err := scanVectorPushRecord(rows)
if err != nil {
return nil, err
}
items = append(items, item)
}
return items, rows.Err()
}
func (s *Store) PushVectorAsset(ctx context.Context, targetCode, vectorCode string) (VectorPushRecord, error) {
target, err := s.GetTargetSystemByCode(ctx, targetCode)
if err != nil {
return VectorPushRecord{}, err
}
if !target.Enabled {
return VectorPushRecord{}, fmt.Errorf("target system %q is disabled", target.Code)
}
vectorCode = normalizeBasemapCode(vectorCode)
asset, err := s.GetVectorAssetByCode(ctx, vectorCode)
if err != nil {
return VectorPushRecord{}, err
}
payload := map[string]any{
"event": "vector.published",
"service": strings.TrimSpace(s.cfg.APIBaseURL),
"vector": asset,
"pushed_at": toRFC3339(nowUTC()),
"targetCode": target.Code,
}
body, err := json.Marshal(payload)
if err != nil {
return VectorPushRecord{}, err
}
request, err := http.NewRequestWithContext(ctx, target.CallbackMethod, target.CallbackURL, bytes.NewReader(body))
if err != nil {
return VectorPushRecord{}, err
}
request.Header.Set("Content-Type", "application/json")
for key, value := range target.CallbackHeaders {
request.Header.Set(key, value)
}
response, err := s.httpClient.Do(request)
record := VectorPushRecord{
ID: uid.New(),
TargetSystemID: target.ID,
TargetSystemCode: target.Code,
VectorAssetID: asset.ID,
VectorCode: asset.Code,
Request: payload,
PushedAt: nowUTC(),
}
if err != nil {
record.Status = "failed"
record.ErrorMessage = err.Error()
} else {
defer response.Body.Close()
record.ResponseBody = limitedReadAll(response.Body, 256*1024)
status := response.StatusCode
record.ResponseStatus = &status
if status >= 200 && status < 300 {
record.Status = "succeeded"
} else {
record.Status = "failed"
record.ErrorMessage = fmt.Sprintf("target returned status %d", status)
}
}
finishedAt := nowUTC()
record.FinishedAt = &finishedAt
_, insertErr := s.db.ExecContext(ctx, `
INSERT INTO vector_push_records (
id, target_system_id, vector_asset_id, status, request_json, response_status, response_body, error_message, pushed_at, finished_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`, record.ID, record.TargetSystemID, record.VectorAssetID, record.Status, writeJSON(record.Request), record.ResponseStatus, record.ResponseBody, record.ErrorMessage, toRFC3339(record.PushedAt), nullableTime(record.FinishedAt))
if insertErr != nil {
return VectorPushRecord{}, insertErr
}
if err != nil {
return record, err
}
if record.Status != "succeeded" {
return record, errors.New(record.ErrorMessage)
}
return record, nil
}