Add official source edition probes
GeoIntel CI / docs-smoke (push) Canceled after 0s
GeoIntel CI / contract-smoke (push) Canceled after 0s

This commit is contained in:
Codex
2026-07-16 17:33:40 +02:00
parent fd05c46e2a
commit 595967f892
27 changed files with 1513 additions and 18 deletions
@@ -0,0 +1,545 @@
from __future__ import annotations
from dataclasses import dataclass, replace
from datetime import datetime, timedelta, timezone
from email.utils import parsedate_to_datetime
from hashlib import sha256
import re
from threading import Lock
from typing import Any, Callable
from urllib.error import HTTPError, URLError
from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit
from urllib.request import Request, urlopen
from uuid import UUID
from xml.etree import ElementTree
from sqlalchemy.orm import Session
from app.core.config import Settings, get_settings
from app.core.errors import AppError
from app.models import Dataset, Project
from app.schemas.source_catalog import (
SourceCatalogProbeItem,
SourceCatalogProbeReport,
SourceCatalogProbeSummary,
)
_GMD = "http://www.isotc211.org/2005/gmd"
_GCO = "http://www.isotc211.org/2005/gco"
_WFS = "http://www.opengis.net/wfs/2.0"
_XLINK = "http://www.w3.org/1999/xlink"
_METADATA_HOST = "metadata.vlaanderen.be"
_VERSION_DATE = re.compile(r"^(?:toestand\s+)?(\d{4}-\d{2}-\d{2})$", re.IGNORECASE)
_ORTHOPHOTO_EDITION = re.compile(r"^\d{4}\.\d{2}$")
class CatalogProbeFailure(RuntimeError):
def __init__(self, code: str, message: str) -> None:
super().__init__(message)
self.code = code
self.message = message
@dataclass(frozen=True)
class ProbeContract:
source_name: str
display_name: str
service_type: str
endpoint_url: str
expected_layers: tuple[str, ...]
@dataclass(frozen=True)
class FetchResult:
content: bytes
content_type: str
etag: str | None
last_modified_at: datetime | None
final_url: str
@dataclass(frozen=True)
class RemoteProbe:
status: str
reachable: bool
checked_at: datetime
expected_layers: tuple[str, ...]
matched_layers: tuple[str, ...] = ()
missing_layers: tuple[str, ...] = ()
advertised_layer_count: int = 0
metadata_url: str | None = None
metadata_identifier: str | None = None
remote_title: str | None = None
remote_version: str | None = None
remote_modified_at: datetime | None = None
remote_published_at: datetime | None = None
capabilities_sha256: str | None = None
capabilities_etag: str | None = None
capabilities_last_modified_at: datetime | None = None
message: str = ""
error_code: str | None = None
cached: bool = False
@dataclass(frozen=True)
class CacheEntry:
expires_at: datetime
probe: RemoteProbe
_REMOTE_CACHE: dict[str, CacheEntry] = {}
_CACHE_LOCK = Lock()
def _utc(value: datetime | None = None) -> datetime:
current = value or datetime.now(timezone.utc)
if current.tzinfo is None:
return current.replace(tzinfo=timezone.utc)
return current.astimezone(timezone.utc)
def _header(headers: Any, name: str) -> str | None:
value = headers.get(name) if headers is not None else None
return str(value).strip() if value is not None and str(value).strip() else None
def _http_date(value: str | None) -> datetime | None:
if not value:
return None
try:
return _utc(parsedate_to_datetime(value))
except (TypeError, ValueError, OverflowError):
return None
def _with_capabilities_query(base_url: str, service_type: str) -> str:
parsed = urlsplit(base_url)
retained = [(key, value) for key, value in parse_qsl(parsed.query) if key.lower() not in {"service", "request", "version"}]
version = "2.0.0" if service_type == "WFS" else "1.3.0"
query = urlencode([*retained, ("SERVICE", service_type), ("VERSION", version), ("REQUEST", "GetCapabilities")])
return urlunsplit((parsed.scheme, parsed.netloc, parsed.path, query, ""))
def _bounded_fetch(url: str, settings: Settings, opener: Callable[..., Any] | None = None) -> FetchResult:
request = Request(
url,
headers={"Accept": "application/xml,text/xml,application/json", "User-Agent": "GeoIntel/0.1 source-catalog-probe"},
)
max_bytes = settings.source_catalog_probe_max_response_mb * 1024 * 1024
try:
with (opener or urlopen)(request, timeout=settings.source_catalog_probe_timeout_seconds) as response:
content_length = _header(response.headers, "Content-Length")
if content_length:
try:
if int(content_length) > max_bytes:
raise CatalogProbeFailure("CATALOG_RESPONSE_TOO_LARGE", "De officiële metadatarespons overschrijdt de ingestelde limiet.")
except ValueError as exc:
raise CatalogProbeFailure("CATALOG_INVALID_RESPONSE", "De officiële metadatarespons bevat een ongeldige Content-Length.") from exc
content = response.read(max_bytes + 1)
result = FetchResult(
content=content,
content_type=_header(response.headers, "Content-Type") or "",
etag=_header(response.headers, "ETag"),
last_modified_at=_http_date(_header(response.headers, "Last-Modified")),
final_url=str(response.geturl()) if hasattr(response, "geturl") else url,
)
except CatalogProbeFailure:
raise
except (HTTPError, URLError, TimeoutError, OSError) as exc:
raise CatalogProbeFailure("CATALOG_PROVIDER_UNAVAILABLE", "De officiële catalogus kon niet tijdig worden gelezen.") from exc
if len(result.content) > max_bytes:
raise CatalogProbeFailure("CATALOG_RESPONSE_TOO_LARGE", "De officiële metadatarespons overschrijdt de ingestelde limiet.")
return result
def _local_name(tag: str) -> str:
return tag.rsplit("}", 1)[-1]
def _direct_text(parent: ElementTree.Element, name: str) -> str | None:
for child in parent:
if _local_name(child.tag) == name and child.text and child.text.strip():
return child.text.strip()
return None
def _metadata_urls(parent: ElementTree.Element) -> list[str]:
urls: list[str] = []
for node in parent.iter():
if _local_name(node.tag) not in {"MetadataURL", "MetadataUrl"}:
continue
href = node.attrib.get(f"{{{_XLINK}}}href")
if href:
urls.append(href.strip())
for child in node.iter():
if _local_name(child.tag) in {"URL", "OnlineResource"}:
nested = child.attrib.get(f"{{{_XLINK}}}href") or (child.text or "").strip()
if nested:
urls.append(nested)
return urls
def _parse_capabilities(content: bytes, contract: ProbeContract) -> tuple[list[str], str | None]:
try:
root = ElementTree.fromstring(content)
except ElementTree.ParseError as exc:
raise CatalogProbeFailure("CATALOG_INVALID_XML", "De capabilities-respons is geen geldige XML.") from exc
layers: list[str] = []
metadata_urls: list[str] = []
node_name = "FeatureType" if contract.service_type == "WFS" else "Layer"
for node in root.iter():
if _local_name(node.tag) != node_name:
continue
raw_name = _direct_text(node, "Name")
if not raw_name:
continue
name = raw_name.rsplit(":", 1)[-1]
layers.append(name)
if name in contract.expected_layers:
metadata_urls.extend(_metadata_urls(node))
expected = set(contract.expected_layers)
if not expected.intersection(layers):
raise CatalogProbeFailure("CATALOG_EXPECTED_LAYERS_MISSING", "De officiële service bevat geen van de verwachte lagen.")
xml_metadata = next((url for url in metadata_urls if "GetRecordById" in url and "OUTPUTSCHEMA" in url.upper()), None)
return sorted(set(layers)), xml_metadata
def _validate_metadata_url(url: str) -> str:
parsed = urlsplit(url)
if parsed.scheme != "https" or parsed.hostname != _METADATA_HOST or parsed.username or parsed.password:
raise CatalogProbeFailure("CATALOG_METADATA_URL_REJECTED", "De capabilities verwijzen niet naar de toegestane officiële metadatahost.")
if not parsed.path.startswith("/srv/dut/csw"):
raise CatalogProbeFailure("CATALOG_METADATA_URL_REJECTED", "De capabilities verwijzen niet naar het toegestane CSW-pad.")
query = {key.lower(): value for key, value in parse_qsl(parsed.query)}
if query.get("request", "").lower() != "getrecordbyid" or not query.get("id"):
raise CatalogProbeFailure("CATALOG_METADATA_URL_REJECTED", "De capabilities bevatten geen begrensde GetRecordById-verwijzing.")
return url
def _validate_capabilities_url(url: str) -> str:
parsed = urlsplit(url)
if parsed.scheme not in {"http", "https"} or not parsed.hostname or parsed.username or parsed.password:
raise CatalogProbeFailure("CATALOG_ENDPOINT_REJECTED", "De ingestelde capabilities-URL moet een geldige HTTP(S)-URL zonder credentials zijn.")
return url
def _node_text(node: ElementTree.Element | None) -> str | None:
if node is None:
return None
for descendant in node.iter():
if descendant is not node and descendant.text and descendant.text.strip():
return descendant.text.strip()
return node.text.strip() if node.text and node.text.strip() else None
def _parse_iso_datetime(value: str | None) -> datetime | None:
if not value:
return None
normalized = value.strip().replace("Z", "+00:00")
try:
return _utc(datetime.fromisoformat(normalized))
except ValueError:
try:
return datetime.strptime(normalized, "%Y-%m-%d").replace(tzinfo=timezone.utc)
except ValueError:
return None
def _parse_metadata(content: bytes) -> dict[str, Any]:
try:
root = ElementTree.fromstring(content)
except ElementTree.ParseError as exc:
raise CatalogProbeFailure("CATALOG_INVALID_METADATA_XML", "Het officiële metadatarecord is geen geldige XML.") from exc
namespaces = {"gmd": _GMD, "gco": _GCO}
metadata = root.find(".//gmd:MD_Metadata", namespaces)
if metadata is None and _local_name(root.tag) == "MD_Metadata":
metadata = root
if metadata is None:
raise CatalogProbeFailure("CATALOG_METADATA_MISSING", "Het CSW-antwoord bevat geen ISO 19139 metadatarecord.")
citation = metadata.find(".//gmd:identificationInfo/*/gmd:citation/gmd:CI_Citation", namespaces)
title = _node_text(citation.find("gmd:title", namespaces) if citation is not None else None)
edition = _node_text(citation.find("gmd:edition", namespaces) if citation is not None else None)
identifier = _node_text(metadata.find("gmd:fileIdentifier", namespaces))
modified = _parse_iso_datetime(_node_text(metadata.find("gmd:dateStamp", namespaces)))
published = None
if citation is not None:
for date_node in citation.findall("gmd:date/gmd:CI_Date", namespaces):
date_type = date_node.find("gmd:dateType/gmd:CI_DateTypeCode", namespaces)
if date_type is not None and date_type.attrib.get("codeListValue") == "publication":
published = _parse_iso_datetime(_node_text(date_node.find("gmd:date", namespaces)))
break
if not title or not edition:
raise CatalogProbeFailure("CATALOG_VERSION_MISSING", "Het officiële metadatarecord bevat geen herkenbare titel en editie.")
return {
"identifier": identifier,
"title": title,
"version": edition,
"modified_at": modified,
"published_at": published,
}
def _probe_remote(
contract: ProbeContract,
settings: Settings,
*,
opener: Callable[..., Any] | None,
now: datetime,
) -> RemoteProbe:
try:
capabilities = _bounded_fetch(_validate_capabilities_url(contract.endpoint_url), settings, opener)
_validate_capabilities_url(capabilities.final_url)
content_type = capabilities.content_type.lower()
if content_type and "xml" not in content_type and "text" not in content_type:
raise CatalogProbeFailure("CATALOG_INVALID_CONTENT_TYPE", "De officiële capabilities-respons is geen XML.")
layers, metadata_url = _parse_capabilities(capabilities.content, contract)
matched = tuple(layer for layer in contract.expected_layers if layer in layers)
missing = tuple(layer for layer in contract.expected_layers if layer not in layers)
digest = sha256(capabilities.content).hexdigest()
if not metadata_url:
return RemoteProbe(
status="degraded",
reachable=True,
checked_at=now,
expected_layers=contract.expected_layers,
matched_layers=matched,
missing_layers=missing,
advertised_layer_count=len(layers),
capabilities_sha256=digest,
capabilities_etag=capabilities.etag,
capabilities_last_modified_at=capabilities.last_modified_at,
message="De service is bereikbaar, maar publiceert geen machineleesbare ISO-metadata voor de verwachte lagen.",
error_code="CATALOG_METADATA_LINK_MISSING",
)
metadata_url = _validate_metadata_url(metadata_url)
metadata_response = _bounded_fetch(metadata_url, settings, opener)
_validate_metadata_url(metadata_response.final_url)
metadata = _parse_metadata(metadata_response.content)
status = "degraded" if missing else "available"
message = (
f"De officiële catalogus is bereikbaar en publiceert editie {metadata['version']}."
if not missing
else f"Editie {metadata['version']} is gevonden, maar niet alle verwachte lagen worden aangeboden."
)
return RemoteProbe(
status=status,
reachable=True,
checked_at=now,
expected_layers=contract.expected_layers,
matched_layers=matched,
missing_layers=missing,
advertised_layer_count=len(layers),
metadata_url=metadata_url,
metadata_identifier=metadata["identifier"],
remote_title=metadata["title"],
remote_version=metadata["version"],
remote_modified_at=metadata["modified_at"],
remote_published_at=metadata["published_at"],
capabilities_sha256=digest,
capabilities_etag=capabilities.etag,
capabilities_last_modified_at=capabilities.last_modified_at,
message=message,
error_code="CATALOG_EXPECTED_LAYERS_INCOMPLETE" if missing else None,
)
except CatalogProbeFailure as exc:
return RemoteProbe(
status="unavailable",
reachable=False,
checked_at=now,
expected_layers=contract.expected_layers,
missing_layers=contract.expected_layers,
message=exc.message,
error_code=exc.code,
)
def _remote_with_cache(
contract: ProbeContract,
settings: Settings,
*,
force: bool,
opener: Callable[..., Any] | None,
now: datetime,
) -> RemoteProbe:
cache_key = f"{contract.source_name}|{contract.endpoint_url}|{','.join(contract.expected_layers)}"
if not force and settings.source_catalog_probe_cache_ttl_seconds > 0:
with _CACHE_LOCK:
entry = _REMOTE_CACHE.get(cache_key)
if entry and entry.expires_at > now:
return replace(entry.probe, cached=True)
probe = _probe_remote(contract, settings, opener=opener, now=now)
if settings.source_catalog_probe_cache_ttl_seconds > 0:
with _CACHE_LOCK:
_REMOTE_CACHE[cache_key] = CacheEntry(
expires_at=now + timedelta(seconds=settings.source_catalog_probe_cache_ttl_seconds),
probe=probe,
)
return probe
def _dataset_source_name(dataset: Dataset) -> str:
return (dataset.source_name or dataset.source or "").strip().lower()
def _latest_local_version(source_name: str, datasets: list[Dataset]) -> str | None:
candidates = [item for item in datasets if _dataset_source_name(item) == source_name and item.source_version]
if source_name == "digitaal_vlaanderen_orthophoto":
explicit_current = [
item
for item in candidates
if any(token in (item.source_version or "").lower() for token in ("most_recent", "latest", "current"))
]
if explicit_current:
candidates = explicit_current
if not candidates:
return None
latest = max(
candidates,
key=lambda item: (
_utc(item.imported_at) if item.imported_at else datetime.min.replace(tzinfo=timezone.utc),
_utc(item.observed_at) if item.observed_at else datetime.min.replace(tzinfo=timezone.utc),
str(item.id),
),
)
return latest.source_version
def _normalized_version(source_name: str, version: str | None) -> str | None:
if not version:
return None
value = version.strip()
if source_name == "grb":
match = _VERSION_DATE.fullmatch(value)
return match.group(1) if match else None
if source_name == "digitaal_vlaanderen_orthophoto" and _ORTHOPHOTO_EDITION.fullmatch(value):
return value
return None
def _comparison(source_name: str, local: str | None, remote: str | None, remote_status: str) -> str:
if remote_status in {"unavailable", "disabled"}:
return "unavailable"
if not local:
return "no_local_data"
normalized_local = _normalized_version(source_name, local)
normalized_remote = _normalized_version(source_name, remote)
if normalized_local is None or normalized_remote is None:
return "not_comparable"
return "same" if normalized_local == normalized_remote else "different"
class SourceCatalogProbeService:
@staticmethod
def clear_cache() -> None:
with _CACHE_LOCK:
_REMOTE_CACHE.clear()
@staticmethod
def audit_project(
db: Session,
project_id: UUID,
*,
force: bool = False,
opener: Callable[..., Any] | None = None,
now: datetime | None = None,
settings: Settings | None = None,
) -> SourceCatalogProbeReport:
if not db.get(Project, project_id):
raise AppError(code="PROJECT_NOT_FOUND", message="Project not found", status_code=404)
active_settings = settings or get_settings()
generated_at = _utc(now)
datasets = db.query(Dataset).filter(Dataset.project_id == project_id).all()
contracts = (
ProbeContract(
source_name="grb",
display_name="Basiskaart Vlaanderen (GRB)",
service_type="WFS",
endpoint_url=_with_capabilities_query(active_settings.source_catalog_grb_wfs_url, "WFS"),
expected_layers=("GBG", "WBN", "WGO", "ADP"),
),
ProbeContract(
source_name="digitaal_vlaanderen_orthophoto",
display_name="Orthofoto Vlaanderen",
service_type="WMS",
endpoint_url=_with_capabilities_query(active_settings.orthophoto_wms_url, "WMS"),
expected_layers=(active_settings.orthophoto_wms_layer, "Vliegdagcontour"),
),
)
items: list[SourceCatalogProbeItem] = []
for contract in contracts:
local_version = _latest_local_version(contract.source_name, datasets)
if active_settings.source_catalog_probe_enabled:
remote = _remote_with_cache(
contract,
active_settings,
force=force,
opener=opener,
now=generated_at,
)
else:
remote = RemoteProbe(
status="disabled",
reachable=False,
checked_at=generated_at,
expected_layers=contract.expected_layers,
missing_layers=contract.expected_layers,
message="Officiële catalogusprobes zijn uitgeschakeld in de runtimeconfiguratie.",
error_code="CATALOG_PROBE_DISABLED",
)
comparison = _comparison(contract.source_name, local_version, remote.remote_version, remote.status)
message = remote.message
if comparison == "different":
message += " De officiële editie verschilt van de lokaal vastgelegde bronversie; controleer dit handmatig vóór een begrensde verversing."
elif comparison == "not_comparable" and local_version:
message += " De lokale waarde is een opname- of importmarkering en kan niet eerlijk als officiële cataloguseditie worden vergeleken."
items.append(
SourceCatalogProbeItem(
source_name=contract.source_name,
display_name=contract.display_name,
service_type=contract.service_type,
endpoint_url=contract.endpoint_url,
status=remote.status,
reachable=remote.reachable,
checked_at=remote.checked_at,
cached=remote.cached,
expected_layers=list(remote.expected_layers),
matched_layers=list(remote.matched_layers),
missing_layers=list(remote.missing_layers),
advertised_layer_count=remote.advertised_layer_count,
metadata_url=remote.metadata_url,
metadata_identifier=remote.metadata_identifier,
remote_title=remote.remote_title,
remote_version=remote.remote_version,
remote_modified_at=remote.remote_modified_at,
remote_published_at=remote.remote_published_at,
local_source_version=local_version,
comparison_status=comparison,
capabilities_sha256=remote.capabilities_sha256,
capabilities_etag=remote.capabilities_etag,
capabilities_last_modified_at=remote.capabilities_last_modified_at,
message=message,
error_code=remote.error_code,
)
)
summary = SourceCatalogProbeSummary(
provider_count=len(items),
available_count=sum(item.status == "available" for item in items),
degraded_count=sum(item.status == "degraded" for item in items),
unavailable_count=sum(item.status == "unavailable" for item in items),
disabled_count=sum(item.status == "disabled" for item in items),
different_version_count=sum(item.comparison_status == "different" for item in items),
)
return SourceCatalogProbeReport(
project_id=project_id,
generated_at=generated_at,
summary=summary,
items=items,
limitations=[
"Deze expliciete controle leest alleen allowlisted WFS/WMS-capabilities en gekoppelde ISO 19139 metadata.",
"Er worden geen features, rasters of modelbestanden opgehaald en geen datasets aangemaakt of overschreven.",
"Een versieverschil is controlesignaal, geen bewijs dat een lokale dataset onbruikbaar is en geen automatische importopdracht.",
],
)