diff --git a/web/services/city_api.py b/web/services/city_api.py index ea0650f9..0d6d24af 100644 --- a/web/services/city_api.py +++ b/web/services/city_api.py @@ -18,8 +18,10 @@ import web.routes as legacy_routes from web.analysis_service import _runway_history_temp_for_city from web.services.canonical_temperature import build_city_weather_from_canonical from web.services.latest_observation_overlay import ( + overlay_latest_amos_observation, overlay_latest_amsc_observation, overlay_latest_jma_amedas_observation, + overlay_latest_mgm_observation, parse_observation_epoch, ) from web.services.request_timing import ServerTimingRecorder @@ -671,6 +673,20 @@ async def _get_city_chart_data(city: str, *, force_refresh: bool) -> Dict[str, A fn=overlay_latest_jma_amedas_observation, args=(legacy_routes._weather, city, payload), ) + payload = await _run_optional_city_chart_overlay( + city=city, + overlay_name="amos_latest_raw", + payload=payload, + fn=overlay_latest_amos_observation, + args=(legacy_routes._CACHE_DB, city, payload), + ) + payload = await _run_optional_city_chart_overlay( + city=city, + overlay_name="mgm_latest_raw", + payload=payload, + fn=overlay_latest_mgm_observation, + args=(legacy_routes._CACHE_DB, city, payload), + ) return await _run_optional_city_chart_overlay( city=city, overlay_name="wunderground_current", @@ -707,6 +723,20 @@ async def _get_city_chart_data(city: str, *, force_refresh: bool) -> Dict[str, A fn=overlay_latest_jma_amedas_observation, args=(legacy_routes._weather, city, payload), ) + payload = await _run_optional_city_chart_overlay( + city=city, + overlay_name="amos_latest_raw", + payload=payload, + fn=overlay_latest_amos_observation, + args=(legacy_routes._CACHE_DB, city, payload), + ) + payload = await _run_optional_city_chart_overlay( + city=city, + overlay_name="mgm_latest_raw", + payload=payload, + fn=overlay_latest_mgm_observation, + args=(legacy_routes._CACHE_DB, city, payload), + ) return await _run_optional_city_chart_overlay( city=city, overlay_name="wunderground_current", diff --git a/web/services/latest_observation_overlay.py b/web/services/latest_observation_overlay.py index d9b295bf..4c581fa3 100644 --- a/web/services/latest_observation_overlay.py +++ b/web/services/latest_observation_overlay.py @@ -708,3 +708,163 @@ def overlay_latest_amsc_observation( changed = True return next_payload if changed else payload + + +# ═══════════════════════════════════════════════════════════════════════════════ +# AMOS (Korean runway sensors — Seoul, Busan) +# ═══════════════════════════════════════════════════════════════════════════════ + + +def _latest_amos_row(db: Any, city: str) -> tuple: + getter = getattr(db, "get_latest_raw_observation", None) + if not callable(getter): + return None, None + try: + row = getter("amos", city) + except Exception as exc: + logger.debug("latest AMOS raw overlay read failed city={}: {}", city, exc) + return None, None + if not isinstance(row, dict): + return None, None + raw_payload = row.get("payload") + if isinstance(raw_payload, dict) and raw_payload and raw_payload.get("temp") is not None: + return row, raw_payload + return None, None + + +def _raw_amos_epoch(row, raw_payload): + payload_values = ( + raw_payload.get("observation_time"), + raw_payload.get("observed_at"), + ) + parsed = [ + epoch + for epoch in (parse_observation_epoch(value) for value in payload_values) + if epoch is not None + ] + if parsed: + return max(parsed) + return parse_observation_epoch(row.get("observed_at")) + + +def overlay_latest_amos_observation(db, city, payload): + normalized_city = str(city or payload.get("name") or payload.get("city") or "").strip().lower() + if not normalized_city or not isinstance(payload, dict) or not payload: + return payload + row, raw_payload = _latest_amos_row(db, normalized_city) + if not isinstance(row, dict) or not isinstance(raw_payload, dict): + return payload + raw_epoch = _raw_amos_epoch(row, raw_payload) + if raw_epoch is None: + return payload + update = _raw_observation_update(normalized_city, row, raw_payload) + if not update: + return payload + + next_payload = deepcopy(payload) + changed = False + amos = next_payload.get("amos") + amos_epoch = _block_epoch(amos) + if amos_epoch is None or raw_epoch > amos_epoch: + next_payload["amos"] = dict(raw_payload) + changed = True + + for key in ("current", "airport_primary", "airport_current"): + changed = _merge_observation_block(next_payload, key, update, raw_epoch) or changed + + canonical = next_payload.get("canonical_temperature") + canonical_epoch = _block_epoch(canonical) + if canonical_epoch is None or raw_epoch > canonical_epoch: + canonical_payload = build_canonical_temperature( + normalized_city, + { + "name": normalized_city, + "temp_symbol": next_payload.get("temp_symbol") or "\u00b0C", + "updated_at": row.get("fetched_at"), + "current": update, + }, + fetched_at=str(row.get("fetched_at") or ""), + ) + if canonical_payload: + next_payload["canonical_temperature"] = canonical_payload + changed = True + + return next_payload if changed else payload + + +# ═══════════════════════════════════════════════════════════════════════════════ +# MGM (Turkish State Meteorological Service — Ankara, Istanbul) +# ═══════════════════════════════════════════════════════════════════════════════ + + +def _latest_mgm_row(db, city): + getter = getattr(db, "get_latest_raw_observation", None) + if not callable(getter): + return None, None + try: + row = getter("mgm", city) + except Exception as exc: + logger.debug("latest MGM raw overlay read failed city={}: {}", city, exc) + return None, None + if not isinstance(row, dict): + return None, None + raw_payload = row.get("payload") + if isinstance(raw_payload, dict) and raw_payload: + temp = raw_payload.get("temp") + if temp is not None: + return row, raw_payload + return None, None + + +def _raw_mgm_epoch(row, raw_payload): + payload_values = ( + raw_payload.get("obs_time"), + ) + parsed = [ + epoch + for epoch in (parse_observation_epoch(value) for value in payload_values) + if epoch is not None + ] + if parsed: + return max(parsed) + return parse_observation_epoch(row.get("observed_at")) + + +def overlay_latest_mgm_observation(db, city, payload): + normalized_city = str(city or payload.get("name") or payload.get("city") or "").strip().lower() + if not normalized_city or not isinstance(payload, dict) or not payload: + return payload + row, raw_payload = _latest_mgm_row(db, normalized_city) + if not isinstance(row, dict) or not isinstance(raw_payload, dict): + return payload + raw_epoch = _raw_mgm_epoch(row, raw_payload) + if raw_epoch is None: + return payload + update = _raw_observation_update(normalized_city, row, raw_payload) + if not update: + return payload + + next_payload = deepcopy(payload) + changed = False + + for key in ("current", "airport_primary", "airport_current"): + changed = _merge_observation_block(next_payload, key, update, raw_epoch) or changed + + canonical = next_payload.get("canonical_temperature") + canonical_epoch = _block_epoch(canonical) + if canonical_epoch is None or raw_epoch > canonical_epoch: + canonical_payload = build_canonical_temperature( + normalized_city, + { + "name": normalized_city, + "temp_symbol": next_payload.get("temp_symbol") or "\u00b0C", + "updated_at": row.get("fetched_at"), + "current": update, + }, + fetched_at=str(row.get("fetched_at") or ""), + ) + if canonical_payload: + next_payload["canonical_temperature"] = canonical_payload + changed = True + + return next_payload if changed else payload