579 lines
16 KiB
Go
579 lines
16 KiB
Go
package basemap
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"path/filepath"
|
|
"sort"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"map-asset-gateway/api-go/internal/uid"
|
|
)
|
|
|
|
type scannedVersion struct {
|
|
BasemapCode string
|
|
Name string
|
|
Type string
|
|
Status string
|
|
Description string
|
|
Version string
|
|
Manifest string
|
|
TileRoot string
|
|
TileFormat string
|
|
TileScheme string
|
|
MinZoom int
|
|
MaxZoom int
|
|
BBox []float64
|
|
Attribution string
|
|
Metadata map[string]any
|
|
IsDefault bool
|
|
}
|
|
|
|
func (s *Store) RunScan(ctx context.Context, sourceCode string) (ScanRun, error) {
|
|
source, err := s.GetScanSourceByCode(ctx, sourceCode)
|
|
if err != nil {
|
|
return ScanRun{}, err
|
|
}
|
|
if !source.Enabled {
|
|
return ScanRun{}, fmt.Errorf("scan source %q is disabled", sourceCode)
|
|
}
|
|
|
|
startedAt := nowUTC()
|
|
runID := uid.New()
|
|
if _, err := s.db.ExecContext(ctx, `
|
|
INSERT INTO scan_runs (
|
|
id, scan_source_id, status, scanned_count, added_count, updated_count, removed_count, summary_json, started_at, finished_at
|
|
) VALUES (?, ?, 'running', 0, 0, 0, 0, '', ?, ?)
|
|
`, runID, source.ID, toRFC3339(startedAt), toRFC3339(startedAt)); err != nil {
|
|
return ScanRun{}, fmt.Errorf("create scan run: %w", err)
|
|
}
|
|
|
|
run, finalErr := s.runScan(ctx, runID, source, startedAt)
|
|
return run, finalErr
|
|
}
|
|
|
|
func (s *Store) runScan(ctx context.Context, runID string, source ScanSource, startedAt time.Time) (ScanRun, error) {
|
|
rootInfo, err := os.Stat(source.RootPath)
|
|
if err != nil {
|
|
return s.finishScanRun(ctx, runID, source, startedAt, "failed", nil, 0, 0, 0, 0, fmt.Errorf("stat scan root: %w", err))
|
|
}
|
|
if !rootInfo.IsDir() {
|
|
return s.finishScanRun(ctx, runID, source, startedAt, "failed", nil, 0, 0, 0, 0, fmt.Errorf("scan root %s is not a directory", source.RootPath))
|
|
}
|
|
|
|
manifestItems, coveredRoots, err := s.findManifestVersions(source)
|
|
if err != nil {
|
|
return s.finishScanRun(ctx, runID, source, startedAt, "failed", nil, 0, 0, 0, 0, err)
|
|
}
|
|
fallbackItems, err := s.findDirectTileRoots(source, coveredRoots)
|
|
if err != nil {
|
|
return s.finishScanRun(ctx, runID, source, startedAt, "failed", nil, 0, 0, 0, 0, err)
|
|
}
|
|
|
|
allItems := append(manifestItems, fallbackItems...)
|
|
sort.Slice(allItems, func(i, j int) bool {
|
|
if allItems[i].BasemapCode == allItems[j].BasemapCode {
|
|
return allItems[i].Version < allItems[j].Version
|
|
}
|
|
return allItems[i].BasemapCode < allItems[j].BasemapCode
|
|
})
|
|
|
|
addedCount, updatedCount, removedCount, err := s.persistScanResult(ctx, source, allItems)
|
|
if err != nil {
|
|
return s.finishScanRun(ctx, runID, source, startedAt, "failed", allItems, len(allItems), addedCount, updatedCount, removedCount, err)
|
|
}
|
|
return s.finishScanRun(ctx, runID, source, startedAt, "succeeded", allItems, len(allItems), addedCount, updatedCount, removedCount, nil)
|
|
}
|
|
|
|
func (s *Store) finishScanRun(
|
|
ctx context.Context,
|
|
runID string,
|
|
source ScanSource,
|
|
startedAt time.Time,
|
|
status string,
|
|
items []scannedVersion,
|
|
scannedCount int,
|
|
addedCount int,
|
|
updatedCount int,
|
|
removedCount int,
|
|
runErr error,
|
|
) (ScanRun, error) {
|
|
finishedAt := nowUTC()
|
|
summary := map[string]any{
|
|
"source_code": source.Code,
|
|
"scan_root": source.RootPath,
|
|
"scanned_count": scannedCount,
|
|
"added_count": addedCount,
|
|
"updated_count": updatedCount,
|
|
"removed_count": removedCount,
|
|
}
|
|
if runErr != nil {
|
|
summary["error"] = runErr.Error()
|
|
}
|
|
if len(items) > 0 {
|
|
codes := map[string]struct{}{}
|
|
for _, item := range items {
|
|
codes[item.BasemapCode] = struct{}{}
|
|
}
|
|
summary["basemap_codes"] = sortedKeys(codes)
|
|
}
|
|
|
|
_, updateErr := s.db.ExecContext(ctx, `
|
|
UPDATE scan_runs
|
|
SET status = ?, scanned_count = ?, added_count = ?, updated_count = ?, removed_count = ?, summary_json = ?, finished_at = ?
|
|
WHERE id = ?
|
|
`, status, scannedCount, addedCount, updatedCount, removedCount, writeJSON(summary), toRFC3339(finishedAt), runID)
|
|
if updateErr != nil {
|
|
if runErr != nil {
|
|
return ScanRun{}, fmt.Errorf("%v; update scan run: %w", runErr, updateErr)
|
|
}
|
|
return ScanRun{}, updateErr
|
|
}
|
|
|
|
run := ScanRun{
|
|
ID: runID,
|
|
ScanSourceID: source.ID,
|
|
SourceCode: source.Code,
|
|
Status: status,
|
|
ScannedCount: scannedCount,
|
|
AddedCount: addedCount,
|
|
UpdatedCount: updatedCount,
|
|
RemovedCount: removedCount,
|
|
Summary: summary,
|
|
StartedAt: startedAt,
|
|
FinishedAt: finishedAt,
|
|
}
|
|
if runErr != nil {
|
|
return run, runErr
|
|
}
|
|
return run, nil
|
|
}
|
|
|
|
func (s *Store) findManifestVersions(source ScanSource) ([]scannedVersion, map[string]struct{}, error) {
|
|
manifestName := strings.TrimSpace(source.ManifestName)
|
|
if manifestName == "" {
|
|
manifestName = defaultManifestName
|
|
}
|
|
|
|
var items []scannedVersion
|
|
coveredRoots := map[string]struct{}{}
|
|
err := filepath.WalkDir(source.RootPath, func(path string, entry os.DirEntry, walkErr error) error {
|
|
if walkErr != nil {
|
|
return walkErr
|
|
}
|
|
if entry.IsDir() {
|
|
return nil
|
|
}
|
|
if !strings.EqualFold(entry.Name(), manifestName) {
|
|
return nil
|
|
}
|
|
|
|
manifest, err := LoadManifest(path)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
tileRoot, err := resolveTileRoot(filepath.Dir(path), manifest.RootPath)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
format, minZoom, maxZoom, found, err := inspectTileRoot(tileRoot, manifest.TileFormat)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if found {
|
|
manifest.TileFormat = format
|
|
manifest.MinZoom = minZoom
|
|
manifest.MaxZoom = maxZoom
|
|
}
|
|
|
|
relPath, err := filepath.Rel(source.RootPath, tileRoot)
|
|
if err == nil {
|
|
first := strings.Split(filepath.ToSlash(relPath), "/")[0]
|
|
if first != "" && first != "." {
|
|
coveredRoots[first] = struct{}{}
|
|
}
|
|
}
|
|
|
|
items = append(items, scannedVersion{
|
|
BasemapCode: manifest.Code,
|
|
Name: defaultText(manifest.Name, manifest.Code),
|
|
Type: defaultText(manifest.Type, "xyz"),
|
|
Status: defaultText(manifest.Status, "ready"),
|
|
Description: manifest.Description,
|
|
Version: defaultText(manifest.Version, "current"),
|
|
Manifest: path,
|
|
TileRoot: tileRoot,
|
|
TileFormat: normalizeTileFormat(manifest.TileFormat),
|
|
TileScheme: defaultText(manifest.TileScheme, "xyz"),
|
|
MinZoom: manifest.MinZoom,
|
|
MaxZoom: manifest.MaxZoom,
|
|
BBox: manifest.BBox,
|
|
Attribution: manifest.Attribution,
|
|
Metadata: manifest.Metadata,
|
|
IsDefault: manifest.IsDefault,
|
|
})
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
return nil, nil, fmt.Errorf("walk manifests: %w", err)
|
|
}
|
|
return items, coveredRoots, nil
|
|
}
|
|
|
|
func (s *Store) findDirectTileRoots(source ScanSource, skip map[string]struct{}) ([]scannedVersion, error) {
|
|
entries, err := os.ReadDir(source.RootPath)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("read scan root: %w", err)
|
|
}
|
|
|
|
var items []scannedVersion
|
|
for _, entry := range entries {
|
|
if !entry.IsDir() {
|
|
continue
|
|
}
|
|
name := entry.Name()
|
|
if _, exists := skip[name]; exists {
|
|
continue
|
|
}
|
|
|
|
rootPath := filepath.Join(source.RootPath, name)
|
|
format, minZoom, maxZoom, found, err := inspectTileRoot(rootPath, "")
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if !found {
|
|
continue
|
|
}
|
|
|
|
code := normalizeBasemapCode(name)
|
|
if code == "" {
|
|
continue
|
|
}
|
|
items = append(items, scannedVersion{
|
|
BasemapCode: code,
|
|
Name: name,
|
|
Type: "xyz",
|
|
Status: "ready",
|
|
Version: "current",
|
|
TileRoot: rootPath,
|
|
TileFormat: format,
|
|
TileScheme: "xyz",
|
|
MinZoom: minZoom,
|
|
MaxZoom: maxZoom,
|
|
Metadata: map[string]any{
|
|
"discovered_from": "directory_scan",
|
|
},
|
|
IsDefault: true,
|
|
})
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
func inspectTileRoot(rootPath string, hintFormat string) (string, int, int, bool, error) {
|
|
entries, err := os.ReadDir(rootPath)
|
|
if err != nil {
|
|
if errors.Is(err, os.ErrNotExist) {
|
|
return "", 0, 0, false, nil
|
|
}
|
|
return "", 0, 0, false, fmt.Errorf("read tile root %s: %w", rootPath, err)
|
|
}
|
|
|
|
var zooms []int
|
|
format := normalizeTileFormat(hintFormat)
|
|
for _, entry := range entries {
|
|
if !entry.IsDir() {
|
|
continue
|
|
}
|
|
zoom, err := strconv.Atoi(entry.Name())
|
|
if err != nil {
|
|
continue
|
|
}
|
|
zooms = append(zooms, zoom)
|
|
if format == "" || format == "png" {
|
|
detected, found, err := detectTileFormat(filepath.Join(rootPath, entry.Name()))
|
|
if err != nil {
|
|
return "", 0, 0, false, err
|
|
}
|
|
if found {
|
|
format = detected
|
|
}
|
|
}
|
|
}
|
|
if len(zooms) == 0 {
|
|
return "", 0, 0, false, nil
|
|
}
|
|
sort.Ints(zooms)
|
|
if format == "" {
|
|
format = "png"
|
|
}
|
|
return format, zooms[0], zooms[len(zooms)-1], true, nil
|
|
}
|
|
|
|
func detectTileFormat(zoomPath string) (string, bool, error) {
|
|
xEntries, err := os.ReadDir(zoomPath)
|
|
if err != nil {
|
|
return "", false, fmt.Errorf("read zoom path %s: %w", zoomPath, err)
|
|
}
|
|
for _, xEntry := range xEntries {
|
|
if !xEntry.IsDir() {
|
|
continue
|
|
}
|
|
tileEntries, err := os.ReadDir(filepath.Join(zoomPath, xEntry.Name()))
|
|
if err != nil {
|
|
return "", false, err
|
|
}
|
|
for _, tileEntry := range tileEntries {
|
|
if tileEntry.IsDir() {
|
|
continue
|
|
}
|
|
ext := strings.TrimPrefix(strings.ToLower(filepath.Ext(tileEntry.Name())), ".")
|
|
if ext == "" {
|
|
continue
|
|
}
|
|
return normalizeTileFormat(ext), true, nil
|
|
}
|
|
}
|
|
return "", false, nil
|
|
}
|
|
|
|
func resolveTileRoot(baseDir, relativePath string) (string, error) {
|
|
if strings.TrimSpace(relativePath) == "" {
|
|
relativePath = "tiles"
|
|
}
|
|
if filepath.IsAbs(relativePath) {
|
|
return filepath.Clean(relativePath), nil
|
|
}
|
|
resolved := filepath.Clean(filepath.Join(baseDir, relativePath))
|
|
return resolved, nil
|
|
}
|
|
|
|
func (s *Store) persistScanResult(ctx context.Context, source ScanSource, items []scannedVersion) (int, int, int, error) {
|
|
tx, err := s.db.BeginTx(ctx, nil)
|
|
if err != nil {
|
|
return 0, 0, 0, err
|
|
}
|
|
defer tx.Rollback()
|
|
|
|
existing := map[string]struct{}{}
|
|
rows, err := tx.QueryContext(ctx, `
|
|
SELECT b.code, v.version
|
|
FROM basemap_versions v
|
|
JOIN basemaps b ON b.id = v.basemap_id
|
|
WHERE v.scan_source_id = ?
|
|
`, source.ID)
|
|
if err != nil {
|
|
return 0, 0, 0, err
|
|
}
|
|
for rows.Next() {
|
|
var code string
|
|
var version string
|
|
if err := rows.Scan(&code, &version); err != nil {
|
|
rows.Close()
|
|
return 0, 0, 0, err
|
|
}
|
|
existing[code+"::"+version] = struct{}{}
|
|
}
|
|
rows.Close()
|
|
if err := rows.Err(); err != nil {
|
|
return 0, 0, 0, err
|
|
}
|
|
|
|
now := nowUTC()
|
|
found := map[string]struct{}{}
|
|
touchedBasemaps := map[string]struct{}{}
|
|
addedCount := 0
|
|
updatedCount := 0
|
|
for _, item := range items {
|
|
key := item.BasemapCode + "::" + item.Version
|
|
found[key] = struct{}{}
|
|
if _, exists := existing[key]; exists {
|
|
updatedCount++
|
|
} else {
|
|
addedCount++
|
|
}
|
|
if err := s.upsertScannedVersionTx(ctx, tx, source, item, now); err != nil {
|
|
return 0, 0, 0, err
|
|
}
|
|
touchedBasemaps[item.BasemapCode] = struct{}{}
|
|
}
|
|
|
|
removedCount := 0
|
|
for key := range existing {
|
|
if _, ok := found[key]; ok {
|
|
continue
|
|
}
|
|
parts := strings.SplitN(key, "::", 2)
|
|
if len(parts) != 2 {
|
|
continue
|
|
}
|
|
if err := s.deleteVersionBySourceTx(ctx, tx, source.ID, parts[0], parts[1]); err != nil {
|
|
return 0, 0, 0, err
|
|
}
|
|
touchedBasemaps[parts[0]] = struct{}{}
|
|
removedCount++
|
|
}
|
|
|
|
for code := range touchedBasemaps {
|
|
if err := s.normalizeBasemapDefaultsTx(ctx, tx, code, now); err != nil {
|
|
return 0, 0, 0, err
|
|
}
|
|
}
|
|
|
|
if err := tx.Commit(); err != nil {
|
|
return 0, 0, 0, err
|
|
}
|
|
return addedCount, updatedCount, removedCount, nil
|
|
}
|
|
|
|
func (s *Store) upsertScannedVersionTx(ctx context.Context, tx *sql.Tx, source ScanSource, item scannedVersion, now time.Time) error {
|
|
basemapID := uid.Deterministic("basemap", item.BasemapCode)
|
|
versionID := uid.Deterministic("basemap-version", item.BasemapCode, item.Version)
|
|
|
|
_, err := tx.ExecContext(ctx, `
|
|
INSERT INTO basemaps (id, code, name, type, status, description, created_at, updated_at)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
|
|
ON CONFLICT(code) DO UPDATE SET
|
|
name = excluded.name,
|
|
type = excluded.type,
|
|
status = excluded.status,
|
|
description = excluded.description,
|
|
updated_at = excluded.updated_at
|
|
`, basemapID, item.BasemapCode, item.Name, item.Type, item.Status, item.Description, toRFC3339(now), toRFC3339(now))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
isDefault := 0
|
|
if item.IsDefault {
|
|
isDefault = 1
|
|
}
|
|
_, err = tx.ExecContext(ctx, `
|
|
INSERT INTO basemap_versions (
|
|
id, basemap_id, scan_source_id, version, status, is_default, manifest_path, tile_root_path, url_template,
|
|
tile_format, tile_scheme, min_zoom, max_zoom, bbox_json, attribution, metadata_json, created_at, updated_at
|
|
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
|
ON CONFLICT(basemap_id, version) DO UPDATE SET
|
|
scan_source_id = excluded.scan_source_id,
|
|
status = excluded.status,
|
|
is_default = CASE WHEN excluded.is_default = 1 THEN 1 ELSE basemap_versions.is_default END,
|
|
manifest_path = excluded.manifest_path,
|
|
tile_root_path = excluded.tile_root_path,
|
|
url_template = excluded.url_template,
|
|
tile_format = excluded.tile_format,
|
|
tile_scheme = excluded.tile_scheme,
|
|
min_zoom = excluded.min_zoom,
|
|
max_zoom = excluded.max_zoom,
|
|
bbox_json = excluded.bbox_json,
|
|
attribution = excluded.attribution,
|
|
metadata_json = excluded.metadata_json,
|
|
updated_at = excluded.updated_at
|
|
`, versionID, basemapID, source.ID, item.Version, item.Status, isDefault, item.Manifest, item.TileRoot, s.buildURLTemplate(item.BasemapCode, item.Version, item.TileFormat), item.TileFormat, item.TileScheme, item.MinZoom, item.MaxZoom, writeJSON(item.BBox), item.Attribution, writeJSON(item.Metadata), toRFC3339(now), toRFC3339(now))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if item.IsDefault {
|
|
_, err = tx.ExecContext(ctx, `
|
|
UPDATE basemap_versions
|
|
SET is_default = CASE WHEN version = ? THEN 1 ELSE 0 END, updated_at = ?
|
|
WHERE basemap_id = ?
|
|
`, item.Version, toRFC3339(now), basemapID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *Store) deleteVersionBySourceTx(ctx context.Context, tx *sql.Tx, sourceID, basemapCode, version string) error {
|
|
result, err := tx.ExecContext(ctx, `
|
|
DELETE FROM basemap_versions
|
|
WHERE scan_source_id = ? AND version = ? AND basemap_id = (SELECT id FROM basemaps WHERE code = ?)
|
|
`, sourceID, version, basemapCode)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if _, err := result.RowsAffected(); err != nil {
|
|
return err
|
|
}
|
|
_, err = tx.ExecContext(ctx, `
|
|
DELETE FROM basemaps
|
|
WHERE code = ? AND NOT EXISTS (
|
|
SELECT 1 FROM basemap_versions WHERE basemap_id = basemaps.id
|
|
)
|
|
`, basemapCode)
|
|
return err
|
|
}
|
|
|
|
func (s *Store) normalizeBasemapDefaultsTx(ctx context.Context, tx *sql.Tx, basemapCode string, now time.Time) error {
|
|
var basemapID string
|
|
err := tx.QueryRowContext(ctx, `SELECT id FROM basemaps WHERE code = ?`, basemapCode).Scan(&basemapID)
|
|
if err != nil {
|
|
if errors.Is(err, sql.ErrNoRows) {
|
|
return nil
|
|
}
|
|
return err
|
|
}
|
|
|
|
rows, err := tx.QueryContext(ctx, `
|
|
SELECT version, is_default
|
|
FROM basemap_versions
|
|
WHERE basemap_id = ?
|
|
ORDER BY is_default DESC, version ASC
|
|
`, basemapID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer rows.Close()
|
|
|
|
type versionFlag struct {
|
|
Version string
|
|
IsDefault bool
|
|
}
|
|
var versions []versionFlag
|
|
for rows.Next() {
|
|
var item versionFlag
|
|
var isDefault int
|
|
if err := rows.Scan(&item.Version, &isDefault); err != nil {
|
|
return err
|
|
}
|
|
item.IsDefault = isDefault == 1
|
|
versions = append(versions, item)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return err
|
|
}
|
|
if len(versions) == 0 {
|
|
return nil
|
|
}
|
|
defaultVersion := versions[0].Version
|
|
for _, item := range versions {
|
|
if item.IsDefault {
|
|
defaultVersion = item.Version
|
|
break
|
|
}
|
|
}
|
|
if _, err := tx.ExecContext(ctx, `
|
|
UPDATE basemap_versions
|
|
SET is_default = CASE WHEN version = ? THEN 1 ELSE 0 END, updated_at = ?
|
|
WHERE basemap_id = ?
|
|
`, defaultVersion, toRFC3339(now), basemapID); err != nil {
|
|
return err
|
|
}
|
|
_, err = tx.ExecContext(ctx, `UPDATE basemaps SET updated_at = ? WHERE id = ?`, toRFC3339(now), basemapID)
|
|
return err
|
|
}
|
|
|
|
func (s *Store) buildURLTemplate(code, version, format string) string {
|
|
base := trimURL(s.cfg.TileBaseURL)
|
|
if base == "" {
|
|
base = trimURL(s.cfg.APIBaseURL)
|
|
}
|
|
format = normalizeTileFormat(format)
|
|
return fmt.Sprintf("%s/tiles/%s/%s/{z}/{x}/{y}.%s", base, code, version, format)
|
|
}
|