Initial import of map-asset-gateway
This commit is contained in:
@@ -0,0 +1,327 @@
|
||||
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
|
||||
}
|
||||
Reference in New Issue
Block a user