From 07e28481e8b3f52ea2e1cc46c4d017137837354a Mon Sep 17 00:00:00 2001 From: "2569718930@qq.com" <2569718930@qq.com> Date: Mon, 15 Jun 2026 20:40:57 +0800 Subject: [PATCH] Require recent runway history coverage --- tests/test_latest_observation_overlay.py | 17 +++--- web/services/latest_observation_overlay.py | 65 +++++++++++++++++++--- 2 files changed, 66 insertions(+), 16 deletions(-) diff --git a/tests/test_latest_observation_overlay.py b/tests/test_latest_observation_overlay.py index 6b8f4890..103812c4 100644 --- a/tests/test_latest_observation_overlay.py +++ b/tests/test_latest_observation_overlay.py @@ -162,17 +162,23 @@ def test_overlay_builds_amsc_runway_history_from_raw_store(): def test_overlay_skips_raw_store_when_runway_history_already_has_points(): + existing_points = [ + {"time": f"2026-06-15T{hour:02d}:00:00+00:00", "temp": 25.0 + hour / 10} + for hour in range(0, 15) + for _ in range(2) + ] + class FakeDB: def get_latest_raw_observation(self, source, city): assert (source, city) == ("amsc_awos", "chengdu") return { - "observed_at": "2026-06-15T11:08:00+00:00", + "observed_at": "2026-06-15T14:30:00+00:00", "station_code": "ZUUU", "payload": { "source": "amsc_awos", "icao": "ZUUU", "temp_c": 25.4, - "observation_time": "2026-06-15T11:08:00+00:00", + "observation_time": "2026-06-15T14:30:00+00:00", "runway_obs": { "point_temperatures": [ { @@ -194,16 +200,13 @@ def test_overlay_skips_raw_store_when_runway_history_already_has_points(): "name": "chengdu", "temp_symbol": "°C", "runway_plate_history": { - "02L/20R": [ - {"time": "2026-06-15T11:04:00+00:00", "temp": 25.6}, - {"time": "2026-06-15T11:06:00+00:00", "temp": 25.5}, - ], + "02L/20R": existing_points, }, }, ) assert result["runway_plate_history"]["02L/20R"][-1] == { - "time": "2026-06-15T11:08:00+00:00", + "time": "2026-06-15T14:30:00+00:00", "temp": 25.4, } diff --git a/web/services/latest_observation_overlay.py b/web/services/latest_observation_overlay.py index 29e42d6d..3562e23a 100644 --- a/web/services/latest_observation_overlay.py +++ b/web/services/latest_observation_overlay.py @@ -2,12 +2,19 @@ from __future__ import annotations from copy import deepcopy from datetime import datetime, timezone +import time from typing import Any, Optional from loguru import logger from web.services.canonical_temperature import build_canonical_temperature +_RAW_AMSC_RUNWAY_HISTORY_CACHE: dict[tuple[str, bool], tuple[float, dict[str, list[dict[str, Any]]]]] = {} +_RAW_AMSC_RUNWAY_HISTORY_CACHE_TTL_SEC = 60.0 +_RAW_AMSC_RUNWAY_HISTORY_RECENT_WINDOW_SEC = 24 * 60 * 60 +_RAW_AMSC_RUNWAY_HISTORY_MIN_POINTS = 30 +_RAW_AMSC_RUNWAY_HISTORY_MIN_SPAN_SEC = 12 * 60 * 60 + def parse_observation_epoch(value: Any) -> Optional[int]: if value is None or value == "": @@ -208,6 +215,33 @@ def _runway_history_temp(point: dict[str, Any]) -> Optional[float]: return None +def _runway_history_has_recent_coverage( + history: Any, + latest_epoch: int, +) -> bool: + if not isinstance(history, dict): + return False + window_start = latest_epoch - _RAW_AMSC_RUNWAY_HISTORY_RECENT_WINDOW_SEC + span_start = latest_epoch - _RAW_AMSC_RUNWAY_HISTORY_MIN_SPAN_SEC + for points in history.values(): + if not isinstance(points, list): + continue + epochs = [ + epoch + for epoch in ( + parse_observation_epoch( + point.get("time") or point.get("timestamp") or point.get("observed_at") + ) + for point in points + if isinstance(point, dict) + ) + if epoch is not None and window_start <= epoch <= latest_epoch + 300 + ] + if len(epochs) >= _RAW_AMSC_RUNWAY_HISTORY_MIN_POINTS and min(epochs) <= span_start: + return True + return False + + def _append_latest_amsc_runway_history( db: Any, city: str, @@ -312,16 +346,19 @@ def _append_amsc_runway_history_from_raw_store( db: Any, city: str, payload: dict[str, Any], + latest_epoch: int, ) -> bool: existing_history = payload.get("runway_plate_history") - if isinstance(existing_history, dict): - existing_points = [ - len(points) - for points in existing_history.values() - if isinstance(points, list) - ] - if max(existing_points or [0]) > 1: - return False + if _runway_history_has_recent_coverage(existing_history, latest_epoch): + return False + + use_fahrenheit = "F" in str(payload.get("temp_symbol") or "").upper() + cache_key = (city, use_fahrenheit) + cached = _RAW_AMSC_RUNWAY_HISTORY_CACHE.get(cache_key) + now = time.monotonic() + if cached and now - cached[0] <= _RAW_AMSC_RUNWAY_HISTORY_CACHE_TTL_SEC: + payload["runway_plate_history"] = deepcopy(cached[1]) + return True lister = getattr(db, "list_raw_observation_history", None) if not callable(lister): @@ -348,6 +385,11 @@ def _append_amsc_runway_history_from_raw_store( persist=False, copy_history=False, ) or changed + if changed and isinstance(payload.get("runway_plate_history"), dict): + _RAW_AMSC_RUNWAY_HISTORY_CACHE[cache_key] = ( + now, + deepcopy(payload["runway_plate_history"]), + ) return changed @@ -380,7 +422,12 @@ def overlay_latest_amsc_observation( for key in ("current", "airport_primary", "airport_current"): changed = _merge_observation_block(next_payload, key, update, raw_epoch) or changed - changed = _append_amsc_runway_history_from_raw_store(db, normalized_city, next_payload) or changed + changed = _append_amsc_runway_history_from_raw_store( + db, + normalized_city, + next_payload, + raw_epoch, + ) or changed changed = _append_latest_amsc_runway_history(db, normalized_city, next_payload, row, raw_payload) or changed canonical = next_payload.get("canonical_temperature")