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 from html.parser import HTMLParser 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.services.outbound_request_guard import guarded_opener from app.models import Dataset, Project from app.schemas.source_catalog import ( SourceCatalogProbeItem, SourceCatalogProbeReport, SourceCatalogProbeSummary, ) from app.services.statbel_catalog_probe import ( StatbelCatalogError, parse_statbel_population_catalog, validate_statbel_catalog_url, ) _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"^(20\d{2})\.(\d{2})$") _ALZ_SOURCE_NAME = "agentschap_landbouw_zeevisserij_agricultural_parcels" _ALZ_RELEASE_HOST = "landbouwcijfers.vlaanderen.be" _ALZ_RELEASE_PATH = "/open-geodata-landbouwgebruikspercelen" _ALZ_DOWNLOAD_HOST = "www.landbouwvlaanderen.be" _ALZ_DOWNLOAD_PATH = re.compile(r"^/bestanden/gis/agpa_(20\d{2})_(\d{4}-\d{2}-\d{2})_public\.zip$") _ALZ_SNAPSHOT = re.compile( r"^Landbouwgebruikspercelen\s+(20\d{2})\s*-\s*(\d+)e\s+snapshot\s*" r"\(extractie\s+(\d{2}-\d{2}-\d{4})\)(?:\s*-\s*GPKG)?$", re.IGNORECASE, ) _ALZ_EDITION = re.compile(r"^(20\d{2})-(?:definitive|v3)$", re.IGNORECASE) _YEAR_EDITION = re.compile(r"^20\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 class _AnchorParser(HTMLParser): def __init__(self) -> None: super().__init__(convert_charrefs=True) self.anchors: list[tuple[str, str]] = [] self._href: str | None = None self._text: list[str] = [] def handle_starttag(self, tag: str, attrs: list[tuple[str, str | None]]) -> None: if tag.lower() != "a" or self._href is not None: return href = dict(attrs).get("href") if href: self._href = href.strip() self._text = [] def handle_data(self, data: str) -> None: if self._href is not None: self._text.append(data) def handle_endtag(self, tag: str) -> None: if tag.lower() == "a" and self._href is not None: self.anchors.append((self._href, "".join(self._text))) self._href = None self._text = [] _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, *, max_response_mb: int | None = None, accept: str = "application/xml,text/xml,text/html,application/json", ) -> FetchResult: request = Request( url, headers={"Accept": accept, "User-Agent": "GeoIntel/0.1 source-catalog-probe"}, ) max_bytes = (max_response_mb or settings.source_catalog_probe_max_response_mb) * 1024 * 1024 try: with (opener or guarded_opener(url))(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 _validate_alz_release_url(url: str) -> str: parsed = urlsplit(url) if ( parsed.scheme != "https" or parsed.hostname != _ALZ_RELEASE_HOST or parsed.port not in {None, 443} or parsed.username or parsed.password or parsed.path.rstrip("/") != _ALZ_RELEASE_PATH or parsed.query or parsed.fragment ): raise CatalogProbeFailure( "CATALOG_ALZ_RELEASE_URL_REJECTED", "De ingestelde ALZ-publicatiepagina valt buiten de toegestane officiële URL.", ) return url def _validate_alz_download_url(url: str) -> tuple[int, datetime]: parsed = urlsplit(url) match = _ALZ_DOWNLOAD_PATH.fullmatch(parsed.path) if ( parsed.scheme != "https" or parsed.hostname != _ALZ_DOWNLOAD_HOST or parsed.port not in {None, 443} or parsed.username or parsed.password or parsed.query or parsed.fragment or not match ): raise CatalogProbeFailure( "CATALOG_ALZ_DOWNLOAD_URL_REJECTED", "De ALZ-publicatiepagina bevat een datasetlink buiten de toegestane officiële URL-structuur.", ) try: published_at = datetime.strptime(match.group(2), "%Y-%m-%d").replace(tzinfo=timezone.utc) except ValueError as exc: raise CatalogProbeFailure( "CATALOG_ALZ_DOWNLOAD_URL_REJECTED", "De ALZ-datasetlink bevat geen geldige publicatiedatum.", ) from exc return int(match.group(1)), published_at def _normalized_html_text(value: str) -> str: return re.sub(r"\s+", " ", value.replace("\xad", "").replace("–", "-").replace("—", "-")).strip() def _parse_alz_release_page(content: bytes) -> dict[str, Any]: try: html = content.decode("utf-8") except UnicodeDecodeError as exc: raise CatalogProbeFailure( "CATALOG_ALZ_INVALID_HTML", "De officiële ALZ-publicatiepagina is niet geldige UTF-8 HTML.", ) from exc parser = _AnchorParser() try: parser.feed(html) parser.close() except Exception as exc: raise CatalogProbeFailure( "CATALOG_ALZ_INVALID_HTML", "De officiële ALZ-publicatiepagina kon niet veilig worden ontleed.", ) from exc definitive: list[tuple[int, datetime]] = [] snapshots: list[tuple[int, int, datetime]] = [] for href, raw_text in parser.anchors: text = _normalized_html_text(raw_text) looks_like_alz_release = ( text.casefold() == "downloaden" or text.casefold().startswith("landbouwgebruikspercelen ") or "agpa_" in href.casefold() ) if not looks_like_alz_release: continue year, file_date = _validate_alz_download_url(href) snapshot = _ALZ_SNAPSHOT.fullmatch(text) if snapshot: snapshot_year = int(snapshot.group(1)) snapshot_number = int(snapshot.group(2)) try: extraction_date = datetime.strptime(snapshot.group(3), "%d-%m-%Y").replace(tzinfo=timezone.utc) except ValueError as exc: raise CatalogProbeFailure( "CATALOG_ALZ_SNAPSHOT_INVALID", "De actuele ALZ-snapshot bevat geen geldige extractiedatum.", ) from exc if snapshot_year != year or extraction_date != file_date or snapshot_number not in {1, 2, 3}: raise CatalogProbeFailure( "CATALOG_ALZ_SNAPSHOT_INVALID", "De actuele ALZ-snapshot is niet consistent met de officiële datasetlink.", ) snapshots.append((year, snapshot_number, extraction_date)) elif text.casefold() == "downloaden": definitive.append((year, file_date)) else: raise CatalogProbeFailure( "CATALOG_ALZ_RELEASE_UNRECOGNIZED", "De ALZ-publicatiepagina bevat een niet-herkende landbouwdatasetpublicatie.", ) if not definitive: raise CatalogProbeFailure( "CATALOG_ALZ_DEFINITIVE_MISSING", "De officiële ALZ-publicatiepagina bevat geen herkenbare definitieve landbouwperceeleditie.", ) latest_definitive = max(definitive, key=lambda item: (item[0], item[1])) latest_snapshot = max(snapshots, key=lambda item: (item[0], item[1], item[2])) if snapshots else None if latest_snapshot and latest_snapshot[1] == 3 and latest_snapshot[0] >= latest_definitive[0]: latest_definitive = (latest_snapshot[0], latest_snapshot[2]) return { "definitive": latest_definitive, "definitive_count": len(definitive), "snapshot": latest_snapshot, } 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_alz_remote( contract: ProbeContract, settings: Settings, *, opener: Callable[..., Any] | None, now: datetime, ) -> RemoteProbe: release_url = _validate_alz_release_url(contract.endpoint_url) response = _bounded_fetch(release_url, settings, opener) _validate_alz_release_url(response.final_url) content_type = response.content_type.lower() if content_type and "html" not in content_type and "text" not in content_type: raise CatalogProbeFailure( "CATALOG_ALZ_INVALID_CONTENT_TYPE", "De officiële ALZ-publicatiepagina is geen HTML-respons.", ) release = _parse_alz_release_page(response.content) definitive_year, definitive_date = release["definitive"] snapshot = release["snapshot"] matched = ["definitive_archive"] missing: list[str] = [] if snapshot: matched.append("current_snapshot") else: missing.append("current_snapshot") definitive_version = f"{definitive_year}-v3" if snapshot: snapshot_year, snapshot_number, snapshot_date = snapshot snapshot_version = f"{snapshot_year}-v{snapshot_number}" if snapshot_number < 3: message = ( f"De officiële ALZ-publicatiepagina bevestigt definitieve editie {definitive_version}. " f"De actuele publicatie {snapshot_version} van {snapshot_date.date().isoformat()} is voorlopig " "en wordt niet als historische vervanging aangemerkt." ) else: message = f"De officiële ALZ-publicatiepagina bevestigt definitieve editie {definitive_version}." remote_title = f"Landbouwgebruikspercelen {definitive_version}; actuele publicatie {snapshot_version}" else: message = ( f"De officiële ALZ-publicatiepagina bevestigt definitieve editie {definitive_version}, " "maar bevat geen herkenbare actuele snapshot." ) remote_title = f"Landbouwgebruikspercelen {definitive_version}" return RemoteProbe( status="available" if snapshot else "degraded", reachable=True, checked_at=now, expected_layers=contract.expected_layers, matched_layers=tuple(matched), missing_layers=tuple(missing), advertised_layer_count=release["definitive_count"] + (1 if snapshot else 0), metadata_url=response.final_url, metadata_identifier="alz-agricultural-use-parcels", remote_title=remote_title, remote_version=definitive_version, remote_modified_at=response.last_modified_at, remote_published_at=definitive_date, capabilities_sha256=sha256(response.content).hexdigest(), capabilities_etag=response.etag, capabilities_last_modified_at=response.last_modified_at, message=message, error_code="CATALOG_ALZ_CURRENT_SNAPSHOT_MISSING" if not snapshot else None, ) def _probe_statbel_remote( contract: ProbeContract, settings: Settings, *, opener: Callable[..., Any] | None, now: datetime, ) -> RemoteProbe: try: catalog_url = validate_statbel_catalog_url(contract.endpoint_url) response = _bounded_fetch( catalog_url, settings, opener, max_response_mb=settings.source_catalog_statbel_max_response_mb, accept="text/turtle,application/x-turtle,application/octet-stream,text/plain", ) validate_statbel_catalog_url(response.final_url) content_type = response.content_type.lower() if content_type and not any(token in content_type for token in ("turtle", "octet-stream", "text/plain")): raise StatbelCatalogError( "CATALOG_STATBEL_INVALID_CONTENT_TYPE", "De officiële Statbel DCAT-catalogus heeft geen ondersteund Turtle-contenttype.", ) release = parse_statbel_population_catalog(response.content) except StatbelCatalogError as exc: raise CatalogProbeFailure(exc.code, exc.message) from exc layout = "nieuwe REDEGEO-sectorindeling" if release.current_distribution_variant == "new" else "actuele sectorindeling" message = f"De officiële Statbel DCAT-catalogus bevestigt bevolkingseditie {release.version} met de {layout}." if release.legacy_distribution_available: message += " De oude 2025-indeling is alleen overgangsevidentie en wordt niet als actuele GeoIntel-editie gebruikt." return RemoteProbe( status="available", reachable=True, checked_at=now, expected_layers=contract.expected_layers, matched_layers=contract.expected_layers, advertised_layer_count=release.distribution_count, metadata_url=release.landing_page, metadata_identifier=release.identifier, remote_title=f"Bevolking per statistische sector {release.version} ({layout})", remote_version=release.version, remote_modified_at=release.catalog_modified_at, capabilities_sha256=sha256(response.content).hexdigest(), capabilities_etag=response.etag, capabilities_last_modified_at=response.last_modified_at, message=message, ) def _probe_remote( contract: ProbeContract, settings: Settings, *, opener: Callable[..., Any] | None, now: datetime, ) -> RemoteProbe: try: if contract.service_type == "HTML": return _probe_alz_remote(contract, settings, opener=opener, now=now) if contract.service_type == "DCAT": return _probe_statbel_remote(contract, settings, opener=opener, now=now) 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": official_editions = [ item for item in candidates if _ORTHOPHOTO_EDITION.fullmatch((item.source_version or "").strip()) ] if official_editions: latest = max( official_editions, key=lambda item: ( tuple(int(value) for value in (item.source_version or "0.0").split(".")), _utc(item.imported_at) if item.imported_at else datetime.min.replace(tzinfo=timezone.utc), str(item.id), ), ) return latest.source_version 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 elif source_name == _ALZ_SOURCE_NAME: definitive = [item for item in candidates if _ALZ_EDITION.fullmatch((item.source_version or "").strip())] if definitive: latest = max( definitive, key=lambda item: ( int(_ALZ_EDITION.fullmatch((item.source_version or "").strip()).group(1)), _utc(item.observed_at) if item.observed_at else datetime.min.replace(tzinfo=timezone.utc), str(item.id), ), ) return latest.source_version elif source_name == "statbel": annual = [item for item in candidates if _YEAR_EDITION.fullmatch((item.source_version or "").strip())] if annual: return max(annual, key=lambda item: int((item.source_version or "0").strip())).source_version 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 if source_name == _ALZ_SOURCE_NAME: match = _ALZ_EDITION.fullmatch(value) return f"{match.group(1)}-v3" if match else None if source_name == "statbel" and _YEAR_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"), ), ProbeContract( source_name="statbel", display_name="Bevolking per statistische sector (Statbel)", service_type="DCAT", endpoint_url=active_settings.source_catalog_statbel_dcat_url, expected_layers=("population_txt_current", "landing_page", "cc_by_4_0"), ), ProbeContract( source_name=_ALZ_SOURCE_NAME, display_name="Landbouwgebruikspercelen (ALZ)", service_type="HTML", endpoint_url=active_settings.source_catalog_alz_release_url, expected_layers=("definitive_archive", "current_snapshot"), ), ) 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, gekoppelde ISO 19139 metadata, de officiële Statbel DCAT-catalogus en de officiële ALZ-publicatiepagina.", "Er worden geen features, rasters of modelbestanden opgehaald en geen datasets aangemaakt of overschreven.", "Statbel distributielinks worden alleen als release-evidentie gevalideerd; de oude en nieuwe 2025-sectorindeling blijven semantisch gescheiden.", "Voor ALZ is alleen de nieuwste definitieve v3-editie vergelijkbaar; voorlopige v1/v2-snapshots zijn uitsluitend informatief.", "Een versieverschil is controlesignaal, geen bewijs dat een lokale dataset onbruikbaar is en geen automatische importopdracht.", ], )