Files
geointel/scripts/provision_mol_context_layers.py
T
Codex c0943fe7d4
GeoIntel CI / docs-smoke (push) Canceled after 0s
GeoIntel CI / contract-smoke (push) Canceled after 0s
feat: add temporal Mol explorer
2026-07-14 15:19:42 +02:00

470 lines
19 KiB
Python

"""Provision official Mol roads, water and parcel context layers.
This explicit operator command reads Digitaal Vlaanderen OGC API Features,
clips every feature to the official Mol municipality boundary and imports the
result through the existing GeoIntel dataset upload API. It is idempotent and
never runs automatically during application startup.
"""
from __future__ import annotations
import argparse
import hashlib
import json
import os
import sys
from dataclasses import dataclass
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Iterable
import requests
from requests.adapters import HTTPAdapter
from shapely.geometry import mapping, shape
from shapely.validation import make_valid
from urllib3.util.retry import Retry
MUNICIPALITY_NAME = "Mol"
MUNICIPALITY_NIS_CODE = "13025"
PROJECT_NAME = "Mol Municipality Workbench"
AREA_NAME = "Gemeente Mol - officiele grens"
BOUNDARY_URL = "https://geo.api.vlaanderen.be/VRBG/ogc/features/v1/collections/Refgem/items"
GRB_COLLECTION_URL = "https://geo.api.vlaanderen.be/GRB/ogc/features/v1/collections/{collection}/items"
ATTRIBUTION = "Bron: Grootschalig Referentie Bestand Vlaanderen, Digitaal Vlaanderen"
DEFAULT_OUTPUT_DIR = Path("/app/storage/operator-data/mol-context")
DEFAULT_BOUNDARY_PATH = Path("/app/storage/operator-data/mol-municipality/mol_municipality_boundary.geojson")
DEFAULT_API_URL = "http://127.0.0.1:8000"
RETRYABLE_STATUS_CODES = (429, 500, 502, 503, 504)
@dataclass(frozen=True)
class LayerDefinition:
key: str
filename: str
collections: tuple[str, ...]
reference_layer_name: str
layer_type: str
LAYERS = (
LayerDefinition("roads", "mol_grb_roads.geojson", ("Wegsegment",), "roads", "road"),
LayerDefinition("water", "mol_grb_water.geojson", ("WTZ", "WLAS", "WGR"), "water", "water"),
LayerDefinition("parcels", "mol_grb_parcels.geojson", ("ADP",), "parcels", "parcel"),
)
def parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description="Provision official GRB context layers for the municipality of Mol.")
parser.add_argument("--base-url", default=os.environ.get("GEOINTEL_INTERNAL_API_URL", DEFAULT_API_URL))
parser.add_argument("--project-name", default=PROJECT_NAME)
parser.add_argument("--output-dir", type=Path, default=Path(os.environ.get("MOL_CONTEXT_OUTPUT_DIR", DEFAULT_OUTPUT_DIR)))
parser.add_argument("--boundary-path", type=Path, default=Path(os.environ.get("MOL_BOUNDARY_PATH", DEFAULT_BOUNDARY_PATH)))
parser.add_argument("--page-limit", type=int, default=int(os.environ.get("MOL_CONTEXT_PAGE_LIMIT", "1000")))
parser.add_argument("--max-features", type=int, default=int(os.environ.get("MOL_CONTEXT_MAX_FEATURES", "150000")))
parser.add_argument("--request-timeout", type=int, default=int(os.environ.get("MOL_CONTEXT_REQUEST_TIMEOUT", "180")))
parser.add_argument("--import-timeout", type=int, default=int(os.environ.get("MOL_CONTEXT_IMPORT_TIMEOUT", "1800")))
parser.add_argument("--layers", default="roads,water,parcels", help="Comma-separated subset: roads,water,parcels")
parser.add_argument("--force", action="store_true", help="Refetch artifacts even when a complete artifact exists.")
parser.add_argument("--fetch-only", action="store_true")
return parser.parse_args()
def utc_now() -> str:
return datetime.now(timezone.utc).isoformat()
def build_session() -> requests.Session:
retry = Retry(
total=5,
connect=5,
read=5,
status=5,
backoff_factor=1.0,
status_forcelist=RETRYABLE_STATUS_CODES,
allowed_methods=frozenset({"GET"}),
raise_on_status=True,
)
session = requests.Session()
session.headers.update({"User-Agent": "GeoIntel-Mol-Context-Operator/1.0"})
adapter = HTTPAdapter(max_retries=retry)
session.mount("https://", adapter)
session.mount("http://", adapter)
return session
def next_page_url(payload: dict[str, Any]) -> str | None:
for link in payload.get("links") or []:
if link.get("rel") == "next" and link.get("href"):
return str(link["href"])
return None
def normalized_geometry(payload: dict[str, Any] | None):
if not payload:
return None
geometry = shape(payload)
if geometry.is_empty:
return None
if not geometry.is_valid:
geometry = make_valid(geometry)
return geometry if not geometry.is_empty and geometry.is_valid else None
def geometry_dimension(geometry) -> int:
geometry_type = geometry.geom_type
if "Polygon" in geometry_type:
return 2
if "LineString" in geometry_type or geometry_type == "LinearRing":
return 1
if "Point" in geometry_type:
return 0
if geometry_type == "GeometryCollection":
return max((geometry_dimension(part) for part in geometry.geoms), default=-1)
return -1
def load_boundary(path: Path, session: requests.Session, timeout: int):
if path.exists():
payload = json.loads(path.read_text(encoding="utf-8"))
features = payload.get("features") or []
if len(features) == 1:
boundary = normalized_geometry(features[0].get("geometry"))
if boundary is not None:
return boundary
response = session.get(
BOUNDARY_URL,
params={"f": "application/geo+json", "limit": "10", "filter": "NAAM='Mol'", "filter-lang": "cql2-text"},
timeout=timeout,
)
response.raise_for_status()
matches = [
feature
for feature in response.json().get("features") or []
if str((feature.get("properties") or {}).get("NISCODE")) == MUNICIPALITY_NIS_CODE
]
if len(matches) != 1:
raise RuntimeError(f"Expected one official Mol boundary, received {len(matches)}")
boundary = normalized_geometry(matches[0].get("geometry"))
if boundary is None:
raise RuntimeError("Official Mol boundary is invalid")
return boundary
def iter_collection_pages(
session: requests.Session,
collection: str,
bbox: tuple[float, float, float, float],
*,
page_limit: int,
timeout: int,
) -> Iterable[tuple[dict[str, Any], str]]:
url: str | None = GRB_COLLECTION_URL.format(collection=collection)
params = {
"f": "application/geo+json",
"limit": str(page_limit),
"bbox": ",".join(f"{value:.8f}" for value in bbox),
}
first = True
seen: set[str] = set()
while url:
if url in seen:
raise RuntimeError(f"Pagination loop for GRB collection {collection}")
seen.add(url)
response = session.get(url, params=params if first else None, timeout=timeout)
first = False
response.raise_for_status()
payload = response.json()
yield payload, response.url
url = next_page_url(payload)
def build_layer(
definition: LayerDefinition,
boundary,
session: requests.Session,
*,
page_limit: int,
max_features: int,
timeout: int,
) -> tuple[dict[str, Any], dict[str, Any]]:
features: list[dict[str, Any]] = []
source_urls: list[str] = []
seen_ids: set[str] = set()
rejected = 0
clipped = 0
for collection in definition.collections:
for page, source_url in iter_collection_pages(
session,
collection,
boundary.bounds,
page_limit=page_limit,
timeout=timeout,
):
source_urls.append(source_url)
for source_feature in page.get("features") or []:
raw_id = str(source_feature.get("id") or "")
feature_id = f"{collection}:{raw_id}" if raw_id else hashlib.sha256(
json.dumps(source_feature.get("geometry"), sort_keys=True).encode("utf-8")
).hexdigest()
if feature_id in seen_ids:
continue
seen_ids.add(feature_id)
geometry = normalized_geometry(source_feature.get("geometry"))
if geometry is None or not geometry.intersects(boundary):
rejected += 1
continue
if not geometry.within(boundary):
source_dimension = geometry_dimension(geometry)
geometry = normalized_geometry(mapping(geometry.intersection(boundary)))
if geometry is not None and geometry_dimension(geometry) < source_dimension:
geometry = None
clipped += 1
if geometry is None:
rejected += 1
continue
if len(features) >= max_features:
raise RuntimeError(
f"{definition.key} exceeds the {max_features} feature cap; refusing a truncated artifact"
)
properties = dict(source_feature.get("properties") or {})
properties.update(
{
"source_name": "grb",
"source_collection": collection,
"source_feature_id": feature_id,
"reference_layer_name": definition.reference_layer_name,
"layer_type": definition.layer_type,
"authority_level": "authoritative",
"coverage_scope": "municipality",
"municipality": MUNICIPALITY_NAME,
"nis_code": MUNICIPALITY_NIS_CODE,
"attribution": ATTRIBUTION,
}
)
features.append(
{"type": "Feature", "id": feature_id, "geometry": mapping(geometry), "properties": properties}
)
if not features:
raise RuntimeError(f"No {definition.key} features intersect Mol")
generated_at = utc_now()
payload = {
"type": "FeatureCollection",
"name": f"GRB {definition.key} - complete municipality Mol",
"features": features,
"source": f"Digitaal Vlaanderen GRB OGC API collections {', '.join(definition.collections)}",
"attribution": ATTRIBUTION,
"municipality": MUNICIPALITY_NAME,
"nis_code": MUNICIPALITY_NIS_CODE,
"coverage_scope": "municipality",
"reference_truncated": False,
"generated_at": generated_at,
}
summary = {
"layer": definition.key,
"collections": list(definition.collections),
"feature_count": len(features),
"rejected_count": rejected,
"clipped_count": clipped,
"pages_fetched": len(source_urls),
"source_urls": source_urls,
"generated_at": generated_at,
"reference_truncated": False,
}
return payload, summary
def response_data(response: requests.Response) -> Any:
try:
payload = response.json()
except ValueError as exc:
raise RuntimeError(f"GeoIntel API returned non-JSON ({response.status_code}): {response.text[:300]}") from exc
if not response.ok:
raise RuntimeError(f"GeoIntel API failed ({response.status_code}): {json.dumps(payload, ensure_ascii=False)[:800]}")
if not isinstance(payload, dict) or "data" not in payload:
raise RuntimeError("GeoIntel API response does not use the canonical data envelope")
return payload["data"]
def locate_workspace(session: requests.Session, base_url: str, project_name: str, timeout: int) -> tuple[str, str, list[dict[str, Any]]]:
projects = response_data(session.get(f"{base_url}/api/v1/projects", params={"limit": 200}, timeout=timeout))
project = next((item for item in projects.get("items") or [] if item.get("name") == project_name), None)
if not project:
raise RuntimeError(f"Project {project_name!r} is missing; run provision_mol_municipality_workspace.py first")
project_id = str(project["id"])
areas = response_data(session.get(f"{base_url}/api/v1/projects/{project_id}/areas", params={"limit": 200}, timeout=timeout))
area_items = list(areas.get("items") or [])
area = next((item for item in area_items if "gemeente mol" in str(item.get("name", "")).lower()), None)
if not area:
area = max(area_items, key=lambda item: float(item.get("area_m2") or 0), default=None)
if not area:
raise RuntimeError("Mol municipality area is missing")
datasets = response_data(session.get(f"{base_url}/api/v1/projects/{project_id}/datasets", params={"limit": 200}, timeout=timeout))
return project_id, str(area["id"]), list(datasets.get("items") or [])
def upload_layer(
session: requests.Session,
base_url: str,
project_id: str,
area_id: str,
definition: LayerDefinition,
path: Path,
summary: dict[str, Any],
timeout: int,
) -> dict[str, Any]:
source_metadata = {
"provider": "Digitaal Vlaanderen",
"collections": list(definition.collections),
"authority_level": "authoritative",
"coverage_scope": "municipality",
"municipality": MUNICIPALITY_NAME,
"nis_code": MUNICIPALITY_NIS_CODE,
"feature_count": summary["feature_count"],
"attribution": ATTRIBUTION,
}
provenance_metadata = {
"operator_tool": "provision_mol_context_layers.py",
"operator_explicit_fetch": True,
"generated_at": summary["generated_at"],
"source_urls": summary["source_urls"],
"reference_truncated": False,
}
with path.open("rb") as handle:
response = session.post(
f"{base_url}/api/v1/projects/{project_id}/datasets/upload",
data={
"dataset_type": "vector",
"source": "operator_official_import",
"dataset_role": "reference",
"source_name": "grb",
"reference_layer_name": definition.reference_layer_name,
"source_metadata_json": json.dumps(source_metadata, ensure_ascii=False),
"provenance_metadata_json": json.dumps(provenance_metadata, ensure_ascii=False),
"area_id": area_id,
"temporal_series_key": f"grb:{definition.key}:mol",
"observed_at": summary["generated_at"],
"temporal_granularity": "snapshot",
"source_version": summary["generated_at"][:10],
},
files={"file": (path.name, handle, "application/geo+json")},
timeout=timeout,
)
return response_data(response)
def ensure_temporal_metadata(
session: requests.Session,
base_url: str,
project_id: str,
dataset: dict[str, Any],
definition: LayerDefinition,
generated_at: str,
timeout: int,
) -> dict[str, Any]:
series_key = f"grb:{definition.key}:mol"
if dataset.get("temporal_series_key") == series_key and dataset.get("observed_at"):
return dataset
response = session.patch(
f"{base_url}/api/v1/projects/{project_id}/datasets/{dataset['id']}/temporal",
json={
"temporal_series_key": series_key,
"observed_at": generated_at,
"temporal_granularity": "snapshot",
"source_version": generated_at[:10],
},
timeout=timeout,
)
return response_data(response)
def main() -> int:
args = parse_args()
requested = {item.strip().lower() for item in args.layers.split(",") if item.strip()}
selected = [definition for definition in LAYERS if definition.key in requested]
unknown = requested - {definition.key for definition in LAYERS}
if unknown or not selected:
print(json.dumps({"status": "error", "message": f"Unsupported layers: {sorted(unknown)}"}), file=sys.stderr)
return 2
if args.page_limit <= 0 or args.max_features <= 0:
print(json.dumps({"status": "error", "message": "Limits must be positive"}), file=sys.stderr)
return 2
args.output_dir.mkdir(parents=True, exist_ok=True)
results: list[dict[str, Any]] = []
try:
with build_session() as source_session:
boundary = load_boundary(args.boundary_path, source_session, args.request_timeout)
prepared: list[tuple[LayerDefinition, Path, dict[str, Any]]] = []
for definition in selected:
path = args.output_dir / definition.filename
summary_path = path.with_suffix(".manifest.json")
if not args.force and path.exists() and summary_path.exists():
summary = json.loads(summary_path.read_text(encoding="utf-8"))
else:
payload, summary = build_layer(
definition,
boundary,
source_session,
page_limit=args.page_limit,
max_features=args.max_features,
timeout=args.request_timeout,
)
path.write_text(json.dumps(payload, ensure_ascii=False, separators=(",", ":")), encoding="utf-8")
summary_path.write_text(json.dumps(summary, ensure_ascii=False, indent=2), encoding="utf-8")
prepared.append((definition, path, summary))
if args.fetch_only:
results = [
{"layer": definition.key, "path": str(path), "feature_count": summary["feature_count"], "status": "prepared"}
for definition, path, summary in prepared
]
else:
base_url = args.base_url.rstrip("/")
with requests.Session() as api_session:
project_id, area_id, existing = locate_workspace(
api_session, base_url, args.project_name, args.import_timeout
)
for definition, path, summary in prepared:
dataset = next((item for item in existing if item.get("original_filename") == path.name), None)
if dataset:
dataset = ensure_temporal_metadata(
api_session,
base_url,
project_id,
dataset,
definition,
summary["generated_at"],
args.import_timeout,
)
results.append(
{"layer": definition.key, "dataset_id": dataset["id"], "feature_count": dataset.get("feature_count"), "status": "existing"}
)
continue
dataset = upload_layer(
api_session,
base_url,
project_id,
area_id,
definition,
path,
summary,
args.import_timeout,
)
results.append(
{"layer": definition.key, "dataset_id": dataset["id"], "feature_count": dataset.get("feature_count"), "status": "imported"}
)
except (OSError, RuntimeError, requests.RequestException, ValueError, KeyError) as exc:
print(json.dumps({"status": "error", "message": str(exc)}, ensure_ascii=False), file=sys.stderr)
return 1
print(json.dumps({"status": "ok", "municipality": MUNICIPALITY_NAME, "layers": results}, ensure_ascii=False, indent=2))
return 0
if __name__ == "__main__":
sys.exit(main())