Files
geointel/backend/app/services/source_registry_service.py
T

1510 lines
58 KiB
Python

"""Server-owned source identity, snapshot and lineage persistence primitives.
This service deliberately does not inspect caller-provided ``source_name``
metadata. A governed adapter must select one of the static definitions below,
create an immutable snapshot and bind that exact snapshot to a dataset/version.
"""
from __future__ import annotations
from dataclasses import dataclass
from datetime import datetime
import re
from typing import Any
from uuid import UUID
from sqlalchemy.orm import Session
from app.core.errors import AppError
from app.models import (
Dataset,
DatasetLineageEdge,
DatasetQuarantine,
DatasetVersion,
SourceRegistry,
SourceSnapshot,
)
_CHECKSUM = re.compile(r"^[a-fA-F0-9]{64}$")
_VALIDATION_STATUSES = {"not_validated", "passed", "failed"}
_PROVENANCE_STATUSES = {"complete", "incomplete", "not_applicable"}
_LINEAGE_STATUSES = {"complete", "incomplete", "not_applicable"}
_FRESHNESS_STATUSES = {
"unknown",
"current",
"due",
"stale",
"not_applicable",
"review_required",
}
_SNAPSHOT_INGEST_STATUSES = {
"registered",
"configured",
"not_configured",
"available",
"ingested",
"failed",
"quarantined",
"legacy_unverified",
}
_MAX_LINEAGE_GRAPH_NODES = 50_000
def _usage_policy(
*,
ground_truth_allowed: bool = False,
training_allowed: bool = False,
allowed_tasks: tuple[str, ...] = (),
validation_authority: dict[str, str] | None = None,
) -> dict[str, Any]:
return {
"automatic_ground_truth": False,
"ground_truth_allowed": ground_truth_allowed,
"training_allowed": training_allowed,
"allowed_tasks": list(allowed_tasks),
"validation_authority": validation_authority or {},
}
@dataclass(frozen=True)
class SourceRegistryDefinition:
source_key: str
display_name: str
classification: str
authority_name: str
authority_scope: dict[str, Any]
default_crs: str = "unknown"
default_units: str = "unknown"
provider_adapter_key: str | None = None
source_url: str | None = None
license_name: str = "Provider terms must be verified for each immutable snapshot."
usage_restrictions: str = (
"Use only according to the source-specific snapshot terms and attribution."
)
spatial_resolution: dict[str, Any] | None = None
temporal_coverage: dict[str, Any] | None = None
geographic_coverage: dict[str, Any] | None = None
expected_geometry_types: tuple[str, ...] = ()
expected_attributes: dict[str, Any] | None = None
usage_policy: dict[str, Any] | None = None
freshness_status: str = "unknown"
ingest_status: str = "registered"
known_limitations: tuple[str, ...] = ()
def as_model_values(self) -> dict[str, Any]:
return {
"source_key": self.source_key,
"display_name": self.display_name,
"classification": self.classification,
"authority_name": self.authority_name,
"authority_scope_json": dict(self.authority_scope),
"provider_adapter_key": self.provider_adapter_key,
"source_url": self.source_url,
"license_name": self.license_name,
"usage_restrictions": self.usage_restrictions,
"default_crs": self.default_crs,
"default_units": self.default_units,
"spatial_resolution_json": dict(
self.spatial_resolution or {"status": "unknown"}
),
"temporal_coverage_json": dict(
self.temporal_coverage or {"status": "unknown"}
),
"geographic_coverage_json": dict(
self.geographic_coverage or {"status": "unknown"}
),
"expected_geometry_types_json": list(self.expected_geometry_types),
"expected_attributes_json": dict(
self.expected_attributes or {"status": "unknown"}
),
"usage_policy_json": dict(self.usage_policy or _usage_policy()),
"freshness_status": self.freshness_status,
"ingest_status": self.ingest_status,
"known_limitations_json": list(
self.known_limitations
or (
"No authority, ground-truth, freshness or training claim is allowed without a governed snapshot and passed contract.",
)
),
"registry_metadata_json": {
"registry_owner": "server",
"definition_version": "phase2-v1",
},
}
def _definition(
source_key: str,
display_name: str,
classification: str,
authority_name: str,
authority_scope: dict[str, Any],
**kwargs: Any,
) -> SourceRegistryDefinition:
return SourceRegistryDefinition(
source_key=source_key,
display_name=display_name,
classification=classification,
authority_name=authority_name,
authority_scope=authority_scope,
**kwargs,
)
_VECTOR_CONTEXT = _usage_policy(allowed_tasks=("reference_context",))
_IMAGERY_CONTEXT = _usage_policy(
training_allowed=True,
allowed_tasks=("imagery", "training_input", "visual_context"),
)
_REGIONAL_BUILDING_LABELS = _usage_policy(
ground_truth_allowed=True,
training_allowed=True,
allowed_tasks=("building_validation", "building_labels"),
validation_authority={"building_validation": "regional_primary_pending_contract"},
)
SERVER_OWNED_SOURCE_DEFINITIONS: dict[str, SourceRegistryDefinition] = {
definition.source_key: definition
for definition in (
_definition(
"grb",
"Grootschalig Referentie Bestand",
"authoritative",
"Digitaal Vlaanderen",
{"zone": "Flanders", "themes": ["buildings", "roads", "water", "parcels"]},
provider_adapter_key="grb",
source_url="https://www.vlaanderen.be/datavindplaats/catalogus/basiskaart-vlaanderen-grb",
default_crs="EPSG:31370",
default_units="metres",
geographic_coverage={"zone": "Flanders"},
expected_geometry_types=(
"Polygon",
"MultiPolygon",
"LineString",
"MultiLineString",
),
expected_attributes={
"required": ["id"],
"layers": ["GBG", "Wegsegment", "WTZ", "WLAS", "WGR", "ADP"],
},
usage_policy=_usage_policy(
ground_truth_allowed=True,
training_allowed=True,
allowed_tasks=(
"building_validation",
"building_labels",
"reference_context",
),
validation_authority={"building_validation": "primary"},
),
ingest_status="configured",
known_limitations=(
"GRB building geometry is authoritative only for a governed, versioned snapshot within Flanders.",
"GRB does not independently establish imagery-time alignment or national model validation.",
),
),
_definition(
"digitaal_vlaanderen",
"Digitaal Vlaanderen (bronportaal)",
"authoritative",
"Digitaal Vlaanderen",
{"zone": "Flanders", "role": "umbrella_catalogue_and_adapter_authority"},
provider_adapter_key="digitaal_vlaanderen",
source_url="https://www.vlaanderen.be/datavindplaats",
default_crs="product_specific",
default_units="product_specific",
geographic_coverage={"zone": "Flanders"},
usage_policy=_usage_policy(
allowed_tasks=("source_catalogue", "reference_context")
),
ingest_status="configured",
known_limitations=(
"This umbrella authority is not a product-level ground-truth source.",
"A governed import must use a product-specific source key whenever one is available.",
),
),
_definition(
"digitaal_vlaanderen_buildings_addresses_register",
"Gebouwen- en adressenregister",
"authoritative",
"Digitaal Vlaanderen",
{"zone": "Flanders", "theme": "buildings_addresses"},
provider_adapter_key="buildings_addresses_register",
default_crs="EPSG:31370",
default_units="metres",
geographic_coverage={"zone": "Flanders"},
expected_geometry_types=("Point", "Polygon", "MultiPolygon"),
expected_attributes={
"required": ["id"],
"role": "administrative_corroboration",
},
usage_policy=_usage_policy(
allowed_tasks=(
"building_validation",
"address_corroboration",
"building_register_validation",
),
validation_authority={
"building_validation": "corroborative",
"building_register_validation": "primary",
},
),
known_limitations=(
"Administrative records do not replace a governed footprint-label contract.",
),
),
_definition(
"sentinel_2",
"Sentinel-2",
"contextual",
"Copernicus Programme",
{"scope": "Belgium and Belgian North Sea", "role": "multispectral_context"},
provider_adapter_key="sentinel_2",
default_crs="product_specific",
default_units="reflectance",
spatial_resolution={"metres": [10, 20, 60]},
temporal_coverage={"cadence_days": 5, "status": "product_specific"},
geographic_coverage={"scope": "Belgium and Belgian North Sea"},
expected_attributes={"required": ["product_id", "sensing_time"]},
usage_policy=_usage_policy(
training_allowed=True,
allowed_tasks=("imagery_context", "change_context"),
),
ingest_status="not_configured",
known_limitations=(
"Sentinel-2 is contextual imagery, never automatic building ground truth.",
),
),
_definition(
"digitaal_vlaanderen_dhmv",
"Digitaal Hoogtemodel Vlaanderen",
"authoritative",
"Digitaal Vlaanderen",
{"zone": "Flanders", "role": "terrain_height_corroboration"},
provider_adapter_key="dhmv",
default_crs="EPSG:31370",
default_units="m TAW",
spatial_resolution={"metres": 1},
temporal_coverage={"period": "2013-2015", "status": "product_specific"},
geographic_coverage={"zone": "Flanders"},
expected_attributes={"bands": 1, "nodata_required": True},
usage_policy=_usage_policy(
training_allowed=True,
allowed_tasks=(
"terrain_context",
"height_corroboration",
"elevation_validation",
),
validation_authority={
"building_validation": "corroborative",
"elevation_validation": "primary",
},
),
ingest_status="configured",
known_limitations=(
"DHMV is height context and cannot independently establish building labels.",
),
),
_definition(
"osm",
"OpenStreetMap",
"contextual",
"OpenStreetMap contributors",
{"scope": "community-maintained", "role": "contextual"},
provider_adapter_key="osm",
source_url="https://www.openstreetmap.org",
license_name="ODbL",
usage_restrictions="OpenStreetMap attribution and ODbL obligations apply.",
default_crs="EPSG:4326",
default_units="mixed",
geographic_coverage={"scope": "global"},
expected_geometry_types=("Point", "LineString", "Polygon", "MultiPolygon"),
expected_attributes={"status": "community_tags"},
usage_policy=_usage_policy(
allowed_tasks=("context", "candidate_discovery")
),
ingest_status="not_configured",
known_limitations=(
"OSM is never automatic ground truth for GeoIntel validation or labels.",
),
),
_definition(
"manual",
"Handmatige upload",
"experimental",
"Operator supplied",
{"scope": "operator_supplied", "trust": "unverified"},
known_limitations=(
"Manual uploads remain untrusted until a passed contract and governed provenance are attached.",
),
ingest_status="configured",
),
_definition(
"fixture",
"Test- en demo fixture",
"experimental",
"GeoIntel test fixture",
{"scope": "test_only"},
default_crs="fixture_specific",
default_units="fixture_specific",
known_limitations=(
"Fixtures must never be presented as official data or used for production training/promotion.",
),
ingest_status="configured",
),
_definition(
"map_selection",
"Afgeleide kaartselectie",
"derived",
"GeoIntel derived operation",
{"scope": "derived_from_registered_input"},
default_crs="EPSG:4326",
default_units="source_dependent",
known_limitations=(
"Derived selections inherit no authority beyond complete source snapshots and lineage edges.",
),
),
_definition(
"derived",
"Afgeleide dataset",
"derived",
"GeoIntel derived operation",
{"scope": "derived_from_registered_input"},
),
_definition(
"training_label",
"Afgeleide trainingslabels",
"derived",
"GeoIntel reviewed label pipeline",
{"scope": "derived_from_reviewed_source_snapshots"},
known_limitations=(
"Training labels require complete source lineage and human-review evidence; they inherit no automatic authority.",
),
),
_definition(
"model",
"Model artifact",
"experimental",
"GeoIntel model pipeline",
{"scope": "internal_model_artifact"},
known_limitations=(
"A model artifact is not a validated capability or promotion decision without its model card and evaluation evidence.",
),
),
_definition(
"experimental",
"Experimentele bron",
"experimental",
"Unverified",
{"scope": "unverified"},
),
_definition(
"legacy_unknown",
"Niet-geclassificeerde historische bron",
"experimental",
"Legacy import — unverified",
{"scope": "legacy", "trust": "unverified"},
ingest_status="legacy_unverified",
known_limitations=(
"Historical source identity is descriptive only until re-ingested through a governed adapter.",
),
),
_definition(
"ngi_adminvector",
"NGI AdminVector",
"authoritative",
"Nationaal Geografisch Instituut",
{"scope": "Belgium"},
default_crs="EPSG:31370",
default_units="metres",
geographic_coverage={"scope": "Belgium"},
expected_geometry_types=("Polygon", "MultiPolygon"),
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"rbins_marine_reporting_units",
"RBINS mariene rapportage-eenheden",
"authoritative",
"RBINS",
{"zone": "Belgian North Sea"},
default_crs="EPSG:4326",
default_units="degrees",
geographic_coverage={"zone": "Belgian North Sea"},
expected_geometry_types=("Polygon", "MultiPolygon"),
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"rbins_msp_2026",
"Belgisch Marien Ruimtelijk Plan 2026-2034",
"authoritative",
"RBINS",
{"zone": "Belgian North Sea", "edition": "2026-2034"},
default_crs="EPSG:4326",
default_units="degrees",
geographic_coverage={"zone": "Belgian North Sea"},
expected_geometry_types=("Polygon", "MultiPolygon"),
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"vrbg",
"Vlaams Wegenregister",
"authoritative",
"Digitaal Vlaanderen",
{"zone": "Flanders", "theme": "roads"},
default_crs="EPSG:31370",
default_units="metres",
geographic_coverage={"zone": "Flanders"},
expected_geometry_types=("LineString", "MultiLineString"),
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"digitaal_vlaanderen_orthophoto",
"Orthofoto Vlaanderen",
"contextual",
"Digitaal Vlaanderen",
{"zone": "Flanders", "role": "imagery"},
default_crs="EPSG:31370",
default_units="pixel",
spatial_resolution={"metres": 0.25},
geographic_coverage={"zone": "Flanders"},
usage_policy=_IMAGERY_CONTEXT,
ingest_status="configured",
),
_definition(
"spw_orthophoto",
"Orthofoto Wallonië",
"contextual",
"Service public de Wallonie",
{"zone": "Wallonia", "role": "imagery"},
license_name="CC BY 4.0",
default_crs="EPSG:31370",
default_units="pixel",
spatial_resolution={"metres": 0.25},
geographic_coverage={"zone": "Wallonia"},
usage_policy=_IMAGERY_CONTEXT,
ingest_status="configured",
),
_definition(
"urbis_orthophoto",
"Orthofoto Brussel",
"contextual",
"UrbIS / Brussels Region",
{"zone": "Brussels-Capital Region", "role": "imagery"},
license_name="CC0",
default_crs="EPSG:31370",
default_units="pixel",
spatial_resolution={"metres": 0.25},
geographic_coverage={"zone": "Brussels-Capital Region"},
usage_policy=_IMAGERY_CONTEXT,
ingest_status="configured",
),
_definition(
"agentschap_landbouw_zeevisserij_agricultural_parcels",
"Landbouwgebruikspercelen",
"authoritative",
"Agentschap Landbouw en Zeevisserij",
{"zone": "Flanders", "theme": "agricultural_parcels"},
default_crs="EPSG:31370",
default_units="metres",
geographic_coverage={"zone": "Flanders"},
expected_geometry_types=("Polygon", "MultiPolygon"),
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"department_omgeving_land_use",
"Landgebruik Vlaanderen",
"authoritative",
"Departement Omgeving",
{"zone": "Flanders", "theme": "land_use"},
default_crs="EPSG:31370",
default_units="metres",
geographic_coverage={"zone": "Flanders"},
expected_geometry_types=("Polygon", "MultiPolygon"),
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"inbo_bwk_natura2000",
"BWK en Natura 2000",
"authoritative",
"INBO",
{"zone": "Flanders", "theme": "nature"},
default_crs="EPSG:31370",
default_units="metres",
geographic_coverage={"zone": "Flanders"},
expected_geometry_types=("Polygon", "MultiPolygon"),
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"statbel",
"Statbel bevolking",
"authoritative",
"Statbel",
{"scope": "Belgium", "theme": "population"},
default_crs="EPSG:31370",
default_units="persons",
geographic_coverage={"scope": "Belgium"},
expected_geometry_types=("Polygon", "MultiPolygon"),
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"waterinfo",
"Waterinfo",
"authoritative",
"Waterinfo Vlaanderen",
{"zone": "Flanders", "theme": "water"},
default_crs="EPSG:31370",
default_units="source_specific",
geographic_coverage={"zone": "Flanders"},
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"department_omgeving_thematic_raster",
"Omgeving thematische rasters",
"authoritative",
"Departement Omgeving",
{"zone": "Flanders", "theme": "thematic_raster"},
default_crs="EPSG:31370",
default_units="source_specific",
geographic_coverage={"zone": "Flanders"},
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"dov_soil_map",
"DOV bodemkaart",
"authoritative",
"Databank Ondergrond Vlaanderen",
{"zone": "Flanders", "theme": "soil"},
default_crs="EPSG:31370",
default_units="metres",
geographic_coverage={"zone": "Flanders"},
expected_geometry_types=("Polygon", "MultiPolygon"),
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"vmm_flood_hazard",
"VMM overstromingskaarten",
"authoritative",
"Vlaamse Milieumaatschappij",
{"zone": "Flanders", "theme": "flood_hazard"},
default_crs="EPSG:31370",
default_units="metres",
geographic_coverage={"zone": "Flanders"},
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"vmm_vha_bathymetry_profiles",
"VHA bathymetrieprofielen",
"authoritative",
"Vlaamse Milieumaatschappij",
{"zone": "Flanders", "theme": "bathymetry_profiles"},
default_crs="EPSG:31370",
default_units="m TAW",
geographic_coverage={"zone": "Flanders"},
expected_geometry_types=("Point",),
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"historical_landuse",
"Historisch landgebruik",
"corroborative",
"Historical archive provider",
{"scope": "Belgium", "theme": "historical_land_use"},
default_crs="source_specific",
default_units="source_specific",
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"spw_geoportail",
"SPW Geoportail (bronportaal)",
"authoritative",
"Service public de Wallonie",
{"zone": "Wallonia", "role": "umbrella_catalogue_and_adapter_authority"},
provider_adapter_key="spw_geoportail",
source_url="https://geoportail.wallonie.be/catalogue",
default_crs="product_specific",
default_units="product_specific",
geographic_coverage={"zone": "Wallonia"},
usage_policy=_usage_policy(
allowed_tasks=("source_catalogue", "reference_context")
),
ingest_status="configured",
known_limitations=(
"This umbrella authority is not a product-level ground-truth source.",
"A governed import must use PICC, WALOUS, terrain, flood, orthophoto or another product-specific key when available.",
),
),
_definition(
"spw_picc",
"PICC",
"authoritative",
"Service public de Wallonie",
{"zone": "Wallonia", "theme": "topography_buildings"},
license_name="CC BY 4.0",
default_crs="EPSG:31370",
default_units="metres",
geographic_coverage={"zone": "Wallonia"},
expected_geometry_types=(
"Polygon",
"MultiPolygon",
"LineString",
"MultiLineString",
),
usage_policy=_REGIONAL_BUILDING_LABELS,
),
_definition(
"urbis",
"UrbIS",
"authoritative",
"Brussels Region",
{"zone": "Brussels-Capital Region", "theme": "topography_buildings"},
license_name="CC0",
default_crs="EPSG:31370",
default_units="metres",
geographic_coverage={"zone": "Brussels-Capital Region"},
expected_geometry_types=(
"Polygon",
"MultiPolygon",
"LineString",
"MultiLineString",
),
usage_policy=_REGIONAL_BUILDING_LABELS,
),
_definition(
"spw_walous_land_cover",
"WALOUS landbedekking",
"authoritative",
"Service public de Wallonie",
{"zone": "Wallonia", "theme": "land_cover"},
license_name="CC BY 4.0",
default_crs="EPSG:3812",
default_units="class_code",
spatial_resolution={"metres": 1},
geographic_coverage={"zone": "Wallonia"},
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"spw_bathymetry",
"SPW bathymetrie",
"authoritative",
"Service public de Wallonie",
{"zone": "Wallonia", "theme": "bathymetry"},
default_crs="EPSG:3812",
default_units="mDNG",
geographic_coverage={"zone": "Wallonia"},
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"spw_terrain",
"SPW terreinmodel",
"corroborative",
"Service public de Wallonie",
{"zone": "Wallonia", "theme": "terrain"},
default_crs="EPSG:3812",
default_units="metres",
spatial_resolution={"metres": 1},
geographic_coverage={"zone": "Wallonia"},
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"spw_flood_hazard",
"SPW overstromingsgevaar",
"authoritative",
"Service public de Wallonie",
{"zone": "Wallonia", "theme": "flood_hazard"},
default_crs="EPSG:3812",
default_units="metres",
geographic_coverage={"zone": "Wallonia"},
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"mdk_bathymetry",
"MDK bathymetrie",
"authoritative",
"Maritieme Dienstverlening en Kust",
{"zone": "Belgian North Sea", "theme": "bathymetry"},
default_crs="EPSG:3812",
default_units="metres",
geographic_coverage={"zone": "Belgian North Sea"},
usage_policy=_VECTOR_CONTEXT,
),
_definition(
"mdk_bcp_bathymetry",
"MDK BCP bathymetrie-probe en verwerving",
"authoritative",
"Maritieme Dienstverlening en Kust",
{
"zone": "Belgian North Sea",
"theme": "bathymetry",
"role": "coverage_probe_and_governed_acquisition",
},
provider_adapter_key="mdk_bcp_bathymetry",
source_url="https://www.vlaanderen.be/datavindplaats",
default_crs="EPSG:3812",
default_units="metres",
geographic_coverage={"zone": "Belgian North Sea"},
usage_policy=_VECTOR_CONTEXT,
ingest_status="not_configured",
known_limitations=(
"A BCP coverage probe is discovery evidence, not a usable bathymetry dataset.",
"Only a successful governed acquisition with an immutable response checksum may create a source snapshot.",
),
),
)
}
class SourceRegistryService:
"""Fail-closed registry operations for governed importers and validators.
All methods intentionally flush but do not commit. The caller owns the
dataset/import transaction, so a source snapshot, validation result and
dataset write can be rolled back together.
"""
@staticmethod
def normalize_source_key(source_key: str) -> str:
normalized = source_key.strip().lower()
if not normalized or len(normalized) > 120:
raise AppError(
code="SOURCE_REGISTRY_KEY_INVALID",
message="Source registry key must be a non-empty value up to 120 characters",
status_code=422,
)
return normalized
@classmethod
def definition_for(cls, source_key: str) -> SourceRegistryDefinition:
normalized = cls.normalize_source_key(source_key)
definition = SERVER_OWNED_SOURCE_DEFINITIONS.get(normalized)
if definition is None:
raise AppError(
code="SOURCE_REGISTRY_ENTRY_NOT_FOUND",
message="Source is not registered as a server-owned source",
details={"source_key": normalized},
status_code=422,
)
return definition
@classmethod
def ensure_server_owned_source(cls, db: Session, source_key: str) -> SourceRegistry:
definition = cls.definition_for(source_key)
existing = (
db.query(SourceRegistry)
.filter(SourceRegistry.source_key == definition.source_key)
.one_or_none()
)
if existing is not None:
return existing
source = SourceRegistry(**definition.as_model_values())
db.add(source)
db.flush()
return source
@staticmethod
def normalize_ingest_key(ingest_key: str) -> str:
normalized = ingest_key.strip()
if not normalized or len(normalized) > 255:
raise AppError(
code="INGEST_KEY_INVALID",
message="Ingest key must be a non-empty value up to 255 characters",
status_code=422,
)
return normalized
@classmethod
def find_dataset_by_ingest_key(
cls, db: Session, project_id: UUID, ingest_key: str
) -> Dataset | None:
normalized = cls.normalize_ingest_key(ingest_key)
return (
db.query(Dataset)
.filter(Dataset.project_id == project_id, Dataset.ingest_key == normalized)
.one_or_none()
)
@classmethod
def find_dataset_version_by_ingest_key(
cls,
db: Session,
dataset_id: UUID,
ingest_key: str,
) -> DatasetVersion | None:
normalized = cls.normalize_ingest_key(ingest_key)
return (
db.query(DatasetVersion)
.filter(
DatasetVersion.dataset_id == dataset_id,
DatasetVersion.ingest_key == normalized,
)
.one_or_none()
)
@staticmethod
def _validate_status(value: str, allowed: set[str], field_name: str) -> str:
normalized = value.strip().lower()
if normalized not in allowed:
raise AppError(
code="SOURCE_REGISTRY_STATUS_INVALID",
message=f"Unsupported {field_name}",
details={
"field": field_name,
"value": value,
"allowed": sorted(allowed),
},
status_code=422,
)
return normalized
@staticmethod
def _validate_checksum(checksum_sha256: str) -> str:
normalized = checksum_sha256.strip().lower()
if not _CHECKSUM.fullmatch(normalized):
raise AppError(
code="SOURCE_SNAPSHOT_CHECKSUM_INVALID",
message="Source snapshot checksum must be a SHA-256 hex digest",
status_code=422,
)
return normalized
@classmethod
def record_snapshot(
cls,
db: Session,
*,
source_key: str,
snapshot_key: str,
checksum_sha256: str,
source_version: str | None = None,
snapshot_at: datetime | None = None,
fetched_at: datetime | None = None,
source_url: str | None = None,
crs: str | None = None,
units: str | None = None,
spatial_resolution: dict[str, Any] | None = None,
temporal_coverage: dict[str, Any] | None = None,
geographic_coverage: dict[str, Any] | None = None,
observed_schema: dict[str, Any] | None = None,
freshness_status: str = "unknown",
ingest_status: str = "ingested",
known_limitations: list[str] | None = None,
snapshot_metadata: dict[str, Any] | None = None,
reuse_existing_snapshot: bool = False,
) -> SourceSnapshot:
"""Record immutable source evidence, or reuse an identical snapshot.
``fetched_at`` belongs to the immutable source snapshot, whereas a
Dataset's ``imported_at`` records each local ingestion event. A
governed importer may therefore replay an already-known source
snapshot for another project. In that narrow replay mode the existing
``fetched_at`` is retained; every other evidence field is still
required to be identical and the row is never updated.
"""
normalized_key = snapshot_key.strip()
if not normalized_key or len(normalized_key) > 255:
raise AppError(
code="SOURCE_SNAPSHOT_KEY_INVALID",
message="Snapshot key must be a non-empty value up to 255 characters",
status_code=422,
)
source = cls.ensure_server_owned_source(db, source_key)
checksum = cls._validate_checksum(checksum_sha256)
normalized_freshness = cls._validate_status(
freshness_status, _FRESHNESS_STATUSES, "freshness_status"
)
normalized_ingest = cls._validate_status(
ingest_status, _SNAPSHOT_INGEST_STATUSES, "ingest_status"
)
existing = (
db.query(SourceSnapshot)
.filter(
SourceSnapshot.source_registry_id == source.id,
SourceSnapshot.snapshot_key == normalized_key,
)
.one_or_none()
)
if existing is not None:
immutable_values = {
"checksum_sha256": checksum,
"source_version": source_version.strip() if source_version else None,
"snapshot_at": snapshot_at,
"source_url": source_url.strip() if source_url else None,
"crs": crs.strip() if crs else None,
"units": units.strip() if units else None,
"spatial_resolution_json": (
dict(spatial_resolution or {"status": "unknown"})
if spatial_resolution is not None
else None
),
"temporal_coverage_json": (
dict(temporal_coverage or {"status": "unknown"})
if temporal_coverage is not None
else None
),
"geographic_coverage_json": (
dict(geographic_coverage or {"status": "unknown"})
if geographic_coverage is not None
else None
),
"observed_schema_json": (
dict(observed_schema or {"status": "unknown"})
if observed_schema is not None
else None
),
"known_limitations_json": list(known_limitations)
if known_limitations is not None
else None,
"snapshot_metadata_json": dict(snapshot_metadata)
if snapshot_metadata is not None
else None,
}
if not reuse_existing_snapshot:
immutable_values["fetched_at"] = fetched_at
conflicts = {
field_name: {
"existing": getattr(existing, field_name),
"incoming": incoming,
}
for field_name, incoming in immutable_values.items()
if incoming is not None and getattr(existing, field_name) != incoming
}
if conflicts:
raise AppError(
code="SOURCE_SNAPSHOT_IMMUTABILITY_CONFLICT",
message="Existing source snapshot key has different immutable evidence",
details={
"source_key": source.source_key,
"snapshot_key": normalized_key,
"conflicting_fields": sorted(conflicts),
},
status_code=409,
)
return existing
snapshot = SourceSnapshot(
source_registry_id=source.id,
snapshot_key=normalized_key,
source_version=source_version.strip() if source_version else None,
snapshot_at=snapshot_at,
fetched_at=fetched_at,
source_url=source_url.strip() if source_url else None,
checksum_sha256=checksum,
crs=crs.strip() if crs else None,
units=units.strip() if units else None,
spatial_resolution_json=dict(spatial_resolution or {"status": "unknown"}),
temporal_coverage_json=dict(temporal_coverage or {"status": "unknown"}),
geographic_coverage_json=dict(geographic_coverage or {"status": "unknown"}),
observed_schema_json=dict(observed_schema or {"status": "unknown"}),
freshness_status=normalized_freshness,
ingest_status=normalized_ingest,
known_limitations_json=list(known_limitations or []),
snapshot_metadata_json=dict(snapshot_metadata or {}),
)
db.add(snapshot)
db.flush()
return snapshot
@classmethod
def bind_dataset_provenance(
cls,
dataset: Dataset,
*,
source: SourceRegistry,
snapshot: SourceSnapshot,
data_contract_key: str,
data_contract_version: str,
validation_status: str,
provenance_status: str,
lineage_status: str,
) -> Dataset:
cls._validate_binding(
source=source,
snapshot=snapshot,
data_contract_key=data_contract_key,
data_contract_version=data_contract_version,
validation_status=validation_status,
provenance_status=provenance_status,
lineage_status=lineage_status,
)
dataset.source_registry_id = source.id
dataset.source_snapshot_id = snapshot.id
dataset.data_contract_key = data_contract_key.strip()
dataset.data_contract_version = data_contract_version.strip()
dataset.validation_status = validation_status.strip().lower()
dataset.provenance_status = provenance_status.strip().lower()
dataset.lineage_status = lineage_status.strip().lower()
return dataset
@classmethod
def bind_dataset_version_provenance(
cls,
dataset_version: DatasetVersion,
*,
source: SourceRegistry,
snapshot: SourceSnapshot,
data_contract_key: str,
data_contract_version: str,
validation_status: str,
provenance_status: str,
lineage_status: str,
) -> DatasetVersion:
cls._validate_binding(
source=source,
snapshot=snapshot,
data_contract_key=data_contract_key,
data_contract_version=data_contract_version,
validation_status=validation_status,
provenance_status=provenance_status,
lineage_status=lineage_status,
)
dataset_version.source_registry_id = source.id
dataset_version.source_snapshot_id = snapshot.id
dataset_version.data_contract_key = data_contract_key.strip()
dataset_version.data_contract_version = data_contract_version.strip()
dataset_version.validation_status = validation_status.strip().lower()
dataset_version.provenance_status = provenance_status.strip().lower()
dataset_version.lineage_status = lineage_status.strip().lower()
return dataset_version
@classmethod
def _validate_binding(
cls,
*,
source: SourceRegistry,
snapshot: SourceSnapshot,
data_contract_key: str,
data_contract_version: str,
validation_status: str,
provenance_status: str,
lineage_status: str,
) -> None:
if source.id != snapshot.source_registry_id:
raise AppError(
code="SOURCE_SNAPSHOT_REGISTRY_MISMATCH",
message="Source snapshot does not belong to the selected source registry entry",
status_code=409,
)
if not data_contract_key.strip() or not data_contract_version.strip():
raise AppError(
code="DATA_CONTRACT_IDENTITY_REQUIRED",
message="Dataset provenance binding requires a contract key and version",
status_code=422,
)
cls._validate_status(
validation_status, _VALIDATION_STATUSES, "validation_status"
)
cls._validate_status(
provenance_status, _PROVENANCE_STATUSES, "provenance_status"
)
cls._validate_status(lineage_status, _LINEAGE_STATUSES, "lineage_status")
@staticmethod
def record_lineage_edge(
db: Session,
*,
parent_dataset_id: UUID,
child_dataset_id: UUID,
relation_type: str,
transformation_name: str,
parent_dataset_version_id: UUID | None = None,
child_dataset_version_id: UUID | None = None,
transformation_version: str | None = None,
parameters: dict[str, Any] | None = None,
input_checksum_sha256: str | None = None,
output_checksum_sha256: str | None = None,
) -> DatasetLineageEdge:
if parent_dataset_id == child_dataset_id:
raise AppError(
code="DATASET_LINEAGE_SELF_REFERENCE",
message="A dataset cannot be its own lineage parent",
status_code=422,
)
normalized_relation = relation_type.strip()
normalized_transform = transformation_name.strip()
if not normalized_relation or not normalized_transform:
raise AppError(
code="DATASET_LINEAGE_IDENTITY_REQUIRED",
message="Lineage relation type and transformation name are required",
status_code=422,
)
normalized_input_checksum = SourceRegistryService._optional_checksum(
input_checksum_sha256
)
normalized_output_checksum = SourceRegistryService._optional_checksum(
output_checksum_sha256
)
existing = (
db.query(DatasetLineageEdge)
.filter(
DatasetLineageEdge.parent_dataset_id == parent_dataset_id,
DatasetLineageEdge.child_dataset_id == child_dataset_id,
DatasetLineageEdge.relation_type == normalized_relation,
DatasetLineageEdge.transformation_name == normalized_transform,
)
.one_or_none()
)
if existing is not None:
if (
existing.input_checksum_sha256 != normalized_input_checksum
or existing.output_checksum_sha256 != normalized_output_checksum
):
raise AppError(
code="DATASET_LINEAGE_IMMUTABILITY_CONFLICT",
message="Existing lineage edge has different artifact checksums",
status_code=409,
)
return existing
if SourceRegistryService._would_create_lineage_cycle(
db,
parent_dataset_id=parent_dataset_id,
child_dataset_id=child_dataset_id,
):
raise AppError(
code="DATASET_LINEAGE_CYCLE_DETECTED",
message="The proposed lineage edge would make the dataset lineage graph cyclic",
details={
"parent_dataset_id": str(parent_dataset_id),
"child_dataset_id": str(child_dataset_id),
},
status_code=409,
)
edge = DatasetLineageEdge(
parent_dataset_id=parent_dataset_id,
child_dataset_id=child_dataset_id,
parent_dataset_version_id=parent_dataset_version_id,
child_dataset_version_id=child_dataset_version_id,
relation_type=normalized_relation,
transformation_name=normalized_transform,
transformation_version=transformation_version.strip()
if transformation_version
else None,
parameters_json=dict(parameters or {}),
input_checksum_sha256=normalized_input_checksum,
output_checksum_sha256=normalized_output_checksum,
)
db.add(edge)
db.flush()
return edge
@staticmethod
def _would_create_lineage_cycle(
db: Session,
*,
parent_dataset_id: UUID,
child_dataset_id: UUID,
) -> bool:
"""Return whether ``parent -> child`` would close an existing DAG path.
A lineage edge is directed from an input/parent Dataset to its derived
child. Adding ``parent -> child`` is unsafe precisely when ``parent``
is already reachable downstream from ``child``. Querying one indexed
parent frontier at a time avoids loading unrelated lineage history and
the visited set makes a pre-existing corrupt cycle finite to inspect.
The graph-size limit is itself fail-closed: a graph too large to audit
may not receive a new edge until it is investigated.
"""
frontier = {child_dataset_id}
visited: set[UUID] = set()
while frontier:
if parent_dataset_id in frontier:
return True
current = frontier - visited
if not current:
return False
visited.update(current)
if len(visited) > _MAX_LINEAGE_GRAPH_NODES:
raise AppError(
code="DATASET_LINEAGE_GRAPH_LIMIT_EXCEEDED",
message="Dataset lineage graph exceeds the safe traversal limit",
details={"max_nodes": _MAX_LINEAGE_GRAPH_NODES},
status_code=409,
)
edges = (
db.query(DatasetLineageEdge)
.filter(DatasetLineageEdge.parent_dataset_id.in_(current))
.all()
)
frontier = {
edge.child_dataset_id
for edge in edges
if edge.child_dataset_id not in visited
}
return False
@classmethod
def _lineage_descendant_dataset_ids(
cls,
db: Session,
*,
root_dataset_ids: set[UUID],
) -> set[UUID]:
"""Return every reachable child Dataset, including the supplied roots.
A quarantine is a lineage safety event, not merely a status update on
the immediately observed asset. Traversal follows the same directed
parent-to-child relation used by the cycle guard and remains bounded so
a corrupt graph cannot make a quarantine operation unobservable.
PostgreSQL enforces the equivalent recursive propagation for all
persisted writes; this application-side traversal keeps the service
fail-closed for normal ORM callers and focused in-memory test doubles.
"""
affected = set(root_dataset_ids)
frontier = set(root_dataset_ids)
visited: set[UUID] = set()
while frontier:
current = frontier - visited
if not current:
break
visited.update(current)
if len(visited) > _MAX_LINEAGE_GRAPH_NODES:
raise AppError(
code="DATASET_LINEAGE_GRAPH_LIMIT_EXCEEDED",
message="Dataset lineage graph exceeds the safe traversal limit during quarantine propagation",
details={"max_nodes": _MAX_LINEAGE_GRAPH_NODES},
status_code=409,
)
edges = (
db.query(DatasetLineageEdge)
.filter(DatasetLineageEdge.parent_dataset_id.in_(current))
.all()
)
frontier = {
edge.child_dataset_id
for edge in edges
if edge.child_dataset_id not in visited
}
affected.update(frontier)
return affected
@staticmethod
def _mark_dataset_quarantined(dataset: Dataset) -> None:
"""Invalidate every Dataset-level consumption gate in one place."""
dataset.status = "quarantined"
dataset.quarantine_status = "quarantined"
dataset.validation_status = "failed"
dataset.provenance_status = "incomplete"
dataset.lineage_status = "incomplete"
@staticmethod
def _mark_dataset_version_quarantined(dataset_version: DatasetVersion) -> None:
"""Invalidate a version that belongs to a quarantined Dataset lineage."""
dataset_version.validation_status = "failed"
dataset_version.provenance_status = "incomplete"
dataset_version.lineage_status = "incomplete"
@classmethod
def _propagate_lineage_quarantine(
cls,
db: Session,
*,
root_datasets: tuple[Dataset, ...],
) -> tuple[Dataset, ...]:
"""Quarantine root datasets and all immutable downstream derivatives.
A source snapshot can be bound by more than one Dataset, so callers
supply every directly affected root. Derived descendants retain their
own source snapshot evidence, but their Dataset and DatasetVersion
state becomes non-consumable until a governed re-ingest establishes a
new valid lineage.
"""
datasets_by_id: dict[UUID, Dataset] = {}
transient_roots: list[Dataset] = []
for candidate in root_datasets:
dataset_id = getattr(candidate, "id", None)
if dataset_id is None:
transient_roots.append(candidate)
else:
datasets_by_id[dataset_id] = candidate
affected_ids = (
cls._lineage_descendant_dataset_ids(
db,
root_dataset_ids=set(datasets_by_id),
)
if datasets_by_id
else set()
)
if affected_ids:
for candidate in (
db.query(Dataset).filter(Dataset.id.in_(affected_ids)).all()
):
datasets_by_id[candidate.id] = candidate
affected_datasets = tuple((*datasets_by_id.values(), *transient_roots))
for candidate in affected_datasets:
cls._mark_dataset_quarantined(candidate)
if affected_ids:
for dataset_version in (
db.query(DatasetVersion)
.filter(DatasetVersion.dataset_id.in_(affected_ids))
.all()
):
cls._mark_dataset_version_quarantined(dataset_version)
return affected_datasets
@classmethod
def quarantine_dataset(
cls,
db: Session,
*,
stage: str,
reason_code: str,
dataset: Dataset | None = None,
dataset_version: DatasetVersion | None = None,
source_snapshot: SourceSnapshot | None = None,
details: dict[str, Any] | None = None,
artifact_path: str | None = None,
artifact_checksum_sha256: str | None = None,
) -> DatasetQuarantine:
normalized_stage = stage.strip()
normalized_reason = reason_code.strip()
if not normalized_stage or not normalized_reason:
raise AppError(
code="QUARANTINE_REASON_REQUIRED",
message="Quarantine stage and reason code are required",
status_code=422,
)
if dataset is None and dataset_version is None and source_snapshot is None:
raise AppError(
code="QUARANTINE_TARGET_REQUIRED",
message="Quarantine requires a dataset, dataset version or source snapshot",
status_code=422,
)
if dataset is None and dataset_version is not None:
dataset = cls._dataset_for_version(db, dataset_version)
if dataset is None:
# A version-only quarantine without the owning Dataset would
# leave the Dataset consumable: every current consumption
# boundary evaluates Dataset, not DatasetVersion. Refuse the
# partial state rather than quietly recording an ineffective
# quarantine.
raise AppError(
code="QUARANTINE_PARENT_DATASET_NOT_FOUND",
message="A version quarantine requires its owning Dataset so the quarantine can propagate.",
details={
"dataset_version_id": str(dataset_version.id),
"dataset_id": str(dataset_version.dataset_id),
},
status_code=409,
)
root_datasets: list[Dataset] = [dataset] if dataset is not None else []
if source_snapshot is not None:
# A snapshot is immutable shared evidence. Its direct bindings are
# roots too, so invalidating one artifact cannot leave a sibling or
# any downstream derivative consumable through a stale lineage.
root_datasets.extend(
db.query(Dataset)
.filter(Dataset.source_snapshot_id == source_snapshot.id)
.all()
)
cls._propagate_lineage_quarantine(db, root_datasets=tuple(root_datasets))
if dataset_version is not None:
cls._mark_dataset_version_quarantined(dataset_version)
if source_snapshot is not None:
# A quarantined source artifact must no longer satisfy the
# authoritative-validation eligibility check for any linked data.
source_snapshot.ingest_status = "quarantined"
record = DatasetQuarantine(
dataset_id=dataset.id if dataset is not None else None,
dataset_version_id=dataset_version.id
if dataset_version is not None
else None,
source_snapshot_id=source_snapshot.id
if source_snapshot is not None
else None,
stage=normalized_stage,
reason_code=normalized_reason,
details_json=dict(details or {}),
artifact_path=artifact_path,
artifact_checksum_sha256=cls._optional_checksum(artifact_checksum_sha256),
status="quarantined",
)
db.add(record)
db.flush()
return record
@staticmethod
def _dataset_for_version(
db: Session, dataset_version: DatasetVersion
) -> Dataset | None:
"""Resolve the Dataset that must share a version's quarantine state."""
related = getattr(dataset_version, "dataset", None)
if related is not None:
return related
getter = getattr(db, "get", None)
if callable(getter):
resolved = getter(Dataset, dataset_version.dataset_id)
if resolved is not None:
return resolved
# Use a normal ORM query as a final path for sessions where the
# relationship is deliberately not loaded. This also keeps bounded
# in-memory persistence fixtures representative of production.
return (
db.query(Dataset)
.filter(Dataset.id == dataset_version.dataset_id)
.one_or_none()
)
@staticmethod
def validation_authority_for_task(source: SourceRegistry, task: str) -> str | None:
policy = (
source.usage_policy_json
if isinstance(source.usage_policy_json, dict)
else {}
)
authority = policy.get("validation_authority")
if not isinstance(authority, dict):
return None
value = authority.get(task)
return str(value) if value else None
@classmethod
def is_dataset_eligible_for_authoritative_validation(
cls,
dataset: Dataset,
*,
source: SourceRegistry,
snapshot: SourceSnapshot,
task: str,
) -> bool:
policy = (
source.usage_policy_json
if isinstance(source.usage_policy_json, dict)
else {}
)
return bool(
source.classification == "authoritative"
and policy.get("ground_truth_allowed") is True
and cls.validation_authority_for_task(source, task) == "primary"
and dataset.status == "ready"
and dataset.quarantine_status == "not_quarantined"
and dataset.validation_status == "passed"
and dataset.provenance_status == "complete"
and dataset.lineage_status in {"complete", "not_applicable"}
and dataset.source_registry_id == source.id
and dataset.source_snapshot_id == snapshot.id
and snapshot.source_registry_id == source.id
and snapshot.ingest_status == "ingested"
and bool(
snapshot.checksum_sha256
and _CHECKSUM.fullmatch(snapshot.checksum_sha256)
)
)
@classmethod
def _optional_checksum(cls, checksum_sha256: str | None) -> str | None:
if checksum_sha256 is None:
return None
return cls._validate_checksum(checksum_sha256)