Files
geointel/scripts/provision_waterinfo_station_history.py
Jens faeb58ef6d
GeoIntel release gates / Compile, test, contracts and builds (push) Successful in 1m49s
GeoIntel release gates / Python and npm vulnerability policy (push) Successful in 21s
GeoIntel release gates / Production AI image, SBOM and container scan (push) Successful in 5m39s
GeoIntel release gates / Deploy exact gated revision to Unraid (push) Failing after 58m43s
Initial public release
2026-08-31 21:56:53 +02:00

589 lines
25 KiB
Python

"""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 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())