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.
This commit is contained in:
@@ -111,9 +111,11 @@ class InfluxDBAdapter(DatabaseAdapter):
|
|||||||
"time": measurement["timestamp"].isoformat(),
|
"time": measurement["timestamp"].isoformat(),
|
||||||
"fields": {
|
"fields": {
|
||||||
"water_level": float(measurement["water_level"]),
|
"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"])
|
"discharge_percent": float(measurement["discharge_percent"])
|
||||||
if measurement["discharge_percent"]
|
if measurement.get("discharge_percent")
|
||||||
else None,
|
else None,
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|||||||
+21
-33
@@ -59,7 +59,7 @@ class MeasurementResponse(BaseModel):
|
|||||||
station_name_en: str
|
station_name_en: str
|
||||||
station_name_th: str
|
station_name_th: str
|
||||||
water_level: float
|
water_level: float
|
||||||
discharge: float
|
discharge: Optional[float] = None
|
||||||
discharge_percent: Optional[float] = None
|
discharge_percent: Optional[float] = None
|
||||||
status: str = "active"
|
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])
|
@app.get("/measurements/latest", response_model=List[MeasurementResponse])
|
||||||
async def get_latest_measurements(limit: int = 100):
|
async def get_latest_measurements(limit: int = 100):
|
||||||
"""Get latest measurements from all stations"""
|
"""Get latest measurements from all stations"""
|
||||||
@@ -490,22 +508,7 @@ async def get_latest_measurements(limit: int = 100):
|
|||||||
try:
|
try:
|
||||||
measurements = scraper.get_latest_data(limit=limit)
|
measurements = scraper.get_latest_data(limit=limit)
|
||||||
|
|
||||||
response = []
|
return [_to_measurement_response(m) for m in measurements]
|
||||||
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
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"Error fetching latest measurements: {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
|
# Limit results
|
||||||
measurements = measurements[:limit]
|
measurements = measurements[:limit]
|
||||||
|
|
||||||
response = []
|
return [_to_measurement_response(m) for m in measurements]
|
||||||
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
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"Error fetching station measurements: {e}")
|
logger.error(f"Error fetching station measurements: {e}")
|
||||||
|
|||||||
Reference in New Issue
Block a user