diff --git a/src/pgstac/.gitignore b/src/pgstac/.gitignore new file mode 100644 index 0000000..d77f3bf --- /dev/null +++ b/src/pgstac/.gitignore @@ -0,0 +1,13 @@ +# Secrets — never commit. Only example.env is tracked. +.env +*/.env +# Re-include the placeholder template (repo-root .gitignore has a broad *.env rule). +!example.env + +# Python +__pycache__/ +*.py[cod] +.venv/ + +# OS +.DS_Store diff --git a/src/pgstac/README.md b/src/pgstac/README.md new file mode 100644 index 0000000..78fe0fa --- /dev/null +++ b/src/pgstac/README.md @@ -0,0 +1,270 @@ +# pgstac + +Dette prosjektet er en eksperimentering for å bli bedre kjent med **pgstac** og hvordan +den brukes. Målet er å katalogisere geospatiale filer som ligger i et lokalt kjørende +S3-lager, og gjøre dem søkbare og browsbare som en STAC-katalog. I stedet for å skrive +katalogen som statiske JSON-filer lagrer vi alt i Postgres med pgstac, og serverer det med +en ekte STAC API. Så scanner vi S3-lageret for filer, leser ut metadata fra hver fil, og +fyller katalogen. + +Tanken er: en mappe i S3 = en STAC Collection. Du peker scanneren mot en bucket og et +prefix, og den bygger opp en collection med ett item per fil. Kjører du den på nytt holder +den katalogen i synk med mappa (legger til nye filer, fjerner de som er borte). + +For å kjøre dette trenger du en Postgres-database med PostGIS og pgstac (`db`), en STAC API +oppå den (`api`), og en viewer (`viewer`) for å browse i nettleser. Scanneren (`scanner`) +kjøres bare når du skal fylle eller oppdatere katalogen. Se [Hva vi har satt opp](#hva-vi-har-satt-opp) +for detaljer. + +### Hvordan dette skiller seg fra en statisk JSON-katalog på S3 + +En enklere variant er å skrive STAC-katalogen som statiske JSON-filer rett på S3. Det +krever ingen database og er billig å hoste, men du får ingen søk eller API ut av det, bare +filer en klient må laste ned og bla gjennom selv. + +pgstac-oppsettet ser ut til å gi flere fordeler. Du får søk og et ekte STAC API rett ut av +boksen, så klienter (STAC Browser, QGIS, egne script) kan filtrere på tid, område og +properties uten å laste ned hele katalogen. Til gjengjeld krever det at du har en +PostgreSQL-database kjørende som brukerne/klientene har tilgang til. Men når den først er +oppe, gjør den jobben kjempebra. + +## Hva vi har satt opp + +Fire containere, der de tre første kjører hele tiden og scanneren kjøres ved behov: + +| Tjeneste | Hva | Levetid | +|----------|-----|---------| +| `db` | PostgreSQL 17 + PostGIS 3 + pgstac (`ghcr.io/stac-utils/pgstac:v0.9.11`) | kjører fast | +| `api` | STAC API (`ghcr.io/stac-utils/stac-fastapi-pgstac:6.2.1`) | kjører fast | +| `viewer` | STAC Browser + en `/s3`-proxy som signerer asset-kall | kjører fast | +| `scanner` | scanner en S3-mappe og laster en Collection inn i pgstac | **ved behov** | + +- **db** er der hele katalogen bor. pgstac legger til sitt eget skjema (collections, items, + søk) oppå PostGIS. Databasen restartes aldri når vi oppdaterer katalogen. +- **api** er stac-fastapi-pgstac. STAC Browser kan ikke snakke med Postgres direkte, så + denne serverer katalogen som JSON over HTTP (med CORS på, så nettleseren kan kalle den). +- **viewer** serverer STAC Browser (SPA-en) og har i tillegg en `/s3`-proxy. Bucketen er + privat, så nettleseren kan ikke hente filene direkte. Proxyen signerer kallene med SigV4 + og støtter Range-requests, slik at f.eks. PMTiles kan leses bit for bit. +- **scanner** er den vi bygger selv. Den lister opp filene i en S3-mappe, leser metadata, + bygger STAC-items og laster dem rett inn i pgstac med `pypgstac`. + +## Slik henger det sammen + +Overordnet er det to veier: en scan-vei som fyller databasen on-demand, og en serve-vei der +katalogen leses ut. Scanneren bygger STAC-items og gjør upsert + prune mot pgstac. API-et +serverer collections og items, og vieweren henter fil-bytene gjennom `/s3`-proxyen med +signerte range-reads. + +Flytdiagram: scan-vei og serve-vei + +Mer detaljert med containere, porter og `.env`: db, api og viewer er den faste stacken på +det delte `pgstac`-nettverket, scanneren kjøres on-demand, og alle deler samme `.env`. +Katalog-JSON går rett fra nettleser til API (CORS), mens fil-bytene går via viewerens +`/s3`-proxy som signerer mot privat S3. + +Arkitekturdiagram med containere, porter, nettverk og .env + +Tjenestene i diagrammene er forklart i [Hva vi har satt opp](#hva-vi-har-satt-opp). + +## Hvordan scanningen fungerer + +Scanneren tar en bucket + prefix og gjør dette per kjøring: + +1. Lister opp filene under prefixet (boto3). +2. For hver fil leser den ut metadata. Den laster ikke ned hele fila der det går an, den + leser bare header/footer (PMTiles-header, parquet-footer, COG-header) for å holde det + billig. +3. Bygger ett STAC-item per fil, med asset-href som peker på viewerens `/s3`-proxy + (`/s3//`). +4. Upserter collection + alle items i pgstac med `pypgstac`. +5. Pruner: sammenligner item-id-ene i scanet mot de som ligger i databasen, og sletter de + som ikke lenger finnes i mappa. (`--no-prune` skrur det av, `--drop` tømmer collectionen + først.) + +Resultatet er idempotent. Kjører du samme `--collection-id` på nytt speiler collectionen +det som faktisk ligger i S3-mappa akkurat nå. Støttede filtyper: COG/GeoTIFF, GeoJSON, +GeoParquet og PMTiles. + +Slik ser katalogen ut i STAC Browser, med en collection per scannet mappe: + +STAC Browser katalogoversikt med to collections + +Inni en collection ligger items med bbox-kart og liste: + +Collection-visning med bbox-kart og item-liste + +### Slik ser det ut i databasen + +Alt dette ligger i Postgres. pgstac lager sitt eget skjema med blant annet `collections`- +og `items`-tabellene. Hver scannet mappe blir en rad i `collections`, og hver fil blir en +rad i `items` (med geometri, collection-referanse og datetime). Du trenger ikke å gå inn i +databasen for å bruke systemet, men det er greit å se hvor dataene faktisk havner: + +pgstac-skjemaet i databasen med collections-tabellen + +items-tabellen med geometri, collection og datetime per fil + +## Ekstra metadata per filtype + +Utover det vanlige (bbox, projeksjon, datetime) leser vi ut felter som er nyttige å se i +browseren. Alle items får `file:size` (filstørrelse i bytes). Resten avhenger av filtype. +Der det finnes en standard STAC-extension bruker vi den, ellers et eget namespace. + +**GeoParquet / GeoJSON** +- `table:columns` (kolonnenavn + datatype) og `table:row_count` (antall rader) +- `vector:geometry_types`, f.eks. `LineString Z`, `Polygon` +- `vector:encoding`, f.eks. `WKB` +- `proj:epsg` utledet fra geo-metadataen +- noen `proc:sample`-rader som forhåndsvisning av dataene + +For parquet leses dette fra footeren, så vi slipper å lese hele fila. + +GeoParquet-item med kolonner, radantall, geometritype og sample-rader + +**PMTiles** +- `pmtiles:tile_type` (vector/raster), min/max zoom, center +- `pmtiles:name` og `pmtiles:vector_layers` (lagene + feltskjemaet i hvert lag) +- `pmtiles:clustered`, `pmtiles:tile_compression` +- teller for tiles (`pmtiles:addressed_tiles_count` osv.). Merk at dette er antall tiles, + ikke antall features (features dupliseres på tvers av zoomnivåer). + +Alt dette leses gratis fra PMTiles-headeren. + +PMTiles-item med navn, zoom, tile type og vector layers + +**COG / GeoTIFF** +- bånd, dtype, nodata og statistikk (via rio-stac) +- `cog:compression`, `cog:blocksize`, `cog:overview_count`, `cog:predictor` fra + COG-headeren + +## MapLibre-preview for PMTiles + +For PMTiles-items får du en egen "Data preview"-fane i vieweren som rendrer tile-ene +direkte med MapLibre. Den henter PMTiles via `/s3`-proxyen (med Range, så den laster bare +de tile-ene den trenger) og tegner vektorlagene oppå et basiskart. Da ser du faktisk +innholdet i fila uten å laste den ned. + +MapLibre-preview som rendrer PMTiles-vektorlag + +Basiskart-stilen settes med `BASEMAP_STYLE_URL` (default er MapLibre sin demo-stil, kan +byttes til f.eks. Norkart-stilen). + +## Bruke katalogen i QGIS + +Siden `api` er en helt vanlig STAC API kan du også koble deg på fra QGIS med STAC +API Browser-pluginen. Legg inn API-URL-en som en STAC-tilkobling, så dukker collections og +items opp i Browser-panelet og kan lastes rett inn i kartet: + +QGIS Browser med STAC-tilkobling som viser collections og items + +Hvert item viser metadataene vi har lagt på, inkludert projeksjon, bbox og extensions: + +STAC Object Details-dialogen i QGIS for et item + +## Kom i gang lokalt + +Du trenger Docker. Så er det tre steg: lag `.env`, lag det delte nettverket, og start +stacken. + +```sh +cp example.env .env # fyll inn verdiene (se under) +docker network create pgstac # engangs, delt nettverk som alle tjenestene henger på +docker compose --env-file .env up -d # starter db + api + viewer (ikke scanner) +``` + +Sjekk at det lever: + +```sh +docker compose ps +curl -s localhost:8082/collections | jq . # STAC API (tom til en scan har kjørt) +curl -s localhost:8080/healthz # viewer +open http://localhost:8080 # STAC Browser +``` + +Nå kjører db, api og viewer. Katalogen er tom til du har kjørt en scan (se neste seksjon). + +### Hva du må fylle ut i `.env` + +`example.env` har alle variablene med kommentarer. De som står til `changeme` må du sette +selv, resten har fornuftige defaults for lokal kjøring og kan stå som de er. + +| Variabel | Må settes? | Hva det er | +|----------|------------|------------| +| `POSTGRES_USER` | nei | DB-bruker (default `pgstac`) | +| `POSTGRES_PASSWORD` | **ja** | DB-passord, bytt fra `changeme` | +| `POSTGRES_DB` | nei | DB-navn (default `postgis`) | +| `DB_HOST_PORT` | nei | host-port for Postgres (default `5439`) | +| `API_HOST_PORT` | nei | host-port for STAC API (default `8082`) | +| `CORS_ORIGINS` | nei | hvilke origins som får kalle API-et fra nettleser. Må inkludere viewer-origin (default `http://localhost:8080`) | +| `S3_ENDPOINT` | **ja** | S3-host uten scheme, f.eks. `s3.example.no` (https legges på internt) | +| `S3_ACCESS_KEY` | **ja** | S3-nøkkel | +| `S3_SECRET_KEY` | **ja** | S3-secret | +| `S3_REGION` | nei | signeringsregion (default `us-east-1`) | +| `ASSET_BASE_URL` | nei | base for asset-href som lagres i DB, `/s3//`. Må matche hvordan nettleseren når vieweren (default `http://localhost:8080`) | +| `STAC_API_URL` | nei | hvordan nettleseren når API-et (default `http://localhost:8082`). Må stå i `CORS_ORIGINS` | +| `VIEWER_HOST_PORT` | nei | host-port for vieweren (default `8080`) | +| `STAC_BROWSER_VERSION` | nei | STAC Browser-versjon som bakes inn (default `v4.0.1`) | +| `BASEMAP_STYLE_URL` | nei | MapLibre basiskart-stil for preview-fanen (default MapLibre sin demo-stil) | + +S3-feltene brukes både av scanneren (lese data) og vieweren (signere asset-kall). +Scanneren utleder DB-koblingen fra `POSTGRES_*`, så den holder seg i synk automatisk. Sett +`PGSTAC_DSN` bare hvis du vil peke scanneren på en helt annen database. + +## Sette opp en scan med scanner-containeren + +Når stacken kjører fyller du katalogen ved å kjøre scanner-containeren mot en S3-mappe. +Den kjøres on-demand (`run --rm`) mot den allerede oppe-kjørende stacken, så databasen +restartes aldri. S3-credentials og DB-kobling kommer fra `.env`, resten er CLI-flagg. + +Den raskeste måten er å peke rett på en bucket + prefix: + +```sh +docker compose -f scanner/docker-compose.yml --env-file .env run --rm scanner \ + --bucket my-bucket --prefix area/sub/ \ + --collection-id area-sub --collection-title "Area Sub" +``` + +- `--bucket` er bucket-navnet (uten skråstrek) +- `--prefix` er mappa inni bucketen som blir til én collection +- `--collection-id` er en stabil id. Kjør samme id igjen for å re-scanne (upsert + prune) +- `--collection-title` er tittelen som vises i browseren + +Nyttige flagg: `--no-recursive` (bare toppnivået), `--no-prune` (behold items som er borte +fra S3), `--drop` (tøm collectionen først), `--datetime` (tving én dato på alle items), +`--fail-on-error`. `--help` viser alle. + +### Faste scan-jobber i en fil + +Vil du lagre oppsettet for en scan i stedet for å skrive flagg hver gang, bruk en +per-collection-fil (én fil = én Collection). Kopier `scanner/collection-nseries.yml` til +`scanner/collection-.yml`, gi servicen et eget navn, og fyll inn `--bucket`, +`--prefix`, `--collection-id` og `--collection-title` i `command`-lista. Secrets og DB +holdes fortsatt i `.env`, aldri i denne fila. Så kjører du uten CLI-args: + +```sh +docker compose -f scanner/collection-.yml --env-file .env run --rm +``` + +Vil du kjøre flere på en gang, legg hver fil inn i +`scanner/docker-compose-all-collections.yml` (én `include:`-linje per fil) og: + +```sh +docker compose -f scanner/docker-compose-all-collections.yml --env-file .env up --abort-on-container-exit +``` + +Kjører du samme `--collection-id` på nytt **upserter** den fant-items og **pruner** items +som ikke lenger ligger i mappa. Refresh så browseren, og den nye collectionen dukker opp +under API-roten. + +## Porter + +- Viewer / STAC Browser: `localhost:${VIEWER_HOST_PORT:-8080}` +- STAC API: `localhost:${API_HOST_PORT:-8082}` +- Postgres: `localhost:${DB_HOST_PORT:-5439}` (container-port 5432) + +## Kjøre tjenestene hver for seg + +`pgstac`-nettverket er delt og eksternt, så hver tjeneste kan kjøre alene, f.eks. +`docker compose -f db/docker-compose.yml --env-file .env up -d`. Scanneren når `db` over +det nettverket uten å restarte noe. diff --git a/src/pgstac/api/docker-compose.yml b/src/pgstac/api/docker-compose.yml new file mode 100644 index 0000000..b06e41e --- /dev/null +++ b/src/pgstac/api/docker-compose.yml @@ -0,0 +1,37 @@ +# STAC API: stac-fastapi-pgstac, talks to the db service over the shared network. +# Browser calls this directly (CORS enabled). Requires pgstac >=0.9,<0.10 (db is v0.9.11). +# +# Run standalone: docker compose -f api/docker-compose.yml --env-file .env up -d +# (Requires db running and the shared network: docker network create pgstac) + +services: + api: + image: ghcr.io/stac-utils/stac-fastapi-pgstac:6.2.1 + environment: + APP_HOST: 0.0.0.0 + APP_PORT: 8082 + ENVIRONMENT: local + # DB connection (db = service name on the shared network). + PGHOST: db + PGPORT: 5432 + PGUSER: ${POSTGRES_USER} + PGPASSWORD: ${POSTGRES_PASSWORD} + PGDATABASE: ${POSTGRES_DB} + DB_MIN_CONN_SIZE: 1 + DB_MAX_CONN_SIZE: 5 + # Browser cross-origin access. + CORS_ORIGINS: ${CORS_ORIGINS} + # Ingest is via pypgstac directly; HTTP transactions off by default. + ENABLE_TRANSACTIONS_EXTENSIONS: ${ENABLE_TRANSACTIONS_EXTENSIONS:-FALSE} + # Skip the image's hardcoded wait-for-it (expects host "database"). Startup + # ordering against db is added by the main compose (depends_on healthcheck). + command: ["python", "-m", "stac_fastapi.pgstac.app"] + ports: + - "${API_HOST_PORT:-8082}:8082" + networks: + - pgstac + +networks: + pgstac: + external: true + name: pgstac diff --git a/src/pgstac/db/docker-compose.yml b/src/pgstac/db/docker-compose.yml new file mode 100644 index 0000000..e04eef7 --- /dev/null +++ b/src/pgstac/db/docker-compose.yml @@ -0,0 +1,37 @@ +# pgstac database: PostgreSQL 17 + PostGIS 3, pgstac schema auto-applied on first +# boot via /docker-entrypoint-initdb.d (bundled in the image). Persistent volume. +# +# Run standalone: docker compose -f db/docker-compose.yml --env-file .env up -d +# (Requires the shared network: docker network create pgstac) + +services: + db: + image: ghcr.io/stac-utils/pgstac:v0.9.11 + environment: + POSTGRES_USER: ${POSTGRES_USER} + POSTGRES_PASSWORD: ${POSTGRES_PASSWORD} + POSTGRES_DB: ${POSTGRES_DB} + # The image's init script also reads the PG* vars. + PGUSER: ${POSTGRES_USER} + PGPASSWORD: ${POSTGRES_PASSWORD} + PGDATABASE: ${POSTGRES_DB} + command: postgres -N 500 + ports: + - "${DB_HOST_PORT:-5439}:5432" + volumes: + - pgstac_data:/var/lib/postgresql/data + healthcheck: + test: ["CMD", "pg_isready", "-U", "${POSTGRES_USER}", "-d", "${POSTGRES_DB}"] + interval: 5s + timeout: 5s + retries: 10 + networks: + - pgstac + +volumes: + pgstac_data: + +networks: + pgstac: + external: true + name: pgstac diff --git a/src/pgstac/docker-compose.yml b/src/pgstac/docker-compose.yml new file mode 100644 index 0000000..3407b7e --- /dev/null +++ b/src/pgstac/docker-compose.yml @@ -0,0 +1,26 @@ +# Main compose: brings up the long-running stack (db + api, viewer added later) as +# one project on the shared `pgstac` network. The scanner is NOT here — it runs +# on-demand (scanner/docker-compose.yml) so updating the catalogue never restarts +# the database. +# +# One-time: docker network create pgstac +# Up: docker compose --env-file .env up -d +# Down: docker compose down (add -v to drop the DB volume) + +include: + - db/docker-compose.yml + - api/docker-compose.yml + - viewer/docker-compose.yml + +services: + # Add startup ordering for the integrated stack (kept out of api/ so that file + # can also run standalone without an undefined-service error). + api: + depends_on: + db: + condition: service_healthy + +networks: + pgstac: + external: true + name: pgstac diff --git a/src/pgstac/example.env b/src/pgstac/example.env new file mode 100644 index 0000000..cf6f53f --- /dev/null +++ b/src/pgstac/example.env @@ -0,0 +1,46 @@ +# Copy to .env and fill in. Secrets only here — never commit .env. +# Shared by all services (db, api, viewer, scanner). + +# --- Postgres / pgstac ------------------------------------------------------- +POSTGRES_USER=pgstac +POSTGRES_PASSWORD=changeme +POSTGRES_DB=postgis +# Host port to expose Postgres on (container port is always 5432). +DB_HOST_PORT=5439 + +# --- STAC API (stac-fastapi-pgstac) ----------------------------------------- +# Host port to expose the API on (container listens on 8082). +API_HOST_PORT=8082 +# Origins allowed to call the API from a browser (comma-separated). Must include +# the viewer origin. Use http://localhost:8080 for the local viewer. +CORS_ORIGINS=http://localhost:8080 +# We ingest via pypgstac directly, so the HTTP transactions extension is off. +ENABLE_TRANSACTIONS_EXTENSIONS=FALSE + +# --- Source S3 (scanner + viewer) ------------------------------------------- +# VersityGW/S3 host, NO scheme (e.g. s3.example.no). https:// is added internally. +S3_ENDPOINT=s3.example.norkart.no +S3_ACCESS_KEY=changeme +S3_SECRET_KEY=changeme +S3_REGION=us-east-1 + +# --- Scanner ingest ---------------------------------------------------------- +# The scanner derives its DB connection from the POSTGRES_* vars above (single source +# of truth — the password can't drift). Set PGSTAC_DSN only to override with an +# external/non-default database, e.g.: +# PGSTAC_DSN=postgresql://user:pass@host:5432/dbname +# Absolute base for asset hrefs stored in the DB: /s3//. +# Point at the viewer's public origin (must match how the browser reaches the viewer). +ASSET_BASE_URL=http://localhost:8080 + +# --- Viewer ------------------------------------------------------------------ +# STAC API root the browser opens (browser-reachable URL, i.e. the API's published +# port — NOT the internal service name). Must be allowed in CORS_ORIGINS on the API. +STAC_API_URL=http://localhost:8082 +# Host port for the viewer (container listens on 8080). +VIEWER_HOST_PORT=8080 +# STAC Browser release tag baked into the SPA. +STAC_BROWSER_VERSION=v4.0.1 +# MapLibre basemap style URL for the inline "Data preview" tab. Defaults to the free +# MapLibre demo style; set to a custom style (e.g. the Norkart style) to override. +# BASEMAP_STYLE_URL=https://demotiles.maplibre.org/style.json diff --git a/src/pgstac/scanner/Dockerfile b/src/pgstac/scanner/Dockerfile new file mode 100644 index 0000000..5494acd --- /dev/null +++ b/src/pgstac/scanner/Dockerfile @@ -0,0 +1,20 @@ +# GDAL official base — system GDAL/GEOS/PROJ for rio-stac + geopandas stack. +# TODO: pin to an explicit version tag (e.g. ubuntu-small-3.10.3) once confirmed. +FROM ghcr.io/osgeo/gdal:ubuntu-small-latest + +ENV PYTHONUNBUFFERED=1 \ + PIP_NO_CACHE_DIR=1 \ + PIP_DISABLE_PIP_VERSION_CHECK=1 + +RUN apt-get update \ + && apt-get install -y --no-install-recommends python3-pip \ + && rm -rf /var/lib/apt/lists/* + +WORKDIR /app + +COPY requirements.txt . +RUN pip install --break-system-packages -r requirements.txt + +COPY stac_scan ./stac_scan + +ENTRYPOINT ["python3", "-m", "stac_scan.cli"] diff --git a/src/pgstac/scanner/collection-nseries.yml b/src/pgstac/scanner/collection-nseries.yml new file mode 100644 index 0000000..d1b5031 --- /dev/null +++ b/src/pgstac/scanner/collection-nseries.yml @@ -0,0 +1,44 @@ +# One scan job = one Collection. Copy this file to collection-.yml, rename the +# service, and edit the values. Run it on its own: +# +# docker compose -f scanner/collection-nseries.yml --env-file .env run --rm nseries +# +# Or add it to scanner/docker-compose-all-collections.yml to run with the others. +# (db must be up + shared network: docker network create pgstac) +# Secrets/endpoint/DB/asset-base come from .env — never put credentials here. + +services: + # Service name = how you invoke this scan (`run --rm nseries`). Keep it unique + # across all collection files so the aggregate file can include them together. + nseries: + build: + context: . + image: pgstac-scanner:dev + env_file: ../.env + networks: + - pgstac + command: + # --- REQUIRED ---------------------------------------------------------- + - --bucket + - geolake # bucket name only — NO trailing slash + - --prefix + - gold/kartverket/nseries/ # folder inside the bucket (trailing / is fine) + - --collection-id + - pmtiles-nseries # stable id; reuse to re-scan (upsert + prune) + - --collection-title + - "PMtiles av Norges kuleste kartserie; n-serien" + # --- OPTIONAL (uncomment + edit) --------------------------------------- + # - --collection-description + # - "Longer description shown in the browser" + # - --datetime # force one ISO datetime on all items + # - "2024-01-01T00:00:00Z" + # - --no-recursive # scan only the top level (single line) + # - --fail-on-error # fail the run if any file errors (single line) + # - --no-prune # keep items removed from S3 (single line) + # - --region # override S3 signing region + # - us-east-1 + +networks: + pgstac: + external: true + name: pgstac diff --git a/src/pgstac/scanner/collection-nvdb-vegnett-pluss.yml b/src/pgstac/scanner/collection-nvdb-vegnett-pluss.yml new file mode 100644 index 0000000..fdda06c --- /dev/null +++ b/src/pgstac/scanner/collection-nvdb-vegnett-pluss.yml @@ -0,0 +1,44 @@ +# One scan job = one Collection. Copy this file to collection-.yml, rename the +# service, and edit the values. Run it on its own: +# +# docker compose -f scanner/collection-nvdb-vegnett-pluss.yml --env-file .env run --rm nvdb-vegnett-pluss +# +# Or add it to scanner/docker-compose-all-collections.yml to run with the others. +# (db must be up + shared network: docker network create pgstac) +# Secrets/endpoint/DB/asset-base come from .env — never put credentials here. + +services: + # Service name = how you invoke this scan (`run --rm nvdb-vegnett-pluss`). Keep it unique + # across all collection files so the aggregate file can include them together. + nvdb-vegnett-pluss: + build: + context: . + image: pgstac-scanner:dev + env_file: ../.env + networks: + - pgstac + command: + # --- REQUIRED ---------------------------------------------------------- + - --bucket + - geolake # bucket name only — NO trailing slash + - --prefix + - gold/kartverket/nvdb_vegnett_pluss/parquet/ # folder inside the bucket (trailing / is fine) + - --collection-id + - parquet-nvdb-vegnett-pluss # stable id; reuse to re-scan (upsert + prune) + - --collection-title + - "NVDB vegnett pluss (parquet)" + # --- OPTIONAL (uncomment + edit) --------------------------------------- + # - --collection-description + # - "Longer description shown in the browser" + # - --datetime # force one ISO datetime on all items + # - "2024-01-01T00:00:00Z" + # - --no-recursive # scan only the top level (single line) + # - --fail-on-error # fail the run if any file errors (single line) + # - --no-prune # keep items removed from S3 (single line) + # - --region # override S3 signing region + # - us-east-1 + +networks: + pgstac: + external: true + name: pgstac diff --git a/src/pgstac/scanner/docker-compose-all-collections.yml b/src/pgstac/scanner/docker-compose-all-collections.yml new file mode 100644 index 0000000..bbeb325 --- /dev/null +++ b/src/pgstac/scanner/docker-compose-all-collections.yml @@ -0,0 +1,15 @@ +# Aggregate of every per-collection scan file. Add one `include:` line per collection. +# `include:` does NOT support globs, so list each file explicitly. +# +# Run ALL collections (each container scans then exits): +# docker compose -f scanner/docker-compose-all-collections.yml --env-file .env up --abort-on-container-exit +# +# Or still run just one by service name: +# docker compose -f scanner/docker-compose-all-collections.yml --env-file .env run --rm nseries +# +# Note: `up` runs them in parallel. For sequential runs (gentler on S3/db), invoke each +# file's `run --rm ` in a loop instead. + +include: + - collection-nseries.yml + # - collection-another.yml diff --git a/src/pgstac/scanner/docker-compose.yml b/src/pgstac/scanner/docker-compose.yml new file mode 100644 index 0000000..d17f520 --- /dev/null +++ b/src/pgstac/scanner/docker-compose.yml @@ -0,0 +1,25 @@ +# On-demand scanner: scans one S3 folder -> one pgstac Collection. NOT part of the +# always-up stack (not in the main compose's `up`), so updating the catalogue never +# restarts the database. Reaches the db service over the shared `pgstac` network. +# +# One-time: docker network create pgstac (and the db must be running) +# Run: +# docker compose -f scanner/docker-compose.yml --env-file .env run --rm scanner \ +# --bucket my-bucket --prefix area/sub/ \ +# --collection-id area-sub --collection-title "Area Sub" +# +# CLI args after `scanner` are passed to the entrypoint (see stac_scan/cli.py --help). + +services: + scanner: + build: + context: . + image: pgstac-scanner:dev + env_file: ../.env + networks: + - pgstac + +networks: + pgstac: + external: true + name: pgstac diff --git a/src/pgstac/scanner/requirements.txt b/src/pgstac/scanner/requirements.txt new file mode 100644 index 0000000..de02801 --- /dev/null +++ b/src/pgstac/scanner/requirements.txt @@ -0,0 +1,15 @@ +boto3 +click +pystac[validation] +rio-stac +rasterio +geopandas +pyogrio +shapely +pyproj +numpy +pandas +pyarrow +pmtiles +# Must match the pgstac DB schema (db image is pgstac v0.9.11). +pypgstac[psycopg]>=0.9,<0.10 diff --git a/src/pgstac/scanner/stac_scan/__init__.py b/src/pgstac/scanner/stac_scan/__init__.py new file mode 100644 index 0000000..5c64735 --- /dev/null +++ b/src/pgstac/scanner/stac_scan/__init__.py @@ -0,0 +1,3 @@ +"""Scan an S3 folder and load a STAC Collection into pgstac.""" + +__version__ = "0.1.0" diff --git a/src/pgstac/scanner/stac_scan/build.py b/src/pgstac/scanner/stac_scan/build.py new file mode 100644 index 0000000..8df1cf5 --- /dev/null +++ b/src/pgstac/scanner/stac_scan/build.py @@ -0,0 +1,47 @@ +"""Assemble a pgstac Collection + Item JSON from extracted pystac Items. + +Unlike the static-catalog variant, we do NOT normalize hrefs into a tree: items go +straight into pgstac, which (via stac-fastapi) regenerates self/root/parent/collection +links on output. We only need each item's `collection` field + a Collection with an +extent unioned from its items. +""" + +from __future__ import annotations + +import pystac + + +def build_collection(collection_id, title, description, items): + """Build a pystac.Collection with spatial+temporal extent unioned from items.""" + collection = pystac.Collection( + id=collection_id, + title=title or collection_id, + description=description or collection_id, + extent=pystac.Extent( + spatial=pystac.SpatialExtent([[-180.0, -90.0, 180.0, 90.0]]), + temporal=pystac.TemporalExtent([[None, None]]), + ), + license="proprietary", + ) + for item in items: + collection.add_item(item) + collection.update_extent_from_items() + return collection + + +def to_jsons(collection, items): + """Return (collection_dict, [item_dict, ...]) ready for pypgstac load. + + Links are stripped: pgstac/stac-fastapi regenerate them, and any tree links we'd + emit here point nowhere (we never normalized hrefs). Asset hrefs are preserved. + """ + coll = collection.to_dict(include_self_link=False) + coll["links"] = [] + + item_dicts = [] + for item in items: + d = item.to_dict(include_self_link=False, transform_hrefs=False) + d["collection"] = collection.id + d["links"] = [] + item_dicts.append(d) + return coll, item_dicts diff --git a/src/pgstac/scanner/stac_scan/cli.py b/src/pgstac/scanner/stac_scan/cli.py new file mode 100644 index 0000000..a63dd5f --- /dev/null +++ b/src/pgstac/scanner/stac_scan/cli.py @@ -0,0 +1,205 @@ +"""CLI: scan an S3/VersityGW prefix and load a STAC Collection into pgstac. + +One run = one source = one Collection. Re-runs upsert found items and prune items no +longer in the S3 folder. Secrets + DB DSN + asset base come from env, never CLI flags. +""" + +from __future__ import annotations + +import os +import tempfile +from datetime import datetime, timezone +from urllib.parse import quote + +import click + +from . import build, discover +from .extract_common import FILE_EXT, SkipFile, item_id_from_key +from .extract_pmtiles import build_pmtiles_item +from .extract_raster import build_raster_item +from .extract_vector import build_vector_item +from .pgload import PgLoader +from .s3io import S3Store + +VECTOR_MEDIA = { + "geojson": "application/geo+json", + "geojson_candidate": "application/geo+json", + "geoparquet": "application/x-parquet", +} + + +def _parse_datetime(value: str) -> datetime: + return datetime.fromisoformat(value.replace("Z", "+00:00")) + + +def _env(name: str) -> str: + val = os.environ.get(name) + if not val: + raise click.ClickException(f"Missing required env var: {name}") + return val + + +def _pgstac_dsn() -> str: + """Connection string for pgstac. + + Single source of truth: derived from the same POSTGRES_* vars the db/api use, so the + password can never drift. PGHOST/PGPORT default to the db service on the shared net. + Set PGSTAC_DSN explicitly only to point at an external/non-default database. + """ + override = os.environ.get("PGSTAC_DSN") + if override: + return override + user = quote(_env("POSTGRES_USER"), safe="") + password = quote(_env("POSTGRES_PASSWORD"), safe="") + db = _env("POSTGRES_DB") + host = os.environ.get("PGHOST", "db") + port = os.environ.get("PGPORT", "5432") + return f"postgresql://{user}:{password}@{host}:{port}/{db}" + + +@click.command() +@click.option("--bucket", required=True, help="S3 bucket to scan.") +@click.option("--prefix", required=True, help="Data folder (key prefix) to scan.") +@click.option("--collection-id", required=True, help="STAC Collection id (the source).") +@click.option("--collection-title", default=None, help="Human-readable Collection title.") +@click.option("--collection-description", default=None, help="Collection description.") +@click.option("--datetime", "datetime_override", default=None, + help="ISO datetime applied to all Items (default: S3 LastModified).") +@click.option("--recursive/--no-recursive", default=True, help="Recurse into subfolders.") +@click.option("--fail-on-error", is_flag=True, default=False, + help="Fail the whole run if any file errors (default: skip + warn).") +@click.option("--no-prune", is_flag=True, default=False, + help="Do not delete items missing from the S3 folder (default: prune).") +@click.option("--drop", "drop", is_flag=True, default=False, + help="Drop the collection (and all its items) before loading, then recreate. " + "Use to fully refresh a previously scanned dataset; makes prune moot.") +@click.option("--region", default=None, help="S3 region (overrides S3_REGION env).") +def main(bucket, prefix, collection_id, collection_title, collection_description, + datetime_override, recursive, fail_on_error, no_prune, drop, region): + prefix = prefix if prefix.endswith("/") else prefix + "/" + + store = S3Store( + endpoint=_env("S3_ENDPOINT"), + access_key=_env("S3_ACCESS_KEY"), + secret_key=_env("S3_SECRET_KEY"), + region=region or os.environ.get("S3_REGION", "us-east-1"), + ) + dsn = _pgstac_dsn() + # Base for asset hrefs stored in the DB: /s3//. + # The viewer's /s3 proxy signs + fetches these; bucket in the path => multi-source. + asset_base = _env("ASSET_BASE_URL").rstrip("/") + + override_dt = _parse_datetime(datetime_override) if datetime_override else None + + # --- Discover ---------------------------------------------------------- + files, skipped = discover.discover(store, bucket, prefix, recursive) + click.echo(f"Discovered {len(files)} candidate file(s), skipped {skipped}.") + + # --- Extract ----------------------------------------------------------- + items = [] + failed = [] + seen_ids: dict[str, int] = {} + with tempfile.TemporaryDirectory() as tmp: + for f in files: + key = f["key"] + base_id = item_id_from_key(key) + n = seen_ids.get(base_id, 0) + seen_ids[base_id] = n + 1 + item_id = base_id if n == 0 else f"{base_id}-{n}" + try: + item = _extract_one(store, bucket, key, item_id, f, tmp, + override_dt, asset_base) + items.append(item) + click.echo(f" built {key}") + except SkipFile as exc: + failed.append({"key": key, "reason": f"skipped: {exc}"}) + click.echo(f" skip {key} ({exc})", err=True) + except Exception as exc: # noqa: BLE001 - skip-and-warn by default + failed.append({"key": key, "reason": f"{type(exc).__name__}: {exc}"}) + click.echo(f" ERROR {key} ({exc})", err=True) + + if not items: + raise click.ClickException("No items built — run unsuccessful.") + if fail_on_error and failed: + raise click.ClickException(f"{len(failed)} file(s) errored and --fail-on-error set.") + + # --- Validate (before assembly) --------------------------------------- + # Validate items while they are standalone. We validate here, BEFORE build_collection + # adds tree links: we never normalize hrefs (pgstac/stac-fastapi regenerate links), so + # added links carry None hrefs that fail schema validation spuriously. Those links are + # stripped before load anyway, so this checks the real content (geometry/proj/etc). + invalid = 0 + for item in items: + try: + item.validate() + except Exception as exc: # noqa: BLE001 + invalid += 1 + click.echo(f" invalid {item.id} ({exc})", err=True) + if invalid: + click.echo(f"Warning: {invalid} item(s) failed STAC validation (loaded anyway).", + err=True) + + # --- Assemble ---------------------------------------------------------- + collection = build.build_collection( + collection_id, collection_title, collection_description, items) + coll_dict, item_dicts = build.to_jsons(collection, items) + + # --- Load into pgstac (upsert) + prune -------------------------------- + loader = PgLoader(dsn) + try: + if drop: + dropped = loader.drop_collection(collection_id) + click.echo( + f"Dropped existing collection '{collection_id}' (and its items)." + if dropped else + f"--drop: collection '{collection_id}' did not exist; creating fresh." + ) + loader.upsert_collection(coll_dict) + loader.upsert_items(item_dicts) + click.echo(f"Upserted collection '{collection_id}' + {len(item_dicts)} item(s).") + if drop: + click.echo("Prune skipped (--drop already recreated the collection).") + elif no_prune: + click.echo("Prune skipped (--no-prune).") + else: + keep = {d["id"] for d in item_dicts} + stale = loader.prune(collection_id, keep) + if stale: + click.echo(f"Pruned {len(stale)} item(s) no longer in S3.") + finally: + loader.close() + + click.echo(f"Done. built={len(items)} skipped={skipped} failed={len(failed)}") + + +def _extract_one(store, bucket, key, item_id, meta, tmp, override_dt, asset_base): + dt = override_dt or meta["last_modified"] + if dt.tzinfo is None: + dt = dt.replace(tzinfo=timezone.utc) + source = "override" if override_dt else "s3_last_modified" + asset_href = f"{asset_base}/s3/{bucket}/{key.lstrip('/')}" + # file:size (File extension) applies to every format; size is known from S3 listing. + extra = {"proc:datetime_source": source, "proc:source_key": key} + if meta.get("size"): + extra["file:size"] = int(meta["size"]) + kind = meta["kind"] + + # PMTiles: header-only via byte-range reads, no full download. + if kind == "pmtiles": + item = build_pmtiles_item(store, bucket, key, item_id, dt, asset_href, extra) + else: + local_path = os.path.join(tmp, os.path.basename(key) or item_id) + store.download(bucket, key, local_path) + if kind == "cog": + item = build_raster_item(local_path, item_id, dt, asset_href, extra) + else: + item = build_vector_item(local_path, kind, item_id, dt, asset_href, + VECTOR_MEDIA[kind], extra) + + if "file:size" in item.properties and FILE_EXT not in item.stac_extensions: + item.stac_extensions.append(FILE_EXT) + return item + + +if __name__ == "__main__": + main() diff --git a/src/pgstac/scanner/stac_scan/discover.py b/src/pgstac/scanner/stac_scan/discover.py new file mode 100644 index 0000000..591f0b1 --- /dev/null +++ b/src/pgstac/scanner/stac_scan/discover.py @@ -0,0 +1,50 @@ +"""List objects under a prefix and classify by extension.""" + +from __future__ import annotations + +import os + +# extension -> kind. ".json" is a *candidate* — validated as GeoJSON at extract time. +EXT_MAP = { + ".tif": "cog", + ".tiff": "cog", + ".geojson": "geojson", + ".json": "geojson_candidate", + ".parquet": "geoparquet", + ".geoparquet": "geoparquet", + ".pmtiles": "pmtiles", +} + + +def discover(store, bucket: str, prefix: str, recursive: bool, exclude_prefix: str | None = None): + """Return (files, skipped_count). files = list of dicts with key/kind/last_modified/size. + + Keys under exclude_prefix (the catalog output prefix) are ignored so the scanner + never ingests its own generated catalog JSON. + """ + files = [] + skipped = 0 + for obj in store.list_objects(bucket, prefix): + key = obj["Key"] + if key.endswith("/"): + continue + if exclude_prefix and key.startswith(exclude_prefix): + continue + rel = key[len(prefix):] + if not recursive and "/" in rel.strip("/"): + skipped += 1 + continue + ext = os.path.splitext(key)[1].lower() + kind = EXT_MAP.get(ext) + if not kind: + skipped += 1 + continue + files.append( + { + "key": key, + "kind": kind, + "last_modified": obj["LastModified"], + "size": obj.get("Size", 0), + } + ) + return files, skipped diff --git a/src/pgstac/scanner/stac_scan/extract_common.py b/src/pgstac/scanner/stac_scan/extract_common.py new file mode 100644 index 0000000..d71582a --- /dev/null +++ b/src/pgstac/scanner/stac_scan/extract_common.py @@ -0,0 +1,25 @@ +"""Shared helpers + extension URLs for extractors.""" + +from __future__ import annotations + +import os + +# STAC extension schema URLs (STAC 1.0.0 compatible). +PROJ_EXT = "https://stac-extensions.github.io/projection/v1.1.0/schema.json" +TABLE_EXT = "https://stac-extensions.github.io/table/v1.2.0/schema.json" +FILE_EXT = "https://stac-extensions.github.io/file/v2.1.0/schema.json" + +# PMTiles spec v3 compression enum -> name. +PMTILES_COMPRESSION = {0: "unknown", 1: "none", 2: "gzip", 3: "brotli", 4: "zstd"} + + +class SkipFile(Exception): + """Raised when a file should be skipped (e.g. .json that isn't GeoJSON).""" + + +def item_id_from_key(key: str) -> str: + """Short, filesystem-safe Item id = the file's basename without extension.""" + base = os.path.basename(key.rstrip("/")) + stem = os.path.splitext(base)[0].replace(" ", "_") + # Fallback to the flattened path if a basename can't be derived. + return stem or key.strip("/").replace("/", "_") diff --git a/src/pgstac/scanner/stac_scan/extract_pmtiles.py b/src/pgstac/scanner/stac_scan/extract_pmtiles.py new file mode 100644 index 0000000..1e0c313 --- /dev/null +++ b/src/pgstac/scanner/stac_scan/extract_pmtiles.py @@ -0,0 +1,99 @@ +"""PMTiles -> pystac.Item by reading only the header + metadata (byte-range reads).""" + +from __future__ import annotations + +from datetime import datetime + +import pystac +from pmtiles.reader import Reader +from shapely.geometry import box, mapping + +from .extract_common import PMTILES_COMPRESSION, PROJ_EXT + +PMTILES_MEDIA_TYPE = "application/vnd.pmtiles" + +# pmtiles header tile_type enum -> name (PMTiles spec v3). +TILE_TYPES = {0: "unknown", 1: "mvt", 2: "png", 3: "jpeg", 4: "webp", 5: "avif", 6: "mlt"} + + +def _enum_int(value) -> int: + """Header enums come back as IntEnum (.value) or plain int depending on lib version.""" + return value.value if hasattr(value, "value") else int(value) + + +def build_item_from_reader(get_bytes, item_id, dt, asset_href, extra_props) -> pystac.Item: + reader = Reader(get_bytes) + header = reader.header() + try: + metadata = reader.metadata() or {} + except Exception: # noqa: BLE001 - metadata is optional + metadata = {} + + min_lon = header["min_lon_e7"] / 1e7 + min_lat = header["min_lat_e7"] / 1e7 + max_lon = header["max_lon_e7"] / 1e7 + max_lat = header["max_lat_e7"] / 1e7 + bbox = [min_lon, min_lat, max_lon, max_lat] + geometry = mapping(box(min_lon, min_lat, max_lon, max_lat)) + + tt = header["tile_type"] + tt_int = tt.value if hasattr(tt, "value") else int(tt) + tile_type = TILE_TYPES.get(tt_int, "unknown") + props = { + "proj:epsg": 4326, + "pmtiles:tile_type": tile_type, + "pmtiles:minzoom": int(header["min_zoom"]), + "pmtiles:maxzoom": int(header["max_zoom"]), + "pmtiles:center": [ + header["center_lon_e7"] / 1e7, + header["center_lat_e7"] / 1e7, + int(header["center_zoom"]), + ], + } + if metadata.get("name"): + props["pmtiles:name"] = metadata["name"] + if metadata.get("vector_layers"): + # TileJSON-style layer + field schema (vector tiles only). + props["pmtiles:vector_layers"] = metadata["vector_layers"] + + # Tile-archive internals (free from the header). NOTE: these are TILE counts, not + # feature counts — vector features are duplicated across zoom levels, so a feature + # count is neither cheap nor meaningful here. + for src_key, prop in ( + ("addressed_tiles_count", "pmtiles:addressed_tiles_count"), + ("tile_entries_count", "pmtiles:tile_entries_count"), + ("tile_contents_count", "pmtiles:tile_contents_count"), + ): + if header.get(src_key) is not None: + props[prop] = int(header[src_key]) + if header.get("tile_compression") is not None: + props["pmtiles:tile_compression"] = PMTILES_COMPRESSION.get( + _enum_int(header["tile_compression"]), "unknown" + ) + if header.get("clustered") is not None: + props["pmtiles:clustered"] = bool(header["clustered"]) + if header.get("tile_data_length") is not None: + props["pmtiles:tile_data_bytes"] = int(header["tile_data_length"]) + + props.update(extra_props) + + item = pystac.Item( + id=item_id, + geometry=geometry, + bbox=bbox, + datetime=dt, + properties=props, + ) + item.stac_extensions.append(PROJ_EXT) + item.add_asset( + "data", + pystac.Asset(href=asset_href, media_type=PMTILES_MEDIA_TYPE, roles=["data"]), + ) + return item + + +def build_pmtiles_item(store, bucket, key, item_id, dt, asset_href, extra_props) -> pystac.Item: + def get_bytes(offset: int, length: int) -> bytes: + return store.get_range(bucket, key, offset, length) + + return build_item_from_reader(get_bytes, item_id, dt, asset_href, extra_props) diff --git a/src/pgstac/scanner/stac_scan/extract_raster.py b/src/pgstac/scanner/stac_scan/extract_raster.py new file mode 100644 index 0000000..1dd0bd9 --- /dev/null +++ b/src/pgstac/scanner/stac_scan/extract_raster.py @@ -0,0 +1,52 @@ +"""COG / GeoTIFF -> pystac.Item via rio-stac (proj + raster extensions).""" + +from __future__ import annotations + +from datetime import datetime + +import pystac +import rasterio +from rio_stac.stac import create_stac_item + + +def _cog_internals(local_path: str) -> dict: + """Cheap header-only COG internals (compression, tiling, overviews). Best-effort.""" + out = {} + try: + with rasterio.open(local_path) as ds: + if ds.compression is not None: + out["cog:compression"] = ds.compression.value + prof = ds.profile + if prof.get("tiled") and prof.get("blockxsize"): + out["cog:blocksize"] = [int(prof["blockxsize"]), int(prof["blockysize"])] + out["cog:overview_count"] = len(ds.overviews(1)) if ds.count else 0 + struct = ds.tags(ns="IMAGE_STRUCTURE") + if struct.get("PREDICTOR"): + out["cog:predictor"] = int(struct["PREDICTOR"]) + except Exception: # noqa: BLE001 - internals are optional, never fail the item + pass + return out + + +def build_raster_item( + local_path: str, + item_id: str, + dt: datetime, + asset_href: str, + extra_props: dict, +) -> pystac.Item: + item = create_stac_item( + source=local_path, + id=item_id, + input_datetime=dt, + asset_name="data", + asset_href=asset_href, + asset_media_type=pystac.MediaType.COG, + with_proj=True, + with_raster=True, + ) + if "data" in item.assets: + item.assets["data"].roles = ["data"] + item.properties.update(_cog_internals(local_path)) + item.properties.update(extra_props) + return item diff --git a/src/pgstac/scanner/stac_scan/extract_vector.py b/src/pgstac/scanner/stac_scan/extract_vector.py new file mode 100644 index 0000000..1cbb69f --- /dev/null +++ b/src/pgstac/scanner/stac_scan/extract_vector.py @@ -0,0 +1,269 @@ +"""GeoJSON / GeoParquet -> pystac.Item with proj + table:columns + geometry breakdown. + +GeoParquet is read CHEAPLY: row count, column schema, geometry types, CRS, encoding and +bbox all come from the parquet footer + the GeoParquet `geo` metadata key — no full row +scan. Only when the footer lacks a bbox/geometry-types do we fall back to reading the single +geometry column. GeoJSON has no footer, so it is still fully loaded (then enriched in-memory). +""" + +from __future__ import annotations + +import json +import math +from datetime import datetime + +import geopandas as gpd +import numpy as np +import pandas as pd +import pyarrow as pa +import pyarrow.parquet as pq +import pystac +from pyproj import CRS, Transformer +from shapely import wkb +from shapely.geometry import box, mapping + +from .extract_common import PROJ_EXT, TABLE_EXT, SkipFile + +SAMPLE_ROWS = 5 +FIELD_CAP = 500 # max chars per sampled field value + + +# --------------------------------------------------------------------------- shared + + +def _safe(value): + if value is None: + return None + try: + if np.isscalar(value) and pd.isna(value): + return None + except (TypeError, ValueError): + pass + if isinstance(value, np.integer): + return int(value) + if isinstance(value, np.floating): + f = float(value) + return None if math.isnan(f) else f + if isinstance(value, (np.bool_, bool)): + return bool(value) + if isinstance(value, (int, float)): + return value + if isinstance(value, (bytes, bytearray)): + # WKB geometry column in a sampled parquet row. + try: + return wkb.loads(bytes(value)).wkt[:FIELD_CAP] + except Exception: # noqa: BLE001 + return None + return str(value)[:FIELD_CAP] + + +def _assemble(item_id, dt, geometry, bbox, epsg, props, asset_href, media_type, extra_props): + props["proj:epsg"] = epsg + props.update(extra_props) + item = pystac.Item( + id=item_id, geometry=geometry, bbox=bbox, datetime=dt, properties=props, + ) + item.stac_extensions.extend([PROJ_EXT, TABLE_EXT]) + item.add_asset( + "data", pystac.Asset(href=asset_href, media_type=media_type, roles=["data"]), + ) + return item + + +def _bbox_to_4326(bbox, epsg): + if epsg == 4326 or bbox is None: + return bbox + t = Transformer.from_crs(epsg, 4326, always_xy=True) + minx, miny, maxx, maxy = bbox + xs, ys = t.transform([minx, maxx, minx, maxx], [miny, miny, maxy, maxy]) + return [min(xs), min(ys), max(xs), max(ys)] + + +# --------------------------------------------------------------------------- geoparquet + + +def _arrow_col_type(t: pa.DataType, is_geom: bool) -> str: + if is_geom: + return "geometry" + if pa.types.is_integer(t): + return "int64" + if pa.types.is_floating(t): + return "float64" + if pa.types.is_boolean(t): + return "boolean" + if pa.types.is_temporal(t): + return "datetime" + if pa.types.is_string(t) or pa.types.is_large_string(t): + return "string" + return str(t) + + +def _epsg_from_geo(col_meta: dict) -> int: + crs = col_meta.get("crs") + if not crs: + return 4326 # GeoParquet default is OGC:CRS84 (== EPSG:4326 lon/lat) + try: + # LIMITATION: a custom CRS with no EPSG code falls back to 4326, which would + # then skip bbox reprojection. Doesn't trigger for our data (EPSG:25833). + return CRS.from_user_input(crs).to_epsg() or 4326 + except Exception: # noqa: BLE001 + return 4326 + + +def _parquet_sample(pf: pq.ParquetFile) -> list[dict]: + """First SAMPLE_ROWS rows only — iter_batches stops after the first batch.""" + try: + batch = next(pf.iter_batches(batch_size=SAMPLE_ROWS)) + except StopIteration: + return [] + df = pa.Table.from_batches([batch]).to_pandas() + return [{col: _safe(val) for col, val in row.items()} for _, row in df.iterrows()] + + +def _build_parquet(local_path, item_id, dt, asset_href, media_type, extra_props): + pf = pq.ParquetFile(local_path) + if pf.metadata.num_rows == 0: + raise SkipFile("no features") + + schema = pf.schema_arrow + md = schema.metadata or {} + geo = json.loads(md[b"geo"]) if b"geo" in md else None + if geo is None: + # Not a GeoParquet (no spatial metadata) — fall back to the geopandas path. + return _build_geojson(local_path, "geoparquet", item_id, dt, asset_href, + media_type, extra_props) + + primary = geo.get("primary_column") or next(iter(geo["columns"])) + geom_cols = set(geo.get("columns", {}).keys()) or {primary} + col_meta = geo["columns"][primary] + epsg = _epsg_from_geo(col_meta) + + columns = [ + {"name": f.name, "type": _arrow_col_type(f.type, f.name in geom_cols)} + for f in schema + ] + + geom_types = col_meta.get("geometry_types") + bbox = col_meta.get("bbox") # native CRS, lon/lat order; both optional in the spec. + + if bbox is None or not geom_types: + # Footer lacks bbox/types — read ONLY the geometry column (cheap vs full table). + gs = gpd.read_parquet(local_path, columns=[primary]) + if gs.crs is not None: + epsg = gs.crs.to_epsg() or epsg + if bbox is None: + # We have the geometries: reproject for an exact lon/lat envelope (a + # 4-corner transform underestimates the envelope for projected CRS). + gs_ll = gs.to_crs(4326) if (gs.crs is not None and epsg != 4326) else gs + bbox = [float(v) for v in gs_ll.total_bounds] + else: + bbox = _bbox_to_4326(bbox, epsg) # STAC bbox is always lon/lat 4326. + if not geom_types: + geom_types = sorted(gs.geom_type.dropna().unique().tolist()) + else: + bbox = _bbox_to_4326(bbox, epsg) # proj:epsg keeps the data's native CRS. + + minx, miny, maxx, maxy = bbox + geometry = mapping(box(minx, miny, maxx, maxy)) + + props = { + "table:columns": columns, + "table:row_count": int(pf.metadata.num_rows), + "table:primary_geometry": primary, + "vector:geometry_types": geom_types, + "proc:sample": _parquet_sample(pf), + } + if col_meta.get("encoding"): + props["vector:encoding"] = col_meta["encoding"] + return _assemble(item_id, dt, geometry, bbox, epsg, props, asset_href, + media_type, extra_props) + + +# --------------------------------------------------------------------------- geojson + + +def _validate_geojson(local_path: str) -> None: + """Raise SkipFile if a .json is not a GeoJSON Feature/FeatureCollection/geometry.""" + try: + with open(local_path, "rb") as fh: + head = fh.read(4096).decode("utf-8", errors="ignore") + except OSError as exc: + raise SkipFile(f"unreadable: {exc}") + valid = ( + '"FeatureCollection"' in head + or '"Feature"' in head + or '"coordinates"' in head + ) + if not valid: + raise SkipFile("not GeoJSON") + + +def _gpd_col_type(dtype, is_geom: bool) -> str: + if is_geom: + return "geometry" + kind = getattr(dtype, "kind", "O") + return { + "i": "int64", "u": "int64", "f": "float64", + "b": "boolean", "O": "string", "M": "datetime", + }.get(kind, str(dtype)) + + +def _gpd_sample(gdf: gpd.GeoDataFrame) -> list[dict]: + geom_col = gdf.geometry.name if gdf.geometry is not None else None + rows = [] + for _, row in gdf.head(SAMPLE_ROWS).iterrows(): + record = {} + for col, val in row.items(): + if col == geom_col: + record[col] = (val.wkt[:FIELD_CAP] if val is not None else None) + else: + record[col] = _safe(val) + rows.append(record) + return rows + + +def _build_geojson(local_path, kind, item_id, dt, asset_href, media_type, extra_props): + if kind == "geojson_candidate": + _validate_geojson(local_path) + + gdf = gpd.read_parquet(local_path) if kind == "geoparquet" else gpd.read_file(local_path) + if gdf.empty: + raise SkipFile("no features") + + epsg = gdf.crs.to_epsg() if gdf.crs is not None else 4326 + gdf_ll = gdf.to_crs(4326) if (gdf.crs is not None and epsg != 4326) else gdf + minx, miny, maxx, maxy = (float(v) for v in gdf_ll.total_bounds) + bbox = [minx, miny, maxx, maxy] + geometry = mapping(box(minx, miny, maxx, maxy)) + + geom_col = gdf.geometry.name + columns = [ + {"name": col, "type": _gpd_col_type(dtype, col == geom_col)} + for col, dtype in gdf.dtypes.items() + ] + props = { + "table:columns": columns, + "table:row_count": int(len(gdf)), + "table:primary_geometry": geom_col, + "vector:geometry_types": sorted(gdf.geom_type.dropna().unique().tolist()), + "proc:sample": _gpd_sample(gdf), + } + return _assemble(item_id, dt, geometry, bbox, epsg, props, asset_href, + media_type, extra_props) + + +# --------------------------------------------------------------------------- entry + + +def build_vector_item( + local_path: str, + kind: str, + item_id: str, + dt: datetime, + asset_href: str, + media_type: str, + extra_props: dict, +) -> pystac.Item: + if kind == "geoparquet": + return _build_parquet(local_path, item_id, dt, asset_href, media_type, extra_props) + return _build_geojson(local_path, kind, item_id, dt, asset_href, media_type, extra_props) diff --git a/src/pgstac/scanner/stac_scan/pgload.py b/src/pgstac/scanner/stac_scan/pgload.py new file mode 100644 index 0000000..31f3446 --- /dev/null +++ b/src/pgstac/scanner/stac_scan/pgload.py @@ -0,0 +1,57 @@ +"""Load a Collection + Items into pgstac via pypgstac, then prune removed items. + +pypgstac's load modes (insert/upsert/delsert) never prune items that are absent from +the load set, so 'upsert + prune' is two steps: upsert everything found, then delete +items present in the Collection but not in this scan (mirrors the S3 folder). +""" + +from __future__ import annotations + +from pypgstac.db import PgstacDB +from pypgstac.load import Loader, Methods + + +class PgLoader: + def __init__(self, dsn: str): + self.db = PgstacDB(dsn=dsn, commit_on_exit=True) + self.loader = Loader(db=self.db) + + def upsert_collection(self, collection: dict) -> None: + self.loader.load_collections([collection], insert_mode=Methods.upsert) + + def upsert_items(self, items: list[dict]) -> None: + # iter() so pypgstac's chunked reader consumes it as a stream. + self.loader.load_items(iter(items), insert_mode=Methods.upsert) + + def drop_collection(self, collection_id: str) -> bool: + """Delete the collection (and all its items, via pgstac) if it exists. + + Returns True if a collection was dropped, False if it didn't exist. Used by + --drop to fully recreate a previously scanned dataset instead of upsert+prune. + """ + rows = self.db.query( + "SELECT 1 FROM pgstac.collections WHERE id = %s", [collection_id] + ) + if not list(rows): + return False + # delete_collection cascades to the collection's items/partitions. + list(self.db.query("SELECT pgstac.delete_collection(%s)", [collection_id])) + return True + + def existing_item_ids(self, collection_id: str) -> set[str]: + rows = self.db.query( + "SELECT id FROM pgstac.items WHERE collection = %s", [collection_id] + ) + return {row[0] for row in rows} + + def prune(self, collection_id: str, keep_ids) -> list[str]: + """Delete items in the Collection whose id is not in keep_ids. Returns deleted ids.""" + keep = set(keep_ids) + stale = [i for i in self.existing_item_ids(collection_id) if i not in keep] + for sid in stale: + # delete_item returns void; consume the generator to execute it. + list(self.db.query("SELECT pgstac.delete_item(%s, %s)", [sid, collection_id])) + return stale + + def close(self) -> None: + self.db.close() diff --git a/src/pgstac/scanner/stac_scan/s3io.py b/src/pgstac/scanner/stac_scan/s3io.py new file mode 100644 index 0000000..95633c0 --- /dev/null +++ b/src/pgstac/scanner/stac_scan/s3io.py @@ -0,0 +1,81 @@ +"""S3 / VersityGW access: list, download, write JSON, delete, build https URLs.""" + +from __future__ import annotations + +import json +from typing import Iterator + +import boto3 +from botocore.config import Config + + +class S3Store: + def __init__(self, endpoint: str, access_key: str, secret_key: str, region: str = "us-east-1"): + host = endpoint.replace("https://", "").replace("http://", "").rstrip("/") + self.endpoint_url = f"https://{host}" + self.client = boto3.client( + "s3", + endpoint_url=self.endpoint_url, + aws_access_key_id=access_key, + aws_secret_access_key=secret_key, + region_name=region, + # VersityGW (and most S3-compatible gateways) need path-style addressing. + config=Config(s3={"addressing_style": "path"}), + ) + + # ---- URLs ----------------------------------------------------------- + def http_url(self, bucket: str, key: str) -> str: + return f"{self.endpoint_url}/{bucket}/{key.lstrip('/')}" + + def key_from_url(self, bucket: str, url: str) -> str: + prefix = f"{self.endpoint_url}/{bucket}/" + return url[len(prefix):] if url.startswith(prefix) else url + + # ---- Reads ---------------------------------------------------------- + def list_objects(self, bucket: str, prefix: str) -> Iterator[dict]: + paginator = self.client.get_paginator("list_objects_v2") + for page in paginator.paginate(Bucket=bucket, Prefix=prefix): + for obj in page.get("Contents", []): + yield obj + + def list_json_keys(self, bucket: str, prefix: str) -> set[str]: + return { + obj["Key"] + for obj in self.list_objects(bucket, prefix) + if obj["Key"].endswith(".json") + } + + def download(self, bucket: str, key: str, dest: str) -> None: + self.client.download_file(bucket, key, dest) + + def presign(self, bucket: str, key: str, expires: int) -> str: + """Time-limited GET URL — fetchable without credentials (for STAC clients).""" + return self.client.generate_presigned_url( + "get_object", Params={"Bucket": bucket, "Key": key}, ExpiresIn=expires + ) + + def get_range(self, bucket: str, key: str, offset: int, length: int) -> bytes: + """Read `length` bytes starting at `offset` (for header-only reads, e.g. PMTiles).""" + resp = self.client.get_object( + Bucket=bucket, Key=key, Range=f"bytes={offset}-{offset + length - 1}" + ) + return resp["Body"].read() + + # ---- Writes --------------------------------------------------------- + def put_json(self, bucket: str, key: str, obj: dict) -> None: + self.client.put_object( + Bucket=bucket, + Key=key, + Body=json.dumps(obj, default=str).encode("utf-8"), + ContentType="application/json", + ) + + def delete_keys(self, bucket: str, keys: list[str]) -> None: + for i in range(0, len(keys), 1000): + batch = keys[i:i + 1000] + if not batch: + continue + self.client.delete_objects( + Bucket=bucket, + Delete={"Objects": [{"Key": k} for k in batch]}, + ) diff --git a/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.17.03.png b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.17.03.png new file mode 100644 index 0000000..17acbd7 Binary files /dev/null and b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.17.03.png differ diff --git a/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.17.49.png b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.17.49.png new file mode 100644 index 0000000..69d0190 Binary files /dev/null and b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.17.49.png differ diff --git a/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.18.30.png b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.18.30.png new file mode 100644 index 0000000..16f61dc Binary files /dev/null and b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.18.30.png differ diff --git a/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.18.57.png b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.18.57.png new file mode 100644 index 0000000..e473f35 Binary files /dev/null and b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.18.57.png differ diff --git a/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.19.15.png b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.19.15.png new file mode 100644 index 0000000..75f302a Binary files /dev/null and b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.19.15.png differ diff --git a/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.36.24.png b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.36.24.png new file mode 100644 index 0000000..de35d0f Binary files /dev/null and b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.36.24.png differ diff --git a/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.36.43.png b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.36.43.png new file mode 100644 index 0000000..86b014a Binary files /dev/null and b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.36.43.png differ diff --git a/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.49.35.png b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.49.35.png new file mode 100644 index 0000000..1bb0c51 Binary files /dev/null and b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.49.35.png differ diff --git a/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.49.51.png b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.49.51.png new file mode 100644 index 0000000..3a54fab Binary files /dev/null and b/src/pgstac/screenshots/Screenshot 2026-06-16 at 10.49.51.png differ diff --git "a/src/pgstac/screenshots/pgstac \342\200\224 scan & view workflow.jpg" "b/src/pgstac/screenshots/pgstac \342\200\224 scan & view workflow.jpg" new file mode 100644 index 0000000..8fe1450 Binary files /dev/null and "b/src/pgstac/screenshots/pgstac \342\200\224 scan & view workflow.jpg" differ diff --git "a/src/pgstac/screenshots/pgstac \342\200\224 scan & view workflow2.jpg" "b/src/pgstac/screenshots/pgstac \342\200\224 scan & view workflow2.jpg" new file mode 100644 index 0000000..e6cb068 Binary files /dev/null and "b/src/pgstac/screenshots/pgstac \342\200\224 scan & view workflow2.jpg" differ diff --git a/src/pgstac/viewer/Dockerfile b/src/pgstac/viewer/Dockerfile new file mode 100644 index 0000000..661ab09 --- /dev/null +++ b/src/pgstac/viewer/Dockerfile @@ -0,0 +1,49 @@ +# Multi-stage: stage 1 builds the STAC Browser SPA, stage 2 serves it + the S3 proxy. + +# ---- stage 1: build STAC Browser (pinned release) ------------------------- +ARG STAC_BROWSER_VERSION=v4.0.1 +ARG PATH_PREFIX=/ + +FROM node:20-alpine AS spa +ARG STAC_BROWSER_VERSION +ARG PATH_PREFIX +RUN apk add --no-cache git python3 make g++ +WORKDIR /src +RUN git clone --depth 1 --branch "${STAC_BROWSER_VERSION}" \ + https://github.com/radiantearth/stac-browser.git . +RUN npm install +# MapLibre + pmtiles for the inline "Data preview" tab (renders the item's PMTiles data). +RUN npm install maplibre-gl@^4.7.1 pmtiles@^3.2.0 +# Custom metadata-field rendering for the scanner's proc:*/pmtiles:* properties. +# Overwrites the upstream stub; imported by src/components/Metadata.vue at build. +COPY stac-browser/fields.config.js /src/fields.config.js +# Make "Open in Protomaps" emit the real S3 URL instead of the viewer proxy URL +# (pmtiles.io fetches client-side; s3Base is injected via /config.js). +COPY stac-browser/actions/Protomaps.js /src/src/actions/assets/Protomaps.js +# Inline MapLibre PMTiles preview: a new component + an Item view patched to mount it as a +# lazy "Data preview" tab. Item.vue is PINNED to v4.0.1 (re-sync on STAC_BROWSER_VERSION bump). +COPY stac-browser/components/MapLibrePreview.vue /src/src/components/MapLibrePreview.vue +COPY stac-browser/views/Item.vue /src/src/views/Item.vue +# DYNAMIC_CONFIG: uncomment the runtime |g' public/index.html +# Hash history mode => no server-side SPA routing; deep-link refresh can't 404. +RUN npm run build -- --historyMode=hash --pathPrefix="${PATH_PREFIX}" +# In DYNAMIC_CONFIG mode the app reads the WHOLE config from window.STAC_BROWSER_CONFIG +# (build-time defaults are NOT embedded). Extract the pinned version's complete default +# config so the runtime config.js can merge our overrides onto a full, valid object. +RUN node -e "process.stdout.write(JSON.stringify(require('./config.js')))" > /src/config.default.json + +# ---- stage 2: FastAPI app serving the SPA + signing/rewriting S3 proxy ----- +FROM python:3.12-slim AS app +WORKDIR /app +COPY requirements.txt ./ +RUN pip install --no-cache-dir -r requirements.txt +COPY app ./app +COPY --from=spa /src/dist ./static +COPY --from=spa /src/config.default.json ./config.default.json +ENV STATIC_DIR=/app/static \ + PORT=8080 \ + PATH_PREFIX=/ +EXPOSE 8080 +CMD ["sh", "-c", "uvicorn app.main:app --host 0.0.0.0 --port ${PORT} --proxy-headers --forwarded-allow-ips='*'"] diff --git a/src/pgstac/viewer/app/__init__.py b/src/pgstac/viewer/app/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/pgstac/viewer/app/config.py b/src/pgstac/viewer/app/config.py new file mode 100644 index 0000000..400c985 --- /dev/null +++ b/src/pgstac/viewer/app/config.py @@ -0,0 +1,60 @@ +"""Runtime configuration, all from environment (k8s-friendly, no baked secrets).""" + +from __future__ import annotations + +import os +from dataclasses import dataclass + + +def _require(name: str) -> str: + val = os.environ.get(name, "").strip() + if not val: + raise RuntimeError(f"required env var {name} is not set") + return val + + +@dataclass(frozen=True) +class Settings: + endpoint_host: str # S3 host, no scheme (e.g. s3.example.no) + endpoint_url: str # https:// + access_key: str + secret_key: str + region: str + stac_api_url: str # STAC API root the browser opens (e.g. http://localhost:8082) + path_prefix: str # sub-path the app is served under (default "/") + port: int + basemap_style_url: str # MapLibre style URL for the inline "Data preview" basemap + + @property + def s3_base(self) -> str: + """Absolute S3 URL prefix (https://host/). Asset keys carry the bucket, so the + proxied /s3// maps back to // for the + client-side 'Open in Protomaps' action.""" + return f"{self.endpoint_url}/" + + +def load_settings() -> Settings: + raw_host = _require("S3_ENDPOINT") + host = raw_host.replace("https://", "").replace("http://", "").rstrip("/") + + prefix = os.environ.get("PATH_PREFIX", "/").strip() or "/" + if not prefix.startswith("/"): + prefix = "/" + prefix + if not prefix.endswith("/"): + prefix = prefix + "/" + + return Settings( + endpoint_host=host, + endpoint_url=f"https://{host}", + access_key=_require("S3_ACCESS_KEY"), + secret_key=_require("S3_SECRET_KEY"), + region=os.environ.get("S3_REGION", "us-east-1").strip() or "us-east-1", + stac_api_url=_require("STAC_API_URL").rstrip("/"), + path_prefix=prefix, + port=int(os.environ.get("PORT", "8080")), + # Free MapLibre demo style by default; override (e.g. the Norkart style) via env. + basemap_style_url=( + os.environ.get("BASEMAP_STYLE_URL", "").strip() + or "https://demotiles.maplibre.org/style.json" + ), + ) diff --git a/src/pgstac/viewer/app/main.py b/src/pgstac/viewer/app/main.py new file mode 100644 index 0000000..aae38f8 --- /dev/null +++ b/src/pgstac/viewer/app/main.py @@ -0,0 +1,89 @@ +"""FastAPI app: serves the STAC Browser SPA + a signing S3 asset proxy. + +The catalog is served by the STAC API (stac-fastapi-pgstac) directly — the browser +talks to it cross-origin (CORS) — so this app only needs to: serve the SPA, inject +runtime config (pointing catalogUrl at the API), and sign asset reads from the private +bucket. There is NO /catalog rewrite route (that was for the static-catalog variant). + +Routes (origin-relative; PATH_PREFIX-aware via uvicorn root_path): + GET /healthz liveness + GET /config.js runtime STAC Browser config injected from env + GET /s3/{bucket}/{key} asset bytes from S3, signed, Range-aware (COG/PMTiles) + GET /* STAC Browser static SPA (hash-routed) +""" + +from __future__ import annotations + +import json +import os + +from fastapi import FastAPI, Header, Request, Response +from fastapi.responses import JSONResponse, StreamingResponse +from fastapi.staticfiles import StaticFiles + +from .config import load_settings +from .s3proxy import S3Reader, stream_body + +settings = load_settings() +reader = S3Reader(settings) + +STATIC_DIR = os.environ.get("STATIC_DIR", "/app/static") + +# Full default config extracted from the pinned STAC Browser at build time. In +# DYNAMIC_CONFIG mode the app reads the entire config from window.STAC_BROWSER_CONFIG, +# so we must serve a complete object (not just our overrides) or it crashes on boot. +_DEFAULTS_PATH = os.environ.get("DEFAULT_CONFIG_PATH", "/app/config.default.json") +try: + with open(_DEFAULTS_PATH, encoding="utf-8") as fh: + DEFAULT_CONFIG = json.load(fh) +except FileNotFoundError: + DEFAULT_CONFIG = {} + +app = FastAPI(title="pgstac viewer", root_path=settings.path_prefix.rstrip("/")) + + +@app.get("/healthz") +def healthz() -> dict: + return {"status": "ok", "api": settings.stac_api_url} + + +@app.get("/config.js") +def config_js() -> Response: + """STAC Browser runtime config. Opens the STAC API directly, hash history mode. + + catalogUrl is the STAC API root (absolute) — the browser fetches catalog/collection/ + item JSON from it cross-origin (the API enables CORS for this origin). Asset hrefs in + that JSON already point at this viewer's /s3 proxy (written by the scanner). + """ + cfg = dict(DEFAULT_CONFIG) + cfg.update({ + "catalogUrl": settings.stac_api_url, + "historyMode": "hash", + "pathPrefix": settings.path_prefix, + # Real S3 base (https:///) so the client-side "Open in Protomaps" + # action can map a proxied /s3// href back to the real S3 URL. + "s3Base": settings.s3_base, + # MapLibre basemap style for the inline "Data preview" tab (MapLibrePreview.vue). + "basemapStyleUrl": settings.basemap_style_url, + }) + body = f"window.STAC_BROWSER_CONFIG = {json.dumps(cfg)};\n" + return Response(content=body, media_type="application/javascript", + headers={"Cache-Control": "no-store"}) + + +@app.get("/s3/{bucket}/{key:path}") +def get_asset(bucket: str, key: str, range: str | None = Header(default=None)) -> Response: + try: + obj = reader.get_object(bucket, key, range) + except KeyError: + return JSONResponse({"error": "not found", "bucket": bucket, "key": key}, + status_code=404) + return StreamingResponse( + stream_body(obj["body"]), + status_code=obj["status"], + headers=obj["headers"], + ) + + +# Static SPA last so explicit API routes win. html=True serves index.html at "/". +app.mount("/", StaticFiles(directory=STATIC_DIR, html=True), name="spa") diff --git a/src/pgstac/viewer/app/s3proxy.py b/src/pgstac/viewer/app/s3proxy.py new file mode 100644 index 0000000..ad7038e --- /dev/null +++ b/src/pgstac/viewer/app/s3proxy.py @@ -0,0 +1,67 @@ +"""Signed reads against VersityGW/S3 — mirrors the scanner's client conventions. + +Asset hrefs stored in pgstac carry the bucket in the path (/s3//), so the +proxy signs reads against whatever bucket the request names — supporting multiple source +buckets behind one viewer. +""" + +from __future__ import annotations + +from typing import Iterator, Optional + +import boto3 +from botocore.config import Config +from botocore.exceptions import ClientError + +from .config import Settings + + +class S3Reader: + def __init__(self, s: Settings): + self.client = boto3.client( + "s3", + endpoint_url=s.endpoint_url, + aws_access_key_id=s.access_key, + aws_secret_access_key=s.secret_key, + region_name=s.region, + # VersityGW (and most S3-compatible gateways) need path-style addressing. + config=Config(s3={"addressing_style": "path"}), + ) + + def get_object(self, bucket: str, key: str, range_header: Optional[str]) -> dict: + """Range-aware passthrough for asset bytes (COG/PMTiles need ranges). + + Returns dict with: body (StreamingBody), status, headers (dict to forward). + """ + kwargs: dict = {"Bucket": bucket, "Key": key} + if range_header: + kwargs["Range"] = range_header + try: + resp = self.client.get_object(**kwargs) + except ClientError as e: + code = e.response.get("Error", {}).get("Code", "") + if code in ("NoSuchKey", "404", "NoSuchBucket"): + raise KeyError(key) from e + raise + + meta = resp["ResponseMetadata"] + status = meta.get("HTTPStatusCode", 200) + headers: dict[str, str] = {"Accept-Ranges": "bytes"} + for src, dst in ( + ("ContentType", "Content-Type"), + ("ContentLength", "Content-Length"), + ("ContentRange", "Content-Range"), + ("ETag", "ETag"), + ): + val = resp.get(src) + if val is not None: + headers[dst] = str(val) + return {"body": resp["Body"], "status": status, "headers": headers} + + +def stream_body(body, chunk_size: int = 64 * 1024) -> Iterator[bytes]: + while True: + chunk = body.read(chunk_size) + if not chunk: + break + yield chunk diff --git a/src/pgstac/viewer/docker-compose.yml b/src/pgstac/viewer/docker-compose.yml new file mode 100644 index 0000000..3c55b1b --- /dev/null +++ b/src/pgstac/viewer/docker-compose.yml @@ -0,0 +1,24 @@ +# Viewer: STAC Browser SPA + signing /s3 asset proxy. Browser opens the STAC API +# directly (catalogUrl); this container only serves the SPA + signs private-bucket +# asset reads. Long-running; on the shared `pgstac` network. +# +# Run standalone: docker compose -f viewer/docker-compose.yml --env-file .env up -d --build +# Open http://localhost:8080 + +services: + viewer: + build: + context: . + args: + STAC_BROWSER_VERSION: ${STAC_BROWSER_VERSION:-v4.0.1} + image: pgstac-viewer:dev + env_file: ../.env + ports: + - "${VIEWER_HOST_PORT:-8080}:8080" + networks: + - pgstac + +networks: + pgstac: + external: true + name: pgstac diff --git a/src/pgstac/viewer/requirements.txt b/src/pgstac/viewer/requirements.txt new file mode 100644 index 0000000..070fee6 --- /dev/null +++ b/src/pgstac/viewer/requirements.txt @@ -0,0 +1,3 @@ +fastapi==0.115.6 +uvicorn[standard]==0.34.0 +boto3==1.35.92 diff --git a/src/pgstac/viewer/stac-browser/actions/Protomaps.js b/src/pgstac/viewer/stac-browser/actions/Protomaps.js new file mode 100644 index 0000000..b6f7dbf --- /dev/null +++ b/src/pgstac/viewer/stac-browser/actions/Protomaps.js @@ -0,0 +1,51 @@ +import AssetActionPlugin from "../AssetActionPlugin"; +import URI from 'urijs'; +import i18n from "../../i18n"; + +// obj & ply files are usually with mime-type text/plain +const PROTOMAPS_SUPPORTED_TYPES = [ + 'application/vnd.pmtiles', +]; + +// Viewer override: the viewer rewrites asset hrefs to its own signing proxy +// (/s3/). pmtiles.io is an EXTERNAL viewer that runs in the +// browser and fetches the archive itself, so it can't use that proxy URL (it would +// hit the viewer's localhost/ingress origin, which it can't reach/authenticate). +// Map the proxied href back to the real S3 URL (//), which +// the user's machine can reach directly. s3Base is injected via /config.js. +function toS3Url(href) { + const cfg = (typeof window !== 'undefined' && window.STAC_BROWSER_CONFIG) || {}; + const base = cfg.s3Base; // e.g. https://s3.example.no/my-bucket/ + if (!base || typeof href !== 'string') { + return href; + } + const marker = '/s3/'; + const idx = href.indexOf(marker); + if (idx === -1) { + return href; // not a proxied asset href — leave as-is + } + const key = href.substring(idx + marker.length); + return (base.endsWith('/') ? base : base + '/') + key; +} + +export default class Protomaps extends AssetActionPlugin { + + get show() { + // Rather check if .pmtiles substring present in this.asset.href or simply this.component.filename.endsWith('pmtiles') + return this.component.isBrowserProtocol && ( + PROTOMAPS_SUPPORTED_TYPES.includes(this.asset.type) + || URI(this.asset.href).suffix() == 'pmtiles' + ); + } + + get uri() { + let uri = new URI("https://pmtiles.io/"); + uri.addQuery("url", toS3Url(this.component.href)); // real S3 URL, not the proxy + return uri; + } + + get text() { + return i18n.t('actions.openIn', {service: 'Protomaps'}); + } + +} diff --git a/src/pgstac/viewer/stac-browser/components/MapLibrePreview.vue b/src/pgstac/viewer/stac-browser/components/MapLibrePreview.vue new file mode 100644 index 0000000..b59b360 --- /dev/null +++ b/src/pgstac/viewer/stac-browser/components/MapLibrePreview.vue @@ -0,0 +1,309 @@ + + + + + + diff --git a/src/pgstac/viewer/stac-browser/fields.config.js b/src/pgstac/viewer/stac-browser/fields.config.js new file mode 100644 index 0000000..9b45636 --- /dev/null +++ b/src/pgstac/viewer/stac-browser/fields.config.js @@ -0,0 +1,154 @@ +// Custom STAC Browser metadata-field rendering for the scanner. +// +// The scanner (stac_scan) writes custom Item properties that stac-fields doesn't +// know about, so by default they render with auto-generated labels and raw JSON. +// Here we register them so they show as proper, labelled rows: +// proc:* provenance written for every item (source object, datetime origin) +// proc:sample first sampled feature rows (vector) -> rendered as an HTML table +// pmtiles:* PMTiles header/metadata (tile type, zoom range, center, layers) +// +// This is imported by STAC Browser's src/components/Metadata.vue at BUILD time, +// so it is baked into the SPA — changing labels here means rebuilding the image, +// not just tweaking runtime /config.js. +// +// Registry/field-spec reference (pinned to stac-fields ~1.5.7, shipped with +// STAC Browser v4.0.1): https://github.com/stac-utils/stac-fields +// A custom `formatter` may return HTML (rendered via v-html in MetadataTable.vue); +// all data-derived strings MUST be escaped with Helper.e() to avoid breaking layout. + +import { Registry, Helper } from '@radiantearth/stac-fields'; + +const e = Helper.e; + +function cell(value) { + if (value === null || typeof value === 'undefined' || value === '') { + return '—'; // em dash for empty + } + return e(String(value)); +} + +// proc:sample is an array of row objects whose columns are the dataset's own +// (dynamic per dataset), so a fixed `items` table schema can't describe it. +// Build the column set as the union of keys across the sampled rows and render +// a plain HTML table ourselves. +function formatSampleRows(rows) { + if (!Array.isArray(rows) || rows.length === 0) { + return cell(null); + } + const cols = []; + for (const row of rows) { + if (row && typeof row === 'object') { + for (const key of Object.keys(row)) { + if (!cols.includes(key)) { + cols.push(key); + } + } + } + } + // Wide values (esp. geometry WKT) would otherwise stretch the table off-screen. Cap each + // cell's width with an ellipsis and put the whole table in a horizontal-scroll box; the + // full value stays available via the cell's title (hover) tooltip. + const wrapStyle = 'max-width:100%;overflow-x:auto;'; + const cellStyle = 'max-width:240px;overflow:hidden;text-overflow:ellipsis;white-space:nowrap;'; + const head = cols.map(c => `${e(c)}`).join(''); + const body = rows.map(row => { + const cells = cols.map(c => { + const raw = row ? row[c] : null; + const title = (raw === null || typeof raw === 'undefined' || raw === '') + ? '' : ` title="${e(String(raw))}"`; + return `${cell(raw)}`; + }).join(''); + return `${cells}`; + }).join(''); + return `
${head}` + + `${body}
`; +} + +// pmtiles:vector_layers is a TileJSON-style list; each layer carries a `fields` +// object (attribute name -> type) that is also dynamic, so render manually. +function formatVectorLayers(layers) { + if (!Array.isArray(layers) || layers.length === 0) { + return cell(null); + } + const rows = layers.map(layer => { + const id = cell(layer && layer.id); + const minz = layer && typeof layer.minzoom !== 'undefined' ? layer.minzoom : '?'; + const maxz = layer && typeof layer.maxzoom !== 'undefined' ? layer.maxzoom : '?'; + const zoom = (layer && (layer.minzoom != null || layer.maxzoom != null)) + ? `${e(String(minz))}–${e(String(maxz))}` : '—'; + const fields = (layer && layer.fields && typeof layer.fields === 'object') + ? Object.entries(layer.fields).map(([k, v]) => `${e(k)}: ${e(String(v))}`).join('
') + : '—'; + const desc = cell(layer && layer.description); + return `${id}${zoom}${fields}${desc}`; + }).join(''); + return `` + + `` + + `${rows}`; +} + +// pmtiles:center is [lon, lat, zoom] from the PMTiles header. +function formatCenter(value) { + if (!Array.isArray(value) || value.length < 2) { + return cell(value); + } + const [lon, lat, zoom] = value; + const ll = `${cell(lat)}, ${cell(lon)}`; + return (zoom !== null && typeof zoom !== 'undefined') + ? `${ll} (zoom ${e(String(zoom))})` : ll; +} + +// --- proc: provenance written by the scanner --------------------- +Registry.addExtension('proc', 'Provenance & processing'); + +Registry.addMetadataField('proc:source_key', { + label: 'Source object', + explain: 'Key of the source file in the bucket this item was extracted from.' +}); + +Registry.addMetadataField('proc:datetime_source', { + label: 'Datetime source', + explain: 'Where the item datetime came from.', + mapping: { + s3_last_modified: 'S3 last-modified', + override: 'Manual override' + } +}); + +Registry.addMetadataField('proc:sample', { + label: 'Sample rows', + explain: 'First feature rows read from the dataset when it was scanned.', + formatter: formatSampleRows +}); + +// --- pmtiles: PMTiles header + metadata ------------------------------------ +Registry.addExtension('pmtiles', 'PMTiles'); + +Registry.addMetadataField('pmtiles:name', { label: 'Name' }); + +Registry.addMetadataField('pmtiles:tile_type', { + label: 'Tile type', + mapping: { + mvt: 'Vector (MVT)', + png: 'Raster (PNG)', + jpeg: 'Raster (JPEG)', + webp: 'Raster (WebP)', + avif: 'Raster (AVIF)', + unknown: 'Unknown' + } +}); + +Registry.addMetadataField('pmtiles:minzoom', { label: 'Min zoom' }); +Registry.addMetadataField('pmtiles:maxzoom', { label: 'Max zoom' }); + +Registry.addMetadataField('pmtiles:center', { + label: 'Center', + explain: 'Default map center from the PMTiles header: longitude, latitude, zoom.', + formatter: formatCenter +}); + +Registry.addMetadataField('pmtiles:vector_layers', { + label: 'Vector layers', + explain: 'Layers and attribute fields declared in the PMTiles metadata.', + formatter: formatVectorLayers +}); diff --git a/src/pgstac/viewer/stac-browser/views/Item.vue b/src/pgstac/viewer/stac-browser/views/Item.vue new file mode 100644 index 0000000..df50ae0 --- /dev/null +++ b/src/pgstac/viewer/stac-browser/views/Item.vue @@ -0,0 +1,164 @@ + + + + + +