798 lines
33 KiB
Python
798 lines
33 KiB
Python
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.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}$")
|
||
_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)
|
||
|
||
|
||
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) -> FetchResult:
|
||
request = Request(
|
||
url,
|
||
headers={"Accept": "application/xml,text/xml,text/html,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 _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_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)
|
||
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
|
||
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
|
||
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
|
||
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=_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 en de officiële ALZ-publicatiepagina.",
|
||
"Er worden geen features, rasters of modelbestanden opgehaald en geen datasets aangemaakt of overschreven.",
|
||
"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.",
|
||
],
|
||
)
|