From 6e78225d00c788caa17dc0f87597fbf78a216b5a Mon Sep 17 00:00:00 2001 From: grabowski Date: Wed, 22 Jul 2026 12:34:46 +0700 Subject: [PATCH] Fix optional-discharge crashes and dedupe measurement mapping - MeasurementResponse.discharge is now Optional[float]; measurements with a null discharge no longer raise a Pydantic ValidationError (HTTP 500) on the /measurements/latest and /measurements/station endpoints. - InfluxDB save_measurements guards float(discharge) against None instead of crashing with TypeError. - Extract the duplicated measurement->response mapping into a single _to_measurement_response helper used by both measurement endpoints. --- src/database_adapters.py | 6 +++-- src/web_api.py | 54 ++++++++++++++++------------------------ 2 files changed, 25 insertions(+), 35 deletions(-) diff --git a/src/database_adapters.py b/src/database_adapters.py index 588a2c5..a7540c0 100644 --- a/src/database_adapters.py +++ b/src/database_adapters.py @@ -111,9 +111,11 @@ class InfluxDBAdapter(DatabaseAdapter): "time": measurement["timestamp"].isoformat(), "fields": { "water_level": float(measurement["water_level"]), - "discharge": float(measurement["discharge"]), + "discharge": float(measurement["discharge"]) + if measurement.get("discharge") is not None + else None, "discharge_percent": float(measurement["discharge_percent"]) - if measurement["discharge_percent"] + if measurement.get("discharge_percent") else None, }, } diff --git a/src/web_api.py b/src/web_api.py index 6944112..82c1a23 100644 --- a/src/web_api.py +++ b/src/web_api.py @@ -59,7 +59,7 @@ class MeasurementResponse(BaseModel): station_name_en: str station_name_th: str water_level: float - discharge: float + discharge: Optional[float] = None discharge_percent: Optional[float] = None status: str = "active" @@ -478,6 +478,24 @@ async def get_station(station_id: int): ) +def _to_measurement_response(measurement: Dict[str, Any]) -> MeasurementResponse: + """Map a raw measurement dict from a DB adapter to the API response model. + + ``discharge`` is optional in the data (some stations report only level), so + it is read with ``.get`` rather than assumed present. + """ + return MeasurementResponse( + timestamp=measurement["timestamp"], + station_code=measurement["station_code"], + station_name_en=measurement["station_name_en"], + station_name_th=measurement["station_name_th"], + water_level=measurement["water_level"], + discharge=measurement.get("discharge"), + discharge_percent=measurement.get("discharge_percent"), + status=measurement.get("status", "active"), + ) + + @app.get("/measurements/latest", response_model=List[MeasurementResponse]) async def get_latest_measurements(limit: int = 100): """Get latest measurements from all stations""" @@ -490,22 +508,7 @@ async def get_latest_measurements(limit: int = 100): try: measurements = scraper.get_latest_data(limit=limit) - response = [] - for measurement in measurements: - response.append( - MeasurementResponse( - timestamp=measurement["timestamp"], - station_code=measurement["station_code"], - station_name_en=measurement["station_name_en"], - station_name_th=measurement["station_name_th"], - water_level=measurement["water_level"], - discharge=measurement["discharge"], - discharge_percent=measurement.get("discharge_percent"), - status=measurement.get("status", "active"), - ) - ) - - return response + return [_to_measurement_response(m) for m in measurements] except Exception as e: logger.error(f"Error fetching latest measurements: {e}") @@ -533,22 +536,7 @@ async def get_station_measurements(station_code: str, hours: int = 24, limit: in # Limit results measurements = measurements[:limit] - response = [] - for measurement in measurements: - response.append( - MeasurementResponse( - timestamp=measurement["timestamp"], - station_code=measurement["station_code"], - station_name_en=measurement["station_name_en"], - station_name_th=measurement["station_name_th"], - water_level=measurement["water_level"], - discharge=measurement["discharge"], - discharge_percent=measurement.get("discharge_percent"), - status=measurement.get("status", "active"), - ) - ) - - return response + return [_to_measurement_response(m) for m in measurements] except Exception as e: logger.error(f"Error fetching station measurements: {e}")