Add governed source freshness audit
This commit is contained in:
@@ -0,0 +1,302 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from collections import defaultdict
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from pathlib import Path
|
||||
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")
|
||||
|
||||
|
||||
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 _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 _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."
|
||||
latest = _latest_dataset(source_datasets)
|
||||
observed_values = {value for item in source_datasets if (value := _as_utc(item.observed_at)) is not None}
|
||||
temporal_keys = {item.temporal_series_key for item in source_datasets if item.temporal_series_key}
|
||||
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=next(
|
||||
(
|
||||
item.source_version
|
||||
for item in sorted(
|
||||
source_datasets,
|
||||
key=lambda value: (
|
||||
_as_utc(value.observed_at) or datetime.min.replace(tzinfo=timezone.utc),
|
||||
_as_utc(value.imported_at) or datetime.min.replace(tzinfo=timezone.utc),
|
||||
),
|
||||
reverse=True,
|
||||
)
|
||||
if item.source_version
|
||||
),
|
||||
latest.source_version,
|
||||
),
|
||||
refresh_policy=policy.refresh_policy,
|
||||
review_interval_days=policy.review_interval_days,
|
||||
next_review_at=next_review_at,
|
||||
status=status,
|
||||
historical_series=len(observed_values) > 1 or len(temporal_keys) > 1,
|
||||
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)
|
||||
Reference in New Issue
Block a user