"""Collector for RID large-dam daily status (app.rid.go.th/reservoir). The Royal Irrigation Department reservoir app exposes an unauthenticated JSON API: ``POST https://app.rid.go.th/reservoir/api/dams`` with form field ``date=YYYY-MM-DD`` (empty = today) returns a daily snapshot of every large dam in Thailand — storage, inflow and outflow in MCM — with archive depth back to at least 2009. GET returns 404 ("Unknown method."); the POST body may be empty but must carry a Content-Length. The reservoir that matters for P.1 flood forecasting is Mae Ngat Somboon Chon (DAM_ID 200103), the only large dam upstream of Chiang Mai: during the Oct 2024 record flood it reached 113% of usable capacity with inflow spikes of ~19 MCM/day. All dams in the payload are stored (same request cost); filtering happens at feature-build time. See docs/DATA_SOURCES.md for the endpoint catalog. """ import datetime import logging import time from typing import Any, Dict, List, Optional import requests logger = logging.getLogger(__name__) RID_DAMS_URL = "https://app.rid.go.th/reservoir/api/dams" MAE_NGAT_DAM_ID = "200103" # Per-column NUMERIC capacity; source junk beyond these becomes NULL instead # of overflowing the insert and discarding the whole daily batch. _MEASURE_BOUNDS = { "storage_mcm": 1e8, "storage_pct": 1e6, "inflow_mcm": 1e8, "outflow_mcm": 1e8, "level_msl": 1e6, } def _bounded(value: Optional[float], limit: float) -> Optional[float]: if value is not None and abs(value) >= limit: return None return value def _to_float(value: Any) -> Optional[float]: """API numerics arrive as strings ('222.01'), ' - ' placeholders, or None.""" if value is None: return None if isinstance(value, str): value = value.replace(",", "").strip() if value in ("", "-"): return None try: return float(value) except (TypeError, ValueError): return None def parse_dam_records(payload: Dict) -> List[Dict]: """Flatten the regions/dams payload into per-dam daily rows.""" records = [] for region in payload.get("regions") or []: for dam in region.get("dams") or []: dam_id = dam.get("DAM_ID") try: date = datetime.date.fromisoformat(dam.get("DMD_Date") or "") except ValueError: continue if not dam_id: continue records.append( { "dam_id": dam_id, "region": region.get("region_name"), "name_th": dam.get("DAM_Name"), "latitude": _to_float(dam.get("DAM_Lat")), "longitude": _to_float(dam.get("DAM_Lon")), "capacity_max_mcm": _to_float(dam.get("DAM_QMax")), "capacity_normal_mcm": _to_float(dam.get("DAM_QStore")), "date": date, "storage_mcm": _to_float(dam.get("DMD_QUse")), "storage_pct": _to_float(dam.get("PERCENT_DMD_QUse")), "inflow_mcm": _to_float(dam.get("DMD_Inflow")), "outflow_mcm": _to_float(dam.get("DMD_Outflow")), "level_msl": _to_float(dam.get("DMD_Q")), } ) return records class RidReservoirClient: """HTTP client for the RID reservoir daily-status API.""" def __init__( self, url: str = RID_DAMS_URL, session: Optional[requests.Session] = None, timeout: int = 60, ): self.url = url self.session = session or requests.Session() self.timeout = timeout def fetch_day(self, date: Optional[datetime.date] = None) -> List[Dict]: """Daily rows for every large dam on `date` (None = today).""" data = {"date": date.isoformat()} if date else {"date": ""} response = self.session.post( self.url, data=data, headers={"X-Requested-With": "XMLHttpRequest"}, timeout=self.timeout, ) response.raise_for_status() return parse_dam_records(response.json()) class RidReservoirStore: """SQL persistence for dam metadata + daily measurements. Shares the app's relational database; writes rid_dams (metadata) and rid_reservoir_daily keyed (dam_id, date) — the composite natural PK keeps the table TimescaleDB-hypertable compatible. """ def __init__(self, connection_string: str, db_type: str): self.db_type = db_type.lower() if self.db_type not in ("sqlite", "postgresql", "mysql"): raise ValueError( f"Reservoir collection requires a SQL database, got '{db_type}'" ) self.connection_string = connection_string self.engine = None def connect(self) -> bool: try: from sqlalchemy import create_engine self.engine = create_engine(self.connection_string, pool_pre_ping=True) self._create_tables() return True except Exception as e: logger.error(f"RidReservoirStore failed to connect: {e}") self.engine = None return False def _create_tables(self): from sqlalchemy import text ddl = [ """ CREATE TABLE IF NOT EXISTS rid_dams ( dam_id VARCHAR(10) PRIMARY KEY, region VARCHAR(40), name_th VARCHAR(255), latitude NUMERIC(10,6), longitude NUMERIC(10,6), capacity_max_mcm NUMERIC(10,2), capacity_normal_mcm NUMERIC(10,2), updated_at TIMESTAMP ) """, """ CREATE TABLE IF NOT EXISTS rid_reservoir_daily ( dam_id VARCHAR(10) NOT NULL, date DATE NOT NULL, storage_mcm NUMERIC(10,2), storage_pct NUMERIC(8,2), inflow_mcm NUMERIC(10,2), outflow_mcm NUMERIC(10,2), level_msl NUMERIC(8,2), created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (dam_id, date) ) """, ] if self.db_type != "mysql": ddl.append( "CREATE INDEX IF NOT EXISTS idx_rid_reservoir_date " "ON rid_reservoir_daily(date)" ) with self.engine.begin() as conn: for statement in ddl: conn.execute(text(statement)) # Widen storage_pct on tables created before 2026-08-13: the source # publishes junk percents (dam 100602 reports 87798%) that overflowed # NUMERIC(6,2) and discarded whole daily batches. if self.db_type == "postgresql": migrations = ( "ALTER TABLE rid_reservoir_daily " "ALTER COLUMN storage_pct TYPE NUMERIC(8,2)", ) elif self.db_type == "mysql": migrations = ( "ALTER TABLE rid_reservoir_daily MODIFY storage_pct NUMERIC(8,2)", ) else: # sqlite: NUMERIC is affinity only, nothing to widen migrations = () for statement in migrations: try: with self.engine.begin() as conn: conn.execute(text(statement)) except Exception as e: logger.warning(f"rid_reservoir_daily migration skipped: {e}") def _upsert(self, table: str, key_cols: List[str], value_cols: List[str]) -> str: cols = key_cols + value_cols col_list = ", ".join(cols) params = ", ".join(f":{c}" for c in cols) if self.db_type == "sqlite": return f"INSERT OR REPLACE INTO {table} ({col_list}) VALUES ({params})" if self.db_type == "postgresql": updates = ", ".join(f"{c} = EXCLUDED.{c}" for c in value_cols) conflict = ", ".join(key_cols) return ( f"INSERT INTO {table} ({col_list}) VALUES ({params}) " f"ON CONFLICT ({conflict}) DO UPDATE SET {updates}" ) updates = ", ".join(f"{c} = VALUES({c})" for c in value_cols) return ( f"INSERT INTO {table} ({col_list}) VALUES ({params}) " f"ON DUPLICATE KEY UPDATE {updates}" ) def save(self, records: List[Dict]) -> int: """Upsert one day's dam rows (metadata + measurements); idempotent.""" if not records: return 0 if not self.engine and not self.connect(): return 0 from sqlalchemy import text dam_cols = [ "region", "name_th", "latitude", "longitude", "capacity_max_mcm", "capacity_normal_mcm", ] measure_cols = [ "storage_mcm", "storage_pct", "inflow_mcm", "outflow_mcm", "level_msl", ] dam_sql = self._upsert("rid_dams", ["dam_id"], dam_cols + ["updated_at"]) measure_sql = self._upsert( "rid_reservoir_daily", ["dam_id", "date"], measure_cols ) now = datetime.datetime.now() dams = {} measurements = [] for record in records: dam_row = {c: record.get(c) for c in dam_cols} dam_row.update({"dam_id": record["dam_id"], "updated_at": now}) dams[record["dam_id"]] = dam_row measure_row = { c: _bounded(record.get(c), _MEASURE_BOUNDS[c]) for c in measure_cols } measure_row.update( {"dam_id": record["dam_id"], "date": record["date"]} ) measurements.append(measure_row) try: with self.engine.begin() as conn: conn.execute(text(dam_sql), list(dams.values())) conn.execute(text(measure_sql), measurements) return len(measurements) except Exception as e: logger.error(f"RidReservoirStore save failed: {e}") return 0 def present_dates( self, start: datetime.date, end: datetime.date ) -> "set[datetime.date]": """Dates in [start, end] that already have rows, for backfill skipping.""" if not self.engine and not self.connect(): return set() from sqlalchemy import text with self.engine.begin() as conn: values = conn.execute( text( "SELECT DISTINCT date FROM rid_reservoir_daily " "WHERE date >= :start AND date <= :end" ), {"start": start, "end": end}, ).fetchall() dates = set() for (value,) in values: if isinstance(value, str): # sqlite returns ISO strings value = datetime.date.fromisoformat(value[:10]) if isinstance(value, datetime.datetime): value = value.date() dates.add(value) return dates class RidReservoirCollector: """Fetch + persist the daily dam snapshot (today and yesterday).""" def __init__(self, db_config: Dict, client: Optional[RidReservoirClient] = None): self.client = client or RidReservoirClient() self.store = RidReservoirStore( connection_string=db_config["connection_string"], db_type=db_config["type"], ) def run_cycle(self) -> int: """Collect today's snapshot plus yesterday's (late daily revisions).""" saved = 0 today = datetime.date.today() for date in (None, today - datetime.timedelta(days=1)): try: saved += self.store.save(self.client.fetch_day(date)) except Exception as e: logger.error(f"RID reservoir collection failed for {date}: {e}") logger.info(f"RID reservoir collection: {saved} dam-day rows saved") return saved def backfill( store: RidReservoirStore, start: datetime.date, end: Optional[datetime.date] = None, client: Optional[RidReservoirClient] = None, throttle_seconds: float = 0.4, ) -> int: """Fetch every MISSING day in [start, end]; one polite request per day. Only dates absent from rid_reservoir_daily are requested, so a rerun repairs holes left by transient failures instead of resuming past them (the hourly collector writes today's rows immediately, which makes any newest-row cursor useless as a resume point). A save that persists nothing counts as a failure too — a broken DB must not burn thousands of requests against the RID API. """ client = client or RidReservoirClient() end = end or datetime.date.today() if not store.engine and not store.connect(): logger.error("backfill aborted: database connection failed") return 0 span = [ start + datetime.timedelta(days=i) for i in range((end - start).days + 1) ] present = store.present_dates(start, end) targets = [d for d in span if d not in present] logger.info( f"backfill: {len(targets)} of {len(span)} days missing in [{start}, {end}]" ) total = 0 failures = 0 for i, date in enumerate(targets): try: records = client.fetch_day(date) saved = store.save(records) if records and not saved: raise RuntimeError("database save persisted 0 rows") total += saved failures = 0 except Exception as e: failures += 1 logger.warning(f"backfill {date} failed ({failures} in a row): {e}") if failures >= 5: logger.error("5 consecutive failures — aborting backfill") break if i % 100 == 0: logger.info(f"backfill progress: {date} ({total} rows)") time.sleep(throttle_seconds) return total def create_collector_from_config() -> Optional[RidReservoirCollector]: """Build a collector from app Config; None when disabled or non-SQL DB.""" from .config import Config if not Config.ENABLE_RESERVOIR_COLLECTION: return None db_config = Config.get_database_config() if db_config["type"] not in ("sqlite", "postgresql", "mysql"): logger.warning( f"Reservoir collection skipped: DB_TYPE '{db_config['type']}' is not SQL" ) return None return RidReservoirCollector(db_config)