"""Provision official Waterinfo station time series for a persisted Area. This operator discovers annual station series through the public Waterinfo KiWIS service, filters station points against the exact persisted Area and imports one immutable GeoJSON point dataset per station and observation year. It never averages different stations and never interprets a point measurement as area-wide water volume. """ from __future__ import annotations import argparse import hashlib import json import math import os import re import sys from dataclasses import dataclass from datetime import datetime, timezone from pathlib import Path from typing import Any import requests from requests.adapters import HTTPAdapter from shapely.geometry import Point, mapping, shape from urllib3.util.retry import Retry DEFAULT_API_URL = "http://127.0.0.1:8000" DEFAULT_PROJECT_NAME = "Kempen Regional Workbench" DEFAULT_AREA_NAME = "Gemeente Mol" DEFAULT_OUTPUT_DIR = Path("/app/storage/operator-data/waterinfo/mol") KIWIS_URL = "https://download.waterinfo.be/tsmdownload/KiWIS/KiWIS" CATALOG_URL = "https://waterinfo.vlaanderen.be/" ATTRIBUTION = "Bron: Waterinfo Vlaanderen / Vlaamse Milieumaatschappij" @dataclass(frozen=True) class ParameterDefinition: key: str group_id: str property_name: str label: str unit: str reference_layer_name: str limitation: str PARAMETERS = { "water_level": ParameterDefinition( key="water_level", group_id="192784", property_name="annual_mean_water_level_m", label="Jaargemiddelde waterstand", unit="m", reference_layer_name="water_level_station", limitation=( "Dit is een jaargemiddelde op één meetstation. De waarde geldt niet voor het volledige geselecteerde gebied " "en levert zonder profiel- of bathymetriegegevens geen watervolume op." ), ), "discharge": ParameterDefinition( key="discharge", group_id="192895", property_name="annual_mean_discharge_m3_s", label="Jaargemiddeld debiet", unit="m³/s", reference_layer_name="discharge_station", limitation=( "Dit is een jaargemiddeld debiet op één meetstation. De waarde geldt niet voor alle waterlopen in het " "geselecteerde gebied en is geen gebiedsdekkend watervolume." ), ), } def parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser(description="Provision official annual Waterinfo station histories.") parser.add_argument("--base-url", default=os.environ.get("GEOINTEL_INTERNAL_API_URL", DEFAULT_API_URL)) parser.add_argument("--project-name", default=DEFAULT_PROJECT_NAME) parser.add_argument("--area-name", default=DEFAULT_AREA_NAME) parser.add_argument("--parameters", default="water_level,discharge") parser.add_argument("--from-year", type=int, default=2013) parser.add_argument("--to-year", type=int, default=datetime.now(timezone.utc).year) parser.add_argument("--min-observations", type=int, default=2) parser.add_argument("--max-stations", type=int, default=50) parser.add_argument("--output-dir", type=Path, default=Path(os.environ.get("WATERINFO_OUTPUT_DIR", DEFAULT_OUTPUT_DIR))) parser.add_argument("--request-timeout", type=int, default=120) parser.add_argument("--import-timeout", type=int, default=300) parser.add_argument("--force", action="store_true", help="Refetch source artifacts; persisted datasets remain immutable.") parser.add_argument("--fetch-only", action="store_true") return parser.parse_args() def build_session() -> requests.Session: retry = Retry( total=5, connect=5, read=5, status=5, backoff_factor=1.0, status_forcelist=(429, 500, 502, 503, 504), allowed_methods=frozenset({"GET"}), raise_on_status=True, ) session = requests.Session() session.headers.update({"User-Agent": "GeoIntel-Waterinfo-Operator/1.0"}) adapter = HTTPAdapter(max_retries=retry) session.mount("https://", adapter) session.mount("http://", adapter) return session def response_data(response: requests.Response) -> Any: try: payload = response.json() except ValueError as exc: raise RuntimeError(f"GeoIntel API returned non-JSON ({response.status_code}): {response.text[:300]}") from exc if not response.ok: raise RuntimeError(f"GeoIntel API failed ({response.status_code}): {json.dumps(payload, ensure_ascii=False)[:800]}") if not isinstance(payload, dict) or "data" not in payload: raise RuntimeError("GeoIntel API response does not use the canonical data envelope") return payload["data"] def list_paginated_items(session: requests.Session, url: str, *, timeout: int) -> list[dict[str, Any]]: items: list[dict[str, Any]] = [] offset = 0 total: int | None = None while total is None or offset < total: page = response_data(session.get(url, params={"limit": 200, "offset": offset}, timeout=timeout)) page_items = page.get("items") if isinstance(page, dict) else None if not isinstance(page_items, list): raise RuntimeError(f"GeoIntel list response for {url} has no items array") if total is None: total = int(page.get("total", len(page_items))) items.extend(page_items) if not page_items: break offset += len(page_items) if total is not None and len(items) != total: raise RuntimeError(f"GeoIntel list response for {url} returned {len(items)} of {total} items") return items def locate_workspace(session: requests.Session, base_url: str, args: argparse.Namespace): projects = list_paginated_items(session, f"{base_url}/api/v1/projects", timeout=args.import_timeout) project = next((item for item in projects if item.get("name") == args.project_name), None) if not project: raise RuntimeError(f"Project {args.project_name!r} is missing") project_id = str(project["id"]) areas = list_paginated_items( session, f"{base_url}/api/v1/projects/{project_id}/areas", timeout=args.import_timeout, ) fragment = args.area_name.strip().casefold() matches = [item for item in areas if fragment in str(item.get("name") or "").casefold()] if len(matches) != 1 or not isinstance(matches[0].get("geometry"), dict): raise RuntimeError(f"Expected one persisted Area with geometry matching {args.area_name!r}, received {len(matches)}") area = matches[0] datasets = list_paginated_items( session, f"{base_url}/api/v1/projects/{project_id}/datasets", timeout=args.import_timeout, ) return project_id, str(area["id"]), shape(area["geometry"]), datasets def fetch_json(session: requests.Session, params: dict[str, Any], *, timeout: int) -> Any: response = session.get(KIWIS_URL, params=params, timeout=timeout) response.raise_for_status() try: return response.json() except ValueError as exc: raise RuntimeError(f"Waterinfo returned non-JSON: {response.text[:300]}") from exc def discover_station_series( session: requests.Session, parameter: ParameterDefinition, area_geometry, *, timeout: int, ) -> tuple[dict[str, Any], list[dict[str, Any]]]: payload = fetch_json( session, { "service": "kisters", "type": "queryServices", "request": "getTimeseriesValueLayer", "datasource": 1, "format": "geojson", "timeseriesgroup_id": parameter.group_id, "metadata": "true", "md_returnfields": ( "custom_attributes,station_id,station_no,station_name,ts_id,ts_name," "stationparameter_name,ts_unitsymbol,parametertype_name" ), "custattr_returnfields": "dataprovider,dataowner", "invalidValue": -9999, "invalidPeriod": "P2Y", }, timeout=timeout, ) if not isinstance(payload, dict) or payload.get("type") != "FeatureCollection": raise RuntimeError(f"Waterinfo group {parameter.group_id} did not return a GeoJSON FeatureCollection") selected: list[dict[str, Any]] = [] for feature in payload.get("features") or []: geometry = feature.get("geometry") if isinstance(feature, dict) else None properties = feature.get("properties") if isinstance(feature, dict) else None if not isinstance(geometry, dict) or geometry.get("type") != "Point" or not isinstance(properties, dict): continue point = shape(geometry) ts_id = properties.get("ts_id") if point.is_empty or ts_id in (None, "") or not area_geometry.covers(point): continue selected.append({"geometry": geometry, "properties": properties, "ts_id": str(ts_id)}) return payload, selected def fetch_annual_values( session: requests.Session, ts_id: str, *, from_year: int, to_year: int, timeout: int, ) -> tuple[Any, dict[int, float]]: payload = fetch_json( session, { "service": "kisters", "type": "queryServices", "request": "getTimeseriesValues", "datasource": 1, "format": "json", "ts_id": ts_id, "metadata": "true", "md_returnfields": "station_name,station_no,stationparameter_name,ts_id,ts_unitsymbol", "from": f"{from_year}-01-01", "to": f"{to_year}-12-31", }, timeout=timeout, ) entries = payload if isinstance(payload, list) else [payload] values: dict[int, float] = {} for entry in entries: if not isinstance(entry, dict): continue for row in entry.get("data") or []: if not isinstance(row, list) or len(row) < 2: continue try: year = int(str(row[0])[:4]) value = float(row[1]) except (TypeError, ValueError): continue if year < from_year or year > to_year or not math.isfinite(value) or value <= -9999: continue if year in values and not math.isclose(values[year], value, rel_tol=0.0, abs_tol=1e-12): raise RuntimeError(f"Waterinfo series {ts_id} contains multiple annual values for {year}") values[year] = value return payload, values def safe_slug(value: str) -> str: normalized = re.sub(r"[^a-z0-9]+", "-", value.casefold()).strip("-") return normalized[:80] or "station" def series_key(parameter: ParameterDefinition, station: dict[str, Any]) -> str: properties = station["properties"] station_identity = str(properties.get("station_no") or properties.get("station_id") or station["ts_id"]) return f"waterinfo:{parameter.key}:annual:{safe_slug(station_identity)}" def build_snapshot( parameter: ParameterDefinition, station: dict[str, Any], year: int, value: float, ) -> dict[str, Any]: properties = station["properties"] station_name = str(properties.get("station_name") or "Waterinfo meetstation") station_no = str(properties.get("station_no") or properties.get("station_id") or station["ts_id"]) feature_id = f"waterinfo-{parameter.key}-{safe_slug(station_no)}-{year}" return { "type": "FeatureCollection", "name": f"{parameter.label} - {station_name} - {year}", "crs": {"type": "name", "properties": {"name": "EPSG:4326"}}, "features": [ { "type": "Feature", "id": feature_id, "geometry": station["geometry"], "properties": { "source_feature_id": feature_id, "source_name": "waterinfo", "theme": "water", "measurement_type": parameter.key, "observation_year": year, parameter.property_name: value, "station_id": properties.get("station_id"), "station_no": properties.get("station_no"), "station_name": station_name, "timeseries_id": station["ts_id"], "timeseries_name": properties.get("ts_name"), "reported_unit": properties.get("ts_unitsymbol"), "data_provider": properties.get("dataprovider"), "data_owner": properties.get("dataowner"), "authority_level": "authoritative", "attribution": ATTRIBUTION, }, } ], } def sha256_bytes(content: bytes) -> str: return hashlib.sha256(content).hexdigest() def write_json(path: Path, payload: Any, *, pretty: bool = False) -> str: content = json.dumps( payload, ensure_ascii=False, indent=2 if pretty else None, separators=None if pretty else (",", ":"), sort_keys=pretty, ).encode("utf-8") temporary = path.with_suffix(f"{path.suffix}.partial") temporary.write_bytes(content) temporary.replace(path) return sha256_bytes(content) def upload_snapshot( session: requests.Session, *, base_url: str, project_id: str, area_id: str, parameter: ParameterDefinition, station: dict[str, Any], year: int, path: Path, raw_sha256: str, coverage: tuple[int, int, int], timeout: int, ) -> dict[str, Any]: properties = station["properties"] station_name = str(properties.get("station_name") or "Waterinfo meetstation") key = series_key(parameter, station) observed_at = f"{year}-01-01T00:00:00Z" source_metadata = { "provider": "Waterinfo Vlaanderen / Vlaamse Milieumaatschappij", "authority_level": "authoritative", "theme": "water", "semantic_metrics": False, "geometry_clipped_to_area": True, "measurement_type": parameter.key, "station_name": station_name, "station_no": properties.get("station_no"), "timeseries_id": station["ts_id"], "attribution": ATTRIBUTION, "catalog_url": CATALOG_URL, "temporal_series_label": f"{station_name} - {parameter.label.lower()}", "observation_date_precision": "year", "selection_aggregation": { "metric_key": parameter.key, "method": "mean", "property": parameter.property_name, "label": parameter.label, "unit": parameter.unit, "warning": parameter.limitation, }, } provenance_metadata = { "operator_tool": "provision_waterinfo_station_history.py", "operator_explicit_fetch": True, "geometry_clipped_to_area": True, "kiwis_url": KIWIS_URL, "timeseries_group_id": parameter.group_id, "timeseries_id": station["ts_id"], "station_id": properties.get("station_id"), "station_no": properties.get("station_no"), "raw_timeseries_sha256": raw_sha256, "coverage_first_year": coverage[0], "coverage_last_year": coverage[1], "coverage_observation_count": coverage[2], "generated_at": datetime.now(timezone.utc).isoformat(), "limitation_message": parameter.limitation, } with path.open("rb") as handle: response = session.post( f"{base_url}/api/v1/projects/{project_id}/datasets/upload", data={ "dataset_type": "vector", "source": "operator_official_import", "dataset_role": "reference", "source_name": "waterinfo", "reference_layer_name": parameter.reference_layer_name, "source_metadata_json": json.dumps(source_metadata, ensure_ascii=False), "provenance_metadata_json": json.dumps(provenance_metadata, ensure_ascii=False), "area_id": area_id, "temporal_series_key": key, "observed_at": observed_at, "valid_from": observed_at, "valid_to": f"{year}-12-31T23:59:59Z", "temporal_granularity": "year", "source_version": f"{station['ts_id']}:{year}", }, files={"file": (path.name, handle, "application/geo+json")}, timeout=timeout, ) return response_data(response) def main() -> int: args = parse_args() if args.from_year > args.to_year or args.min_observations < 2 or args.max_stations < 1: print(json.dumps({"status": "error", "message": "Invalid year, observation or station limits"}), file=sys.stderr) return 2 requested_keys = [value.strip() for value in args.parameters.split(",") if value.strip()] unsupported = [key for key in requested_keys if key not in PARAMETERS] if unsupported or not requested_keys: print(json.dumps({"status": "error", "message": f"Unsupported parameters: {unsupported}"}), file=sys.stderr) return 2 args.output_dir.mkdir(parents=True, exist_ok=True) base_url = args.base_url.rstrip("/") results: list[dict[str, Any]] = [] try: with requests.Session() as api_session: project_id, area_id, area_geometry, existing = locate_workspace(api_session, base_url, args) with build_session() as source_session: for parameter_key in requested_keys: parameter = PARAMETERS[parameter_key] station_layer, stations = discover_station_series( source_session, parameter, area_geometry, timeout=args.request_timeout, ) layer_path = args.output_dir / f"waterinfo_{parameter.key}_station_layer.json" layer_sha256 = write_json(layer_path, station_layer) if not stations: results.append( { "parameter": parameter.key, "timeseries_group_id": parameter.group_id, "status": "no_stations_in_area", "observation_count": 0, } ) if len(stations) > args.max_stations: raise RuntimeError( f"Waterinfo returned {len(stations)} in-area {parameter.key} stations; safety limit is {args.max_stations}" ) for station in stations: raw_path = args.output_dir / f"waterinfo_{parameter.key}_{station['ts_id']}_annual.json" if args.force or not raw_path.exists(): raw_payload, values = fetch_annual_values( source_session, station["ts_id"], from_year=args.from_year, to_year=args.to_year, timeout=args.request_timeout, ) raw_sha256 = write_json(raw_path, raw_payload) else: raw_payload = json.loads(raw_path.read_text(encoding="utf-8")) raw_sha256 = sha256_bytes(raw_path.read_bytes()) entries = raw_payload if isinstance(raw_payload, list) else [raw_payload] values = {} for entry in entries: for row in entry.get("data", []) if isinstance(entry, dict) else []: try: year = int(str(row[0])[:4]) value = float(row[1]) except (IndexError, TypeError, ValueError): continue if args.from_year <= year <= args.to_year and math.isfinite(value) and value > -9999: values[year] = value if len(values) < args.min_observations: results.append( { "parameter": parameter.key, "timeseries_id": station["ts_id"], "status": "insufficient_observations", "observation_count": len(values), } ) continue coverage = (min(values), max(values), len(values)) key = series_key(parameter, station) for year, value in sorted(values.items()): path = args.output_dir / f"{safe_slug(key)}_{year}.geojson" snapshot = build_snapshot(parameter, station, year, value) write_json(path, snapshot) existing_dataset = next( ( item for item in existing if item.get("temporal_series_key") == key and str(item.get("observed_at") or "").startswith(str(year)) ), None, ) if existing_dataset: results.append( { "parameter": parameter.key, "timeseries_id": station["ts_id"], "year": year, "dataset_id": existing_dataset["id"], "status": "existing", } ) elif args.fetch_only: results.append( { "parameter": parameter.key, "timeseries_id": station["ts_id"], "year": year, "path": str(path), "status": "prepared", } ) else: dataset = upload_snapshot( api_session, base_url=base_url, project_id=project_id, area_id=area_id, parameter=parameter, station=station, year=year, path=path, raw_sha256=raw_sha256, coverage=coverage, timeout=args.import_timeout, ) results.append( { "parameter": parameter.key, "timeseries_id": station["ts_id"], "year": year, "dataset_id": dataset["id"], "status": "imported", } ) manifest = { "schema_version": 1, "parameter": parameter.key, "timeseries_group_id": parameter.group_id, "area_name": args.area_name, "station_count": len(stations), "station_layer_path": str(layer_path), "station_layer_sha256": layer_sha256, "from_year": args.from_year, "to_year": args.to_year, "generated_at": datetime.now(timezone.utc).isoformat(), } write_json(args.output_dir / f"waterinfo_{parameter.key}_manifest.json", manifest, pretty=True) except (OSError, RuntimeError, requests.RequestException, ValueError, KeyError) as exc: print(json.dumps({"status": "error", "message": str(exc)}, ensure_ascii=False), file=sys.stderr) return 1 print( json.dumps( { "status": "ok", "project": args.project_name, "area": args.area_name, "results": results, }, ensure_ascii=False, indent=2, ) ) return 0 if __name__ == "__main__": sys.exit(main())