Files
geointel/backend/app/services/grb_refresh_plan_service.py
T
Codex eed6ee796e
GeoIntel CI / docs-smoke (push) Canceled after 0s
GeoIntel CI / contract-smoke (push) Canceled after 0s
Add governed GRB refresh workflow
2026-07-16 17:59:00 +02:00

187 lines
8.4 KiB
Python

from __future__ import annotations
import re
from dataclasses import dataclass
from datetime import date, datetime, timezone
from uuid import UUID
from sqlalchemy.orm import Session
from app.core.errors import AppError
from app.models import Dataset, Project
from app.schemas.grb_refresh import GrbRefreshLayerPlan, GrbRefreshPlan, GrbRefreshPlanSummary
from app.services.source_catalog_probe_service import SourceCatalogProbeService
@dataclass(frozen=True)
class _LayerDefinition:
theme: str
display_name: str
collections: tuple[str, ...]
@property
def series_key(self) -> str:
return f"grb:{self.theme}:kempen-transport-region"
class GrbRefreshPlanService:
SCOPE = "kempen-transport-region"
LAYERS = (
_LayerDefinition("buildings", "Gebouwen", ("GBG",)),
_LayerDefinition("roads", "Wegen", ("Wegsegment",)),
_LayerDefinition("water", "Water", ("WTZ", "WLAS", "WGR")),
_LayerDefinition("parcels", "Percelen", ("ADP",)),
)
_EDITION_DATE = re.compile(r"(?<!\d)(20\d{2}-\d{2}-\d{2})(?!\d)")
@classmethod
def _parse_edition_date(cls, value: str | None) -> date | None:
match = cls._EDITION_DATE.search(value or "")
if not match:
return None
try:
return date.fromisoformat(match.group(1))
except ValueError:
return None
@staticmethod
def _dataset_feature_count(dataset: Dataset) -> int | None:
metadata = dataset.metadata_json if isinstance(dataset.metadata_json, dict) else {}
value = metadata.get("feature_count")
try:
return int(value) if value is not None else None
except (TypeError, ValueError):
return None
@staticmethod
def _latest_dataset(rows: list[Dataset], series_key: str) -> Dataset | None:
candidates = [
row
for row in rows
if row.temporal_series_key == series_key and row.status == "ready"
]
if not candidates:
return None
minimum = datetime.min.replace(tzinfo=timezone.utc)
return max(candidates, key=lambda row: (row.observed_at or row.imported_at or row.created_at or minimum, str(row.id)))
@classmethod
def build(
cls,
db: Session,
project_id: UUID,
*,
scope: str = SCOPE,
refresh_catalog: bool = False,
now: datetime | None = None,
) -> GrbRefreshPlan:
if scope != cls.SCOPE:
raise AppError(
code="GRB_REFRESH_SCOPE_UNSUPPORTED",
message=f"Only the governed scope '{cls.SCOPE}' is supported",
status_code=400,
)
if not db.get(Project, project_id):
raise AppError(code="PROJECT_NOT_FOUND", message="Project not found", status_code=404)
generated_at = now or datetime.now(timezone.utc)
catalog = SourceCatalogProbeService.audit_project(db, project_id, force=refresh_catalog)
grb_probe = next((item for item in catalog.items if item.source_name == "grb"), None)
remote_version = grb_probe.remote_version if grb_probe else None
remote_edition_date = cls._parse_edition_date(remote_version)
remote_available = bool(grb_probe and grb_probe.status == "available" and grb_probe.reachable)
rows = (
db.query(Dataset)
.filter(Dataset.project_id == project_id, Dataset.source_name == "grb")
.all()
)
layer_plans: list[GrbRefreshLayerPlan] = []
for definition in cls.LAYERS:
local = cls._latest_dataset(rows, definition.series_key)
local_date = cls._parse_edition_date(local.source_version if local else None)
if not remote_available:
status = "remote_unavailable"
action = "De officiële catalogus is niet bereikbaar; er wordt geen vernieuwingsbeslissing genomen."
elif remote_edition_date is None:
status = "review_required"
action = "De officiële editie bevat geen herkenbare datum en vereist menselijke beoordeling."
elif local is None:
status = "not_loaded"
action = "Deze laag kan als nieuwe, afzonderlijke GRB-snapshot worden voorbereid."
elif local_date is None:
status = "review_required"
action = "De lokale editie is niet veilig datumvergelijkbaar; controleer de provenance vóór staging."
elif local_date == remote_edition_date:
status = "current"
action = "De lokale snapshot gebruikt dezelfde officiële editie; geen import nodig."
elif local_date < remote_edition_date:
status = "update_available"
action = "Stage eerst alle bronartifacts en controleer exacte aantallen en checksums vóór import."
else:
status = "review_required"
action = "De lokale editie lijkt nieuwer dan de catalogus; automatische terugval is verboden."
layer_plans.append(
GrbRefreshLayerPlan(
theme=definition.theme,
display_name=definition.display_name,
collections=list(definition.collections),
temporal_series_key=definition.series_key,
status=status,
local_dataset_id=local.id if local else None,
local_source_version=local.source_version if local else None,
local_observed_at=local.observed_at if local else None,
local_imported_at=local.imported_at if local else None,
local_feature_count=cls._dataset_feature_count(local) if local else None,
local_size_bytes=local.size_bytes if local else None,
action_message=action,
)
)
counts = {status: sum(item.status == status for item in layer_plans) for status in (
"current", "update_available", "not_loaded", "review_required", "remote_unavailable"
)}
actionable = counts["update_available"] + counts["not_loaded"]
summary = GrbRefreshPlanSummary(
layer_count=len(layer_plans),
current_count=counts["current"],
update_available_count=counts["update_available"],
not_loaded_count=counts["not_loaded"],
review_required_count=counts["review_required"],
remote_unavailable_count=counts["remote_unavailable"],
new_dataset_count_if_applied=actionable,
retained_dataset_count=sum(item.local_dataset_id is not None for item in layer_plans),
current_feature_count=sum(item.local_feature_count or 0 for item in layer_plans),
current_size_bytes=sum(item.local_size_bytes or 0 for item in layer_plans),
)
if actionable:
message = (
f"{actionable} GRB-laag{' is' if actionable == 1 else 'en zijn'} voorbereidbaar voor editie "
f"{remote_edition_date.isoformat() if remote_edition_date else remote_version}. "
"Staging berekent eerst de exacte impact; import vereist daarna de plan-checksum."
)
elif counts["current"] == len(layer_plans):
message = "Alle beheerde regionale GRB-lagen gebruiken de officiële cataloguseditie."
else:
message = "Er is menselijke beoordeling nodig voordat een GRB-staging kan starten."
return GrbRefreshPlan(
project_id=project_id,
scope=scope,
generated_at=generated_at,
remote_status=grb_probe.status if grb_probe else "unavailable",
remote_version=remote_version,
remote_edition_date=remote_edition_date,
catalog_checked_at=grb_probe.checked_at if grb_probe else None,
summary=summary,
layers=layer_plans,
message=message,
limitations=[
"Dit endpoint is read-only en start geen download, import of databasejob.",
"Staging bewaart bronartifacts buiten PostGIS; apply vereist de exacte staged plan-checksum.",
"Een refresh maakt nieuwe immutable Datasets en verwijdert of overschrijft oude snapshots niet.",
"Exacte feature- en opslagverschillen zijn pas bekend nadat alle regionale partitions staged en gevalideerd zijn.",
],
)