Files
geointel/backend/app/services/source_freshness_service.py
T
Codex 0befca02f8
GeoIntel CI / docs-smoke (push) Canceled after 0s
GeoIntel CI / contract-smoke (push) Canceled after 0s
Prefer governed orthophoto editions in source status
2026-07-17 02:33:15 +02:00

318 lines
14 KiB
Python

from __future__ import annotations
from collections import defaultdict
from dataclasses import dataclass
from datetime import datetime, timedelta, timezone
from pathlib import Path
import re
from typing import Iterable
from uuid import UUID
from sqlalchemy.orm import Session
from app.core.errors import AppError
from app.models import Dataset, DatasetVersion, Project
from app.schemas.source_freshness import (
SourceFreshnessItem,
SourceFreshnessReport,
SourceFreshnessSummary,
SourceIntegritySummary,
)
@dataclass(frozen=True)
class SourcePolicy:
display_name: str
refresh_policy: str
review_interval_days: int | None = None
SOURCE_POLICIES: dict[str, SourcePolicy] = {
"grb": SourcePolicy("GRB gebouwen en context", "rolling_snapshot", 90),
"vrbg": SourcePolicy("VRBG wegenregister", "rolling_snapshot", 90),
"digitaal_vlaanderen_buildings_addresses_register": SourcePolicy(
"Gebouwen- en adressenregister", "rolling_snapshot", 90
),
"digitaal_vlaanderen_orthophoto": SourcePolicy("Orthofoto Vlaanderen", "rolling_snapshot", 180),
"agentschap_landbouw_zeevisserij_agricultural_parcels": SourcePolicy(
"Landbouwgebruikspercelen", "annual_release"
),
"department_omgeving_land_use": SourcePolicy("Landgebruik Vlaanderen", "annual_release"),
"inbo_bwk_natura2000": SourcePolicy("BWK en Natura 2000", "annual_release"),
"statbel": SourcePolicy("Statbel bevolking", "annual_release"),
"waterinfo": SourcePolicy("Waterinfo meetreeksen", "annual_release"),
"digitaal_vlaanderen_dhmv": SourcePolicy("Digitaal Hoogtemodel Vlaanderen", "edition"),
"department_omgeving_thematic_raster": SourcePolicy("Omgeving thematische rasters", "edition"),
"dov_soil_map": SourcePolicy("DOV bodemkaart", "edition"),
"vmm_flood_hazard": SourcePolicy("VMM overstromingskaarten", "scenario"),
"historical_landuse": SourcePolicy("Historisch landgebruik", "archive"),
"manual": SourcePolicy("Handmatig ingeladen gegevens", "local"),
"fixture": SourcePolicy("Test- en demonstratiegegevens", "local"),
"map_selection": SourcePolicy("Bewaarde kaartselecties", "local"),
}
DEFAULT_POLICY = SourcePolicy("Niet-geclassificeerde bron", "edition")
_ORTHOPHOTO_EDITION = re.compile(r"^20\d{2}\.\d{2}$")
def _as_utc(value: datetime | None) -> datetime | None:
if value is None:
return None
if value.tzinfo is None:
return value.replace(tzinfo=timezone.utc)
return value.astimezone(timezone.utc)
def _source_key(dataset: Dataset) -> str:
return (dataset.source_name or dataset.source or "unknown").strip().lower() or "unknown"
def _latest_datetime(values: Iterable[datetime | None]) -> datetime | None:
normalized = [_as_utc(value) for value in values if value is not None]
return max(normalized) if normalized else None
def _latest_dataset(datasets: list[Dataset]) -> Dataset:
return max(
datasets,
key=lambda item: (
_as_utc(item.observed_at) or datetime.min.replace(tzinfo=timezone.utc),
_as_utc(item.imported_at) or datetime.min.replace(tzinfo=timezone.utc),
str(item.id),
),
)
def _latest_source_version(source_name: str, policy: SourcePolicy, datasets: list[Dataset]) -> str | None:
versioned = [item for item in datasets if item.source_version]
if not versioned:
return None
if source_name == "digitaal_vlaanderen_orthophoto":
official_editions = [
item for item in versioned if _ORTHOPHOTO_EDITION.fullmatch((item.source_version or "").strip())
]
if official_editions:
return _latest_dataset(official_editions).source_version
if policy.refresh_policy == "rolling_snapshot":
named_current = [
item
for item in versioned
if any(token in (item.source_version or "").lower() for token in ("most_recent", "latest", "current"))
]
if named_current:
return _latest_dataset(named_current).source_version
return _latest_dataset(versioned).source_version
def _is_local_storage_path(storage_path: str) -> bool:
normalized = storage_path.strip().lower()
return bool(normalized) and "://" not in normalized and not normalized.startswith("/vsi")
def _has_historical_series(datasets: list[Dataset]) -> bool:
observations_by_series: dict[str, set[datetime]] = defaultdict(set)
for dataset in datasets:
observed_at = _as_utc(dataset.observed_at)
if dataset.temporal_series_key and observed_at is not None:
observations_by_series[dataset.temporal_series_key].add(observed_at)
return any(len(observations) > 1 for observations in observations_by_series.values())
def _integrity_summary(datasets: list[Dataset], versions_by_dataset: dict[UUID, list[DatasetVersion]]) -> SourceIntegritySummary:
summary = SourceIntegritySummary()
for dataset in datasets:
versions = versions_by_dataset.get(dataset.id, [])
latest_version = max(versions, key=lambda item: item.version) if versions else None
if dataset.status == "ready" and latest_version is None:
summary.missing_version_count += 1
if (
latest_version is not None
and dataset.checksum_sha256
and latest_version.checksum_sha256
and dataset.checksum_sha256 != latest_version.checksum_sha256
):
summary.checksum_mismatch_count += 1
if dataset.storage_path and _is_local_storage_path(dataset.storage_path):
path = Path(dataset.storage_path)
try:
if not path.is_file():
summary.missing_storage_file_count += 1
elif dataset.size_bytes is not None and path.stat().st_size != dataset.size_bytes:
summary.size_mismatch_count += 1
except OSError:
summary.missing_storage_file_count += 1
return summary
def _classify_source(
policy: SourcePolicy,
datasets: list[Dataset],
integrity: SourceIntegritySummary,
now: datetime,
) -> tuple[str, datetime | None, str, str]:
latest_imported = _latest_datetime(item.imported_at for item in datasets)
latest_observed = _latest_datetime(item.observed_at for item in datasets)
has_source_version = any(bool((item.source_version or "").strip()) for item in datasets)
if integrity.issue_count:
return (
"review_required",
None,
"De bewaarde dataset- en versie-evidentie bevat een integriteitsafwijking.",
"Controleer opslag, checksum en datasetversies voordat deze bron opnieuw wordt gebruikt.",
)
if policy.refresh_policy == "local":
return (
"local",
None,
"Deze bron is lokaal aangemaakt en heeft geen externe publicatiecyclus.",
"Geen bronverversing nodig; beheer de lokale dataset via de bestaande werkstroom.",
)
if policy.refresh_policy == "rolling_snapshot":
if latest_imported is None:
return (
"review_required",
None,
"De importdatum voor deze rollende bron ontbreekt.",
"Controleer de provenance voordat een nieuwe begrensde import wordt gestart.",
)
next_review = latest_imported + timedelta(days=policy.review_interval_days or 90)
if next_review <= now:
return (
"due",
next_review,
"De lokale snapshot heeft zijn geplande controledatum bereikt.",
"Vergelijk de broncatalogus en voer alleen daarna een begrensde, expliciete verversing uit.",
)
return (
"current",
next_review,
"De lokale snapshot valt binnen de afgesproken controleperiode.",
"Geen actie nodig tot de volgende controledatum.",
)
if policy.refresh_policy == "annual_release":
if latest_observed is None:
return (
"review_required",
None,
"De recentste waarnemings- of editieperiode ontbreekt.",
"Vul eerst officiële tijds- en versieprovenance aan; download niets automatisch.",
)
next_review = datetime(latest_observed.year + 2, 1, 1, tzinfo=timezone.utc)
if latest_observed.year < now.year - 1:
return (
"due",
next_review,
f"De recentste bewaarde jaargang is {latest_observed.year}.",
"Controleer of de officiële bron een recentere definitieve jaargang publiceerde.",
)
return (
"current",
next_review,
f"De recentste bewaarde jaargang is {latest_observed.year}.",
"Controleer bij de volgende publicatiecyclus of een nieuwe definitieve jaargang beschikbaar is.",
)
if not has_source_version:
return (
"review_required",
None,
"Deze vaste publicatie heeft geen herkenbare bronversie.",
"Leg de officiële editie of scenarioversie vast voordat de bron als gecontroleerd geldt.",
)
return (
"current",
None,
"Dit is een vaste editie, scenario- of archiefpublicatie met vastgelegde bronversie.",
"Vervang deze editie niet automatisch; voeg een nieuwe officiële editie als afzonderlijke versie toe.",
)
class SourceFreshnessService:
@staticmethod
def build_report(
project_id: UUID,
datasets: list[Dataset],
versions: list[DatasetVersion],
*,
now: datetime | None = None,
) -> SourceFreshnessReport:
generated_at = _as_utc(now) or datetime.now(timezone.utc)
versions_by_dataset: dict[UUID, list[DatasetVersion]] = defaultdict(list)
for version in versions:
versions_by_dataset[version.dataset_id].append(version)
datasets_by_source: dict[str, list[Dataset]] = defaultdict(list)
for dataset in datasets:
datasets_by_source[_source_key(dataset)].append(dataset)
items: list[SourceFreshnessItem] = []
for source_name, source_datasets in datasets_by_source.items():
policy = SOURCE_POLICIES.get(source_name, DEFAULT_POLICY)
integrity = _integrity_summary(source_datasets, versions_by_dataset)
status, next_review_at, reason, recommended_action = _classify_source(
policy, source_datasets, integrity, generated_at
)
if policy is DEFAULT_POLICY and not integrity.issue_count:
status = "review_required"
next_review_at = None
reason = "Voor deze bron is nog geen expliciete publicatie- of controlecyclus vastgelegd."
recommended_action = "Classificeer de bron eerst als snapshot, jaargang, vaste editie, scenario, archief of lokaal."
items.append(
SourceFreshnessItem(
source_name=source_name,
display_name=policy.display_name if policy is not DEFAULT_POLICY else source_name.replace("_", " ").title(),
dataset_count=len(source_datasets),
ready_count=sum(item.status == "ready" for item in source_datasets),
version_count=sum(len(versions_by_dataset.get(item.id, [])) for item in source_datasets),
latest_imported_at=_latest_datetime(item.imported_at for item in source_datasets),
latest_observed_at=_latest_datetime(item.observed_at for item in source_datasets),
latest_source_version=_latest_source_version(source_name, policy, source_datasets),
refresh_policy=policy.refresh_policy,
review_interval_days=policy.review_interval_days,
next_review_at=next_review_at,
status=status,
historical_series=_has_historical_series(source_datasets),
reason=reason,
recommended_action=recommended_action,
integrity=integrity,
)
)
status_rank = {"review_required": 0, "due": 1, "current": 2, "local": 3}
items.sort(key=lambda item: (status_rank[item.status], item.display_name.lower()))
integrity_issue_count = sum(item.integrity.issue_count for item in items)
summary = SourceFreshnessSummary(
source_count=len(items),
dataset_count=len(datasets),
current_count=sum(item.status == "current" for item in items),
due_count=sum(item.status == "due" for item in items),
review_required_count=sum(item.status == "review_required" for item in items),
local_count=sum(item.status == "local" for item in items),
sources_with_integrity_issues=sum(item.integrity.issue_count > 0 for item in items),
integrity_issue_count=integrity_issue_count,
)
return SourceFreshnessReport(
project_id=project_id,
generated_at=generated_at,
summary=summary,
items=items,
limitations=[
"Deze controle leest uitsluitend lokale dataset-, versie- en opslaggegevens.",
"Er worden geen externe catalogi bevraagd, bestanden gedownload of datasets overschreven.",
"Een vaste editie of scenario-publicatie wordt niet verouderd genoemd alleen omdat de publicatiedatum oud is.",
],
)
@staticmethod
def audit_project(db: Session, project_id: UUID, *, now: datetime | None = None) -> SourceFreshnessReport:
if not db.get(Project, project_id):
raise AppError(code="PROJECT_NOT_FOUND", message="Project not found", status_code=404)
datasets = db.query(Dataset).filter(Dataset.project_id == project_id).all()
dataset_ids = [dataset.id for dataset in datasets]
versions = (
db.query(DatasetVersion).filter(DatasetVersion.dataset_id.in_(dataset_ids)).all()
if dataset_ids
else []
)
return SourceFreshnessService.build_report(project_id, datasets, versions, now=now)