from __future__ import annotations from dataclasses import dataclass from datetime import UTC, datetime, timedelta import hashlib import json import math from pathlib import Path import re from typing import Any, Callable from urllib.error import HTTPError, URLError from urllib.parse import urlencode, urlparse from urllib.request import HTTPRedirectHandler, Request, build_opener from uuid import UUID from geoalchemy2.shape import to_shape from pyproj import Transformer from shapely.geometry import MultiPolygon, Polygon, box, mapping, shape from shapely.ops import transform, unary_union from shapely.validation import make_valid from app.core.config import Settings, get_settings from app.core.errors import AppError from app.models import Area, Dataset, Project from app.schemas.official_vector import ( OfficialVectorAcquireRequest, OfficialVectorAcquisitionResult, OfficialVectorProductRead, ) from app.services.dataset_service import DatasetService class _RejectRedirects(HTTPRedirectHandler): def redirect_request(self, req, fp, code, msg, headers, newurl): del req, fp, code, msg, headers, newurl return None _NO_REDIRECT_OPENER = build_opener(_RejectRedirects()) _TO_LAMBERT72 = Transformer.from_crs("EPSG:4326", "EPSG:31370", always_xy=True) _TO_WGS84 = Transformer.from_crs("EPSG:31370", "EPSG:4326", always_xy=True) @dataclass(frozen=True) class OfficialVectorProduct: key: str display_name: str theme: str provider: str source_name: str reference_layer_name: str service_type: str collection: str source_crs: str source_version: str observation_label: str authority_level: str catalog_url: str attribution: str license_note: str limitation_message: str source: str observed_at: datetime valid_from: datetime | None valid_to: datetime | None primary_metric: dict[str, Any] selection_metrics: tuple[dict[str, Any], ...] class OfficialVectorAcquisitionService: _BWK_EVALUATION_LABELS = { "z": "Biologisch zeer waardevol", "w": "Biologisch waardevol", "m": "Biologisch minder waardevol", "wz": "Complex van biologisch waardevolle en zeer waardevolle elementen", "mwz": "Complex van minder waardevolle, waardevolle en zeer waardevolle elementen", "mz": "Complex van minder waardevolle en zeer waardevolle elementen", "mw": "Complex van minder waardevolle en waardevolle elementen", } @staticmethod def _products() -> dict[str, OfficialVectorProduct]: share_warning = ( "PHAB-aandelen gelden voor het volledige bronpolygoon. Bij een gedeeltelijke selectie worden " "ze evenredig geschaald en blijven ze dus een oppervlakte-inschatting." ) products = ( OfficialVectorProduct( key="bwk_natura2000_2025", display_name="BWK en Natura 2000-habitatkaart 2025", theme="nature_value", provider="INBO / Digitaal Vlaanderen", source_name="inbo_bwk_natura2000", reference_layer_name="nature_value", service_type="WFS 2.0", collection="BWK:Bwkhab", source_crs="EPSG:31370", source_version="2025", observation_label="Toestand 2025", authority_level="authoritative", catalog_url=( "https://www.vlaanderen.be/datavindplaats/catalogus/" "biologische-waarderingskaart-en-natura-2000-habitatkaart-toestand-2025" ), attribution="Bron: INBO", license_note="Hergebruik volgens de open-datavoorwaarden en bronvermelding van INBO.", limitation_message=( "De BWK is een gebiedsdekkende kartering, geen terreinmeting op aanvraag. " "PHAB-oppervlakten zijn proportionele schattingen binnen bronpolygonen." ), source="INBO BWK WFS", observed_at=datetime(2025, 12, 10, tzinfo=UTC), valid_from=None, valid_to=None, primary_metric={ "metric_key": "nature_mapped_area", "method": "intersection_area", "label": "Gekarteerde natuuroppervlakte", "unit": "ha", "geometry_dimension": 2, "is_estimate": False, }, selection_metrics=( { "metric_key": "bwk_very_valuable_area", "method": "intersection_area", "label": "Biologisch zeer waardevol", "unit": "ha", "geometry_dimension": 2, "filter_property": "bwk_evaluation_code", "filter_values": ["z"], }, { "metric_key": "bwk_valuable_area", "method": "intersection_area", "label": "Biologisch waardevol", "unit": "ha", "geometry_dimension": 2, "filter_property": "bwk_evaluation_code", "filter_values": ["w"], }, { "metric_key": "bwk_less_valuable_area", "method": "intersection_area", "label": "Biologisch minder waardevol", "unit": "ha", "geometry_dimension": 2, "filter_property": "bwk_evaluation_code", "filter_values": ["m"], }, { "metric_key": "bwk_mixed_value_area", "method": "intersection_area", "label": "Gemengde BWK-waardering", "unit": "ha", "geometry_dimension": 2, "filter_property": "bwk_evaluation_code", "filter_values": ["wz", "mwz", "mz", "mw"], }, { "metric_key": "natura2000_area", "method": "area_weighted_sum", "property": "natura2000_area_ha", "label": "Natura 2000-habitat", "unit": "ha", "is_estimate": True, "warning": share_warning, "warning_only_when_estimate": False, }, { "metric_key": "regional_biotope_area", "method": "area_weighted_sum", "property": "regional_biotope_area_ha", "label": "Regionaal belangrijk biotoop", "unit": "ha", "is_estimate": True, "warning": share_warning, "warning_only_when_estimate": False, }, ), ), OfficialVectorProduct( key="dov_soil_types", display_name="Digitale bodemkaart Vlaanderen - bodemtypes", theme="soil", provider="Databank Ondergrond Vlaanderen", source_name="dov_soil_map", reference_layer_name="soil", service_type="WFS 2.0", collection="bodemkaart:bodemtypes", source_crs="EPSG:31370", source_version="Digitale uitgave juni 2017", observation_label="Veldkartering 1949-1971", authority_level="authoritative_historical_baseline", catalog_url=( "https://www.vlaanderen.be/datavindplaats/catalogus/" "digitale-bodemkaart-van-het-vlaams-gewest-bodemtypes" ), attribution="Databank Ondergrond Vlaanderen - Digitale bodemkaart: bodemtypes", license_note="DOV-bronvermelding en de publieke GDI-hergebruikvoorwaarden zijn van toepassing.", limitation_message=( "Historische bodemkartering op schaal 1:20.000 op basis van veldwerk 1949-1971. " "De huidige drainage en lokale bodemtoestand kunnen afwijken; dit is geen terreinonderzoek." ), source="DOV WFS bodemtypes", observed_at=datetime(1971, 12, 31, 23, 59, 59, tzinfo=UTC), valid_from=datetime(1949, 1, 1, tzinfo=UTC), valid_to=datetime(1971, 12, 31, 23, 59, 59, tzinfo=UTC), primary_metric={ "metric_key": "soil_mapped_area", "method": "intersection_area", "label": "Bodemkaartoppervlakte", "unit": "ha", "geometry_dimension": 2, "is_estimate": False, }, selection_metrics=( { "metric_key": "soil_dry_sand_area", "method": "intersection_area", "label": "Gekarteerd als droog zand", "unit": "ha", "geometry_dimension": 2, "filter_property": "soil_generalized_legend", "filter_values": ["Droog zand", "Zeer droog zand"], }, { "metric_key": "soil_moist_sand_area", "method": "intersection_area", "label": "Gekarteerd als vochtig zand", "unit": "ha", "geometry_dimension": 2, "filter_property": "soil_generalized_legend", "filter_values": ["Vochtig zand"], }, { "metric_key": "soil_wet_sand_area", "method": "intersection_area", "label": "Gekarteerd als nat zand", "unit": "ha", "geometry_dimension": 2, "filter_property": "soil_generalized_legend", "filter_values": ["Nat zand", "Zeer nat zand"], }, { "metric_key": "soil_anthropogenic_area", "method": "intersection_area", "label": "Antropogene bodemklasse", "unit": "ha", "geometry_dimension": 2, "filter_property": "soil_generalized_legend", "filter_values": ["Antropogeen"], }, ), ), ) return {product.key: product for product in products} @staticmethod def list_products() -> list[dict[str, Any]]: return [ OfficialVectorProductRead( key=product.key, display_name=product.display_name, theme=product.theme, provider=product.provider, source_name=product.source_name, reference_layer_name=product.reference_layer_name, service_type=product.service_type, collection=product.collection, geometry_types=["Polygon", "MultiPolygon"], source_crs=product.source_crs, source_version=product.source_version, observation_label=product.observation_label, authority_level=product.authority_level, catalog_url=product.catalog_url, attribution=product.attribution, license_note=product.license_note, limitation_message=product.limitation_message, ).model_dump() for product in OfficialVectorAcquisitionService._products().values() ] @staticmethod def _product(product_key: str) -> OfficialVectorProduct: product = OfficialVectorAcquisitionService._products().get(product_key.strip().lower()) if product is None: raise AppError( code="OFFICIAL_VECTOR_PRODUCT_NOT_SUPPORTED", message="Select a governed official vector product", details={"product_key": product_key}, status_code=422, ) return product @staticmethod def _polygonal(geometry: Any) -> Any | None: if geometry is None or geometry.is_empty: return None if not geometry.is_valid: geometry = make_valid(geometry) parts: list[Polygon] = [] def collect(item: Any) -> None: if item is None or item.is_empty: return if isinstance(item, Polygon): parts.append(item) elif isinstance(item, MultiPolygon): parts.extend(part for part in item.geoms if not part.is_empty) elif hasattr(item, "geoms"): for part in item.geoms: collect(part) collect(geometry) if not parts: return None result = unary_union(parts) if not result.is_valid: result = make_valid(result) return result if not result.is_empty and result.is_valid else None @staticmethod def _validate_scope( db, project_id: UUID, payload: OfficialVectorAcquireRequest, settings: Settings, ) -> tuple[Any, Any, list[float], list[float]]: if not settings.official_vector_enabled: raise AppError( code="OFFICIAL_VECTOR_NOT_CONFIGURED", message="Bounded official vector acquisition is disabled", status_code=503, ) if not db.get(Project, project_id): raise AppError(code="PROJECT_NOT_FOUND", message="Project not found", status_code=404) if payload.bbox.crs.upper() != "EPSG:4326": raise AppError( code="OFFICIAL_VECTOR_INVALID_CRS", message="Official vector acquisition requires EPSG:4326", status_code=400, ) values = (payload.bbox.min_x, payload.bbox.min_y, payload.bbox.max_x, payload.bbox.max_y) if ( not all(math.isfinite(value) for value in values) or values[0] >= values[2] or values[1] >= values[3] or values[0] < -180 or values[2] > 180 or values[1] < -90 or values[3] > 90 ): raise AppError( code="OFFICIAL_VECTOR_INVALID_BBOX", message="Bounding box is invalid for EPSG:4326", status_code=400, ) metric_bounds = _TO_LAMBERT72.transform_bounds(*values, densify_pts=21) width_m = metric_bounds[2] - metric_bounds[0] height_m = metric_bounds[3] - metric_bounds[1] if width_m < settings.official_vector_min_side_m or height_m < settings.official_vector_min_side_m: raise AppError( code="OFFICIAL_VECTOR_SELECTION_TOO_SMALL", message=f"Select at least {settings.official_vector_min_side_m:g} by " f"{settings.official_vector_min_side_m:g} metres", status_code=422, ) if width_m > settings.official_vector_max_side_m or height_m > settings.official_vector_max_side_m: raise AppError( code="OFFICIAL_VECTOR_SELECTION_TOO_LARGE", message=f"Select no more than {settings.official_vector_max_side_m:g} by " f"{settings.official_vector_max_side_m:g} metres", details={"width_m": width_m, "height_m": height_m}, status_code=422, ) scope_wgs84 = box(*values) if payload.area_id: area = db.get(Area, payload.area_id) if area is None: raise AppError(code="AREA_NOT_FOUND", message="Area not found", status_code=404) if area.project_id != project_id: raise AppError( code="INVALID_DATASET_SCOPE", message="Area does not belong to this project", status_code=400, ) scope_wgs84 = OfficialVectorAcquisitionService._polygonal( scope_wgs84.intersection(to_shape(area.geometry)) ) if scope_wgs84 is None: raise AppError( code="OFFICIAL_VECTOR_SCOPE_EMPTY", message="The selection does not intersect the selected area", status_code=400, ) scope_metric = OfficialVectorAcquisitionService._polygonal( transform(_TO_LAMBERT72.transform, scope_wgs84) ) if scope_metric is None: raise AppError( code="OFFICIAL_VECTOR_SCOPE_INVALID", message="The selection could not be transformed to EPSG:31370", status_code=400, ) return scope_wgs84, scope_metric, [float(value) for value in values], [ float(value) for value in scope_metric.bounds ] @staticmethod def _read_page( url: str, settings: Settings, opener: Callable[..., Any] | None, ) -> tuple[dict[str, Any], str, int]: request = Request( url, headers={ "Accept": "application/geo+json, application/json", "User-Agent": "GeoIntel/1.0 bounded-official-vector-acquisition", }, ) try: with (opener or _NO_REDIRECT_OPENER.open)( request, timeout=settings.official_vector_timeout_seconds, ) as response: limit = settings.official_vector_max_response_mb * 1024 * 1024 content = response.read(limit + 1) except HTTPError as exc: raise AppError( code="OFFICIAL_VECTOR_PROVIDER_HTTP_ERROR", message="The official vector provider returned an HTTP error", details={"status_code": exc.code}, status_code=502, ) from exc except (TimeoutError, URLError, OSError) as exc: raise AppError( code="OFFICIAL_VECTOR_PROVIDER_UNAVAILABLE", message="The official vector provider is unavailable", status_code=502, ) from exc if len(content) > limit: raise AppError( code="OFFICIAL_VECTOR_PROVIDER_RESPONSE_TOO_LARGE", message="An official vector response page exceeded the configured limit", status_code=502, ) try: payload = json.loads(content.decode("utf-8")) except (UnicodeDecodeError, json.JSONDecodeError) as exc: raise AppError( code="OFFICIAL_VECTOR_PROVIDER_INVALID_RESPONSE", message="The official vector provider returned invalid GeoJSON", status_code=502, ) from exc if not isinstance(payload, dict) or payload.get("type") != "FeatureCollection": raise AppError( code="OFFICIAL_VECTOR_PROVIDER_INVALID_RESPONSE", message="The official vector provider returned a non-FeatureCollection response", status_code=502, ) return payload, hashlib.sha256(content).hexdigest(), len(content) @staticmethod def _nature_url( settings: Settings, bbox_values: tuple[float, ...], start_index: int, ) -> str: query = urlencode( { "service": "WFS", "version": "2.0.0", "request": "GetFeature", "typeNames": "BWK:Bwkhab", "srsName": "EPSG:4326", "bbox": ",".join(f"{value:.8f}" for value in bbox_values) + ",EPSG:4326", "count": settings.official_vector_page_size, "startIndex": start_index, "sortBy": "UIDN", "outputFormat": "application/json", } ) return f"{settings.bwk_wfs_url.rstrip('?')}?{query}" @staticmethod def _soil_url(settings: Settings, metric_bbox: tuple[float, ...], start_index: int) -> str: query = urlencode( { "service": "WFS", "version": "2.0.0", "request": "GetFeature", "typeNames": "bodemkaart:bodemtypes", "srsName": "EPSG:4326", "bbox": ",".join(f"{value:.3f}" for value in metric_bbox) + ",EPSG:31370", "count": settings.official_vector_page_size, "startIndex": start_index, "sortBy": "gid", "outputFormat": "application/json", } ) return f"{settings.dov_soil_wfs_url.rstrip('?')}?{query}" @staticmethod def _habitat_breakdown(properties: dict[str, Any]) -> tuple[list[dict[str, Any]], float, float, float]: entries: list[dict[str, Any]] = [] natura_share = regional_share = uncertain_share = 0.0 for index in range(1, 6): code = str(properties.get(f"HAB{index}") or "").strip() if not code: continue raw_share = properties.get(f"PHAB{index}") try: share = max(0.0, min(100.0, float(raw_share or 0))) except (TypeError, ValueError): share = 0.0 entries.append({"code": code, "share_percent": share}) lowered = code.lower() if re.match(r"^\d", code): natura_share += share elif lowered.startswith("rbb"): regional_share += share elif lowered.startswith("ohab"): uncertain_share += share if str(properties.get("HABLEGENDE") or "").strip().lower() == "ohab" and uncertain_share <= 0: uncertain_share = 100.0 return entries, min(100.0, natura_share), min(100.0, regional_share), min(100.0, uncertain_share) @staticmethod def _normalize_feature( product: OfficialVectorProduct, feature: dict[str, Any], scope_metric: Any, coverage_scope: str, ) -> dict[str, Any] | None: try: source_wgs84 = OfficialVectorAcquisitionService._polygonal(shape(feature.get("geometry"))) except Exception as exc: raise AppError( code="OFFICIAL_VECTOR_PROVIDER_INVALID_GEOMETRY", message=f"{product.display_name} returned invalid geometry", status_code=502, ) from exc if source_wgs84 is None: return None source_metric = OfficialVectorAcquisitionService._polygonal( transform(_TO_LAMBERT72.transform, source_wgs84) ) if source_metric is None or not source_metric.intersects(scope_metric): return None clipped_metric = OfficialVectorAcquisitionService._polygonal(source_metric.intersection(scope_metric)) if clipped_metric is None or clipped_metric.area <= 0: return None clipped_wgs84 = OfficialVectorAcquisitionService._polygonal( transform(_TO_WGS84.transform, clipped_metric) ) if clipped_wgs84 is None: return None raw = dict(feature.get("properties") or {}) if product.theme == "nature_value": raw_id = str(feature.get("id") or raw.get("UIDN") or raw.get("OIDN") or "").strip() if not raw_id: raw_id = hashlib.sha256(json.dumps(feature.get("geometry"), sort_keys=True).encode()).hexdigest() feature_id = f"BWK:Bwkhab:{raw.get('UIDN') or raw_id}" evaluation = str(raw.get("EVAL") or "").strip().lower() habitats, natura_share, regional_share, uncertain_share = ( OfficialVectorAcquisitionService._habitat_breakdown(raw) ) area_ha = float(clipped_metric.area) / 10_000.0 properties = { **raw, "source_name": product.source_name, "source_collection": product.collection, "source_feature_id": feature_id, "reference_layer_name": product.reference_layer_name, "theme": product.theme, "authority_level": product.authority_level, "coverage_scope": coverage_scope, "source_version": product.source_version, "attribution": product.attribution, "bwk_evaluation_code": evaluation or "unknown", "bwk_evaluation_label": OfficialVectorAcquisitionService._BWK_EVALUATION_LABELS.get( evaluation, "Onbekende of ontbrekende BWK-waardering" ), "bwk_label": str(raw.get("BWKLABEL") or "").strip(), "bwk_units": ", ".join( str(raw.get(f"EENH{index}") or "").strip() for index in range(1, 9) if str(raw.get(f"EENH{index}") or "").strip() ), "habitat_entries": habitats, "clipped_area_ha": round(area_ha, 8), "natura2000_share_percent": natura_share, "regional_biotope_share_percent": regional_share, "uncertain_habitat_share_percent": uncertain_share, "natura2000_area_ha": round(area_ha * natura_share / 100.0, 8), "regional_biotope_area_ha": round(area_ha * regional_share / 100.0, 8), "uncertain_habitat_area_ha": round(area_ha * uncertain_share / 100.0, 8), "geometry_clipped_to_selection": not scope_metric.covers(source_metric), } else: gid = raw.get("gid") map_polygon_id = raw.get("id_kaartvlak") feature_id = str(feature.get("id") or f"{product.collection}:{gid or map_polygon_id}").strip() if not feature_id: raise AppError( code="OFFICIAL_VECTOR_PROVIDER_INVALID_RESPONSE", message="DOV returned a soil polygon without an official identity", status_code=502, ) properties = { **raw, "source_name": product.source_name, "source_collection": product.collection, "source_feature_id": feature_id, "source_gid": gid, "source_map_polygon_id": map_polygon_id, "reference_layer_name": product.reference_layer_name, "theme": product.theme, "authority_level": product.authority_level, "coverage_scope": coverage_scope, "source_version": product.source_version, "survey_period": "1949-1971", "soil_type_code": raw.get("Bodemtype"), "unified_soil_type_code": raw.get("Unibodemtype"), "soil_series_code": raw.get("Bodemserie"), "soil_series_description": raw.get("Beknopte_omschrijving_bodemserie"), "soil_generalized_legend": raw.get("Gegeneraliseerde_legende"), "soil_texture_class_code": raw.get("Textuurklasse_code"), "soil_texture_class": raw.get("Textuurklasse"), "soil_drainage_class_code": raw.get("Drainageklasse_code"), "soil_drainage_class": raw.get("Drainageklasse"), "soil_profile_group_code": raw.get("Profielontwikkelingsgroep_code"), "soil_profile_group": raw.get("Profielontwikkelingsgroep"), "soil_substrate_code": raw.get("Substraat_code"), "soil_substrate": raw.get("Substraat_Vlaanderen") or raw.get("Substraat_legende"), "soil_region": raw.get("Streek"), "classification_type": raw.get("Type_classificatie"), "soil_map_title": raw.get("Eenduidige_legende_titel"), "clipped_area_ha": round(float(clipped_metric.area) / 10_000.0, 8), "attribution": product.attribution, "geometry_clipped_to_selection": not scope_metric.covers(source_metric), "historical_drainage_limitation": ( "Drainage class derives from field data collected between 1949 and 1971 and may differ today." ), } return { "type": "Feature", "id": feature_id, "geometry": mapping(clipped_wgs84), "properties": properties, } @staticmethod def _fetch_features( product: OfficialVectorProduct, scope_wgs84: Any, scope_metric: Any, coverage_scope: str, settings: Settings, opener: Callable[..., Any] | None, ) -> tuple[list[dict[str, Any]], dict[str, Any]]: retained: list[dict[str, Any]] = [] seen_ids: set[str] = set() request_urls: list[str] = [] response_hashes: list[str] = [] total_bytes = candidate_count = 0 expected_total: int | None = None next_url = ( OfficialVectorAcquisitionService._nature_url(settings, tuple(scope_wgs84.bounds), 0) if product.theme == "nature_value" else OfficialVectorAcquisitionService._soil_url(settings, tuple(scope_metric.bounds), 0) ) start_index = 0 seen_pages: set[str] = set() while next_url: parsed = urlparse(next_url) configured_url = ( settings.bwk_wfs_url if product.theme == "nature_value" else settings.dov_soil_wfs_url ) base = urlparse(configured_url) if ( parsed.scheme != "https" or base.scheme != "https" or parsed.netloc.casefold() != base.netloc.casefold() or parsed.path != base.path ): raise AppError( code="OFFICIAL_VECTOR_PROVIDER_INVALID_PAGINATION", message="The official WFS request escaped the governed endpoint", status_code=502, ) if next_url in seen_pages: raise AppError( code="OFFICIAL_VECTOR_PROVIDER_PAGINATION_LOOP", message="The official provider repeated a pagination URL", status_code=502, ) if len(request_urls) >= settings.official_vector_max_pages: raise AppError( code="OFFICIAL_VECTOR_SELECTION_TOO_LARGE", message="Official vector acquisition exceeded the configured page limit", status_code=422, ) seen_pages.add(next_url) payload, response_hash, response_size = OfficialVectorAcquisitionService._read_page( next_url, settings, opener ) request_urls.append(next_url) response_hashes.append(response_hash) total_bytes += response_size if total_bytes > settings.official_vector_max_total_response_mb * 1024 * 1024: raise AppError( code="OFFICIAL_VECTOR_PROVIDER_RESPONSE_TOO_LARGE", message="The complete official vector response exceeded the configured transfer limit", status_code=502, ) source_features = payload.get("features") if not isinstance(source_features, list): raise AppError( code="OFFICIAL_VECTOR_PROVIDER_INVALID_RESPONSE", message="The official FeatureCollection has no feature list", status_code=502, ) raw_matched = payload.get("numberMatched", payload.get("totalFeatures")) if raw_matched not in (None, "unknown"): try: matched = int(raw_matched) except (TypeError, ValueError) as exc: raise AppError( code="OFFICIAL_VECTOR_PROVIDER_INVALID_RESPONSE", message="The official WFS returned an invalid numberMatched value", status_code=502, ) from exc if expected_total is None: expected_total = matched elif expected_total != matched: raise AppError( code="OFFICIAL_VECTOR_PROVIDER_UNSTABLE_PAGINATION", message="The official WFS numberMatched changed during pagination", status_code=502, ) for source_feature in source_features: candidate_count += 1 if not isinstance(source_feature, dict): continue normalized = OfficialVectorAcquisitionService._normalize_feature( product, source_feature, scope_metric, coverage_scope ) if normalized is None or normalized["id"] in seen_ids: continue seen_ids.add(normalized["id"]) if len(retained) >= settings.official_vector_max_features: raise AppError( code="OFFICIAL_VECTOR_SELECTION_TOO_LARGE", message="The selection exceeds the configured feature limit; draw a smaller rectangle", details={"max_features": settings.official_vector_max_features}, status_code=422, ) retained.append(normalized) returned = payload.get("numberReturned", len(source_features)) try: returned_count = int(returned) except (TypeError, ValueError) as exc: raise AppError( code="OFFICIAL_VECTOR_PROVIDER_INVALID_RESPONSE", message="The official WFS returned an invalid numberReturned value", status_code=502, ) from exc if returned_count != len(source_features): raise AppError( code="OFFICIAL_VECTOR_PROVIDER_INVALID_RESPONSE", message="The official WFS numberReturned does not match its feature payload", status_code=502, ) start_index += returned_count if returned_count == 0 or ( expected_total is not None and start_index >= expected_total ) or (expected_total is None and returned_count < settings.official_vector_page_size): if expected_total is not None and start_index != expected_total: raise AppError( code="OFFICIAL_VECTOR_PROVIDER_INCOMPLETE_RESPONSE", message="The official WFS did not return every matched feature", details={"received": start_index, "expected": expected_total}, status_code=502, ) next_url = None elif product.theme == "nature_value": next_url = OfficialVectorAcquisitionService._nature_url( settings, tuple(scope_wgs84.bounds), start_index ) else: next_url = OfficialVectorAcquisitionService._soil_url( settings, tuple(scope_metric.bounds), start_index ) return retained, { "candidate_feature_count": candidate_count, "feature_count": len(retained), "page_count": len(request_urls), "request_urls": request_urls, "response_sha256": response_hashes, "response_size_bytes": total_bytes, "reference_truncated": False, } @staticmethod def _cached_dataset( db, project_id: UUID, product: OfficialVectorProduct, request_hash: str, settings: Settings, ) -> Dataset | None: if settings.official_vector_cache_ttl_hours <= 0: return None candidates = ( db.query(Dataset) .filter( Dataset.project_id == project_id, Dataset.source_name == product.source_name, Dataset.reference_layer_name == product.reference_layer_name, Dataset.status == "ready", ) .order_by(Dataset.imported_at.desc()) .all() ) cutoff = datetime.now(UTC) - timedelta(hours=settings.official_vector_cache_ttl_hours) for candidate in candidates: provenance = candidate.provenance_metadata if isinstance(candidate.provenance_metadata, dict) else {} if ( provenance.get("request_hash") == request_hash and candidate.storage_path and Path(candidate.storage_path).is_file() and candidate.imported_at is not None and candidate.imported_at >= cutoff ): return candidate return None @staticmethod def _result( dataset: Dataset, product: OfficialVectorProduct, *, reused: bool, bbox_values: list[float], ) -> dict[str, Any]: source_metadata = dataset.source_metadata if isinstance(dataset.source_metadata, dict) else {} provenance = dataset.provenance_metadata if isinstance(dataset.provenance_metadata, dict) else {} metadata = dataset.metadata_json if isinstance(dataset.metadata_json, dict) else {} return OfficialVectorAcquisitionResult( output_dataset_id=dataset.id, reused=reused, product_key=product.key, display_name=product.display_name, theme=product.theme, provider=product.provider, source_name=product.source_name, reference_layer_name=product.reference_layer_name, service_type=product.service_type, collection=product.collection, feature_count=int(metadata.get("feature_count", source_metadata.get("feature_count", 0))), candidate_feature_count=int(provenance.get("candidate_feature_count", 0)), page_count=int(provenance.get("page_count", 0)), bbox_epsg4326=bbox_values, source_version=str(dataset.source_version or product.source_version), attribution=product.attribution, limitation_message=product.limitation_message, ).model_dump(mode="json") @staticmethod def acquire( db, project_id: UUID, payload: OfficialVectorAcquireRequest, *, settings: Settings | None = None, opener: Callable[..., Any] | None = None, ) -> dict[str, Any]: resolved_settings = settings or get_settings() product = OfficialVectorAcquisitionService._product(payload.product_key) scope_wgs84, scope_metric, bbox_values, metric_bounds = ( OfficialVectorAcquisitionService._validate_scope( db, project_id, payload, resolved_settings ) ) request_identity = { "product_key": product.key, "bbox_epsg4326": [round(value, 8) for value in bbox_values], "area_id": str(payload.area_id) if payload.area_id else None, "source_version": product.source_version, } request_hash = hashlib.sha256( json.dumps(request_identity, sort_keys=True).encode() ).hexdigest() if not payload.force_refresh: cached = OfficialVectorAcquisitionService._cached_dataset( db, project_id, product, request_hash, resolved_settings ) if cached is not None: return OfficialVectorAcquisitionService._result( cached, product, reused=True, bbox_values=bbox_values ) area = db.get(Area, payload.area_id) if payload.area_id else None coverage_scope = ( "municipality" if area is not None and area.name.strip().lower().startswith("gemeente ") else "bounded_selection" ) features, transfer = OfficialVectorAcquisitionService._fetch_features( product, scope_wgs84, scope_metric, coverage_scope, resolved_settings, opener, ) acquired_at = datetime.now(UTC) artifact = json.dumps( { "type": "FeatureCollection", "name": product.display_name, "crs": {"type": "name", "properties": {"name": "EPSG:4326"}}, "features": features, }, ensure_ascii=False, separators=(",", ":"), ).encode("utf-8") filename = ( f"{product.source_name}_{product.key}_{request_hash[:12]}.geojson" ) source_metadata = { "provider": product.provider, "service": product.service_type, "product_key": product.key, "product_display_name": product.display_name, "source_collection": product.collection, "authority_level": product.authority_level, "theme": product.theme, "layer_type": product.reference_layer_name, "coverage_scope": coverage_scope, "geometry_clipped_to_area": payload.area_id is not None, "geometry_clipped_to_selection": True, "bbox_epsg4326": bbox_values, "bbox_epsg31370": metric_bounds, "feature_count": len(features), "identity_stable": True, "source_storage_crs": product.source_crs, "persisted_crs": "EPSG:4326", "selection_aggregation": product.primary_metric, "selection_metrics": list(product.selection_metrics), "attribution": product.attribution, "license_note": product.license_note, "catalog_url": product.catalog_url, "limitation_message": product.limitation_message, } if product.theme == "soil": source_metadata.update( { "survey_period": "1949-1971", "source_scale": "1:20,000", "semantic_metrics": False, } ) provenance_metadata = { "acquisition": f"explicit_bounded_{product.service_type.lower().replace(' ', '_')}", "acquired_at": acquired_at.isoformat(), "request_hash": request_hash, "request_urls": transfer["request_urls"], "response_sha256": transfer["response_sha256"], "response_size_bytes": transfer["response_size_bytes"], "page_count": transfer["page_count"], "candidate_feature_count": transfer["candidate_feature_count"], "exact_feature_count": transfer["feature_count"], "reference_truncated": False, "artifact_sha256": hashlib.sha256(artifact).hexdigest(), "clipped_to_area_id": str(payload.area_id) if payload.area_id else None, "scope_geometry_type": scope_wgs84.geom_type, "catalog_url": product.catalog_url, "limitation_message": product.limitation_message, } try: dataset_response = DatasetService.import_vector_bytes( db, project_id=project_id, area_id=payload.area_id, filename=filename, content=artifact, source=product.source, source_name=product.source_name, dataset_role="reference", reference_layer_name=product.reference_layer_name, temporal_series_key=f"{product.source_name}:{product.key}:{request_hash[:24]}", observed_at=product.observed_at, valid_from=product.valid_from, valid_to=product.valid_to, temporal_granularity="period" if product.valid_from else "snapshot", source_version=product.source_version, source_metadata=source_metadata, provenance_metadata=provenance_metadata, ) except AppError: raise except Exception as exc: raise AppError( code="OFFICIAL_VECTOR_PERSISTENCE_FAILED", message="The validated official vector selection could not be persisted", details={"reason": str(exc)}, status_code=500, ) from exc persisted = db.get(Dataset, dataset_response.id) if persisted is None: raise AppError( code="OFFICIAL_VECTOR_PERSISTENCE_FAILED", message="The persisted official vector dataset could not be reloaded", status_code=500, ) return OfficialVectorAcquisitionService._result( persisted, product, reused=False, bbox_values=bbox_values )