Fix Ankara MGM observation freshness

This commit is contained in:
2569718930@qq.com
2026-06-15 22:47:13 +08:00
parent c0b050d2a7
commit fd9d768c6d
9 changed files with 300 additions and 5 deletions
+14 -2
View File
@@ -14,6 +14,7 @@ from web.realtime_patch_schema import normalize_observation_patch
_SETTLEMENT_SOURCE_ADAPTERS = {
"hko": {"hko_obs", "cowin_obs"},
"cwa": {"cwa"},
"mgm": {"mgm"},
"noaa": {"madis_hfmetar", "metar", "noaa"},
"wunderground": {"amsc_awos", "amos", "madis_hfmetar", "metar", "wunderground"},
}
@@ -23,6 +24,7 @@ _SOURCE_WEIGHTS = {
"amos": 700,
"hko_obs": 680,
"cowin_obs": 660,
"mgm": 640,
"madis_hfmetar": 500,
"metar": 460,
}
@@ -35,6 +37,11 @@ _FRESHNESS_WEIGHTS = {
"stale": 5,
}
_TURKISH_MGM_STATION_CODES = {
"ankara": "17128",
"istanbul": "17058",
}
def _normalized_city(city: Any) -> str:
return str(city or "").strip().lower()
@@ -116,8 +123,13 @@ def _candidate_canonical(city: str, row: dict[str, Any]) -> Optional[dict[str, A
def _score(city: str, row: dict[str, Any], canonical: dict[str, Any]) -> tuple[int, float]:
meta = CITY_REGISTRY.get(city) or {}
station_code = _station_code(row.get("station_code"))
settlement_station_code = _station_code(meta.get("settlement_station_code") or meta.get("icao"))
settlement_source = _normalized_source(meta.get("settlement_source"))
is_turkish_mgm_city = city in _TURKISH_MGM_STATION_CODES
settlement_station_code = _station_code(
_TURKISH_MGM_STATION_CODES.get(city)
if is_turkish_mgm_city
else (meta.get("settlement_station_code") or meta.get("icao"))
)
settlement_source = "mgm" if is_turkish_mgm_city else _normalized_source(meta.get("settlement_source"))
expected_sources = _SETTLEMENT_SOURCE_ADAPTERS.get(settlement_source, {settlement_source})
station_name = str(row.get("station_name") or "").strip().lower()
candidates = {
+2
View File
@@ -218,6 +218,8 @@ def build_city_weather_from_canonical(city: str, canonical: Dict[str, Any]) -> O
"max_so_far": _first_float(canonical.get("max_so_far"), value),
"source_code": source,
"source_label": source_label,
"station_code": canonical.get("station_code"),
"station_name": canonical.get("station_name"),
"obs_time": observed_at or observed_at_local,
"freshness": freshness,
}
+72 -2
View File
@@ -16,7 +16,10 @@ from loguru import logger
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_amsc_observation
from web.services.latest_observation_overlay import (
overlay_latest_amsc_observation,
parse_observation_epoch,
)
from web.services.request_timing import ServerTimingRecorder
_RECENT_DEB_CACHE: Optional[Dict[str, Dict[str, object]]] = None
@@ -297,11 +300,12 @@ def _start_city_cache_stale_refresh(
async def _overlay_cached_wunderground(city: str, payload: Dict[str, Any]) -> Dict[str, Any]:
latest_payload = await _overlay_cached_canonical_observation(city, payload)
latest_payload = await run_in_threadpool(
overlay_latest_amsc_observation,
legacy_routes._CACHE_DB,
city,
payload,
latest_payload,
)
return await run_in_threadpool(
legacy_routes._overlay_latest_wunderground_current,
@@ -310,6 +314,70 @@ async def _overlay_cached_wunderground(city: str, payload: Dict[str, Any]) -> Di
)
def _observation_block_epoch(block: Any) -> Optional[int]:
if not isinstance(block, dict):
return None
freshness = block.get("freshness") if isinstance(block.get("freshness"), dict) else {}
values = (
block.get("observed_at"),
block.get("observation_time"),
block.get("obs_time"),
freshness.get("observed_at"),
)
epochs = [epoch for epoch in (parse_observation_epoch(value) for value in values) if epoch is not None]
return max(epochs) if epochs else None
def _payload_observation_epoch(payload: Dict[str, Any]) -> Optional[int]:
if not isinstance(payload, dict):
return None
blocks = (
payload.get("current"),
payload.get("airport_primary"),
payload.get("airport_current"),
)
epochs = [epoch for epoch in (_observation_block_epoch(block) for block in blocks) if epoch is not None]
return max(epochs) if epochs else None
def _merge_latest_observation_payload(
payload: Dict[str, Any],
latest_payload: Dict[str, Any],
) -> Dict[str, Any]:
if not isinstance(payload, dict) or not isinstance(latest_payload, dict):
return payload
latest_epoch = _payload_observation_epoch(latest_payload)
if latest_epoch is None:
return payload
current_epoch = _payload_observation_epoch(payload)
if current_epoch is not None and current_epoch >= latest_epoch:
return payload
next_payload = deepcopy(payload)
for key in ("current", "airport_primary", "airport_current"):
latest_block = latest_payload.get(key)
if not isinstance(latest_block, dict) or not latest_block:
continue
base_block = next_payload.get(key) if isinstance(next_payload.get(key), dict) else {}
next_payload[key] = {**base_block, **latest_block}
if isinstance(latest_payload.get("canonical_temperature"), dict):
next_payload["canonical_temperature"] = latest_payload["canonical_temperature"]
if latest_payload.get("updated_at"):
next_payload["updated_at"] = latest_payload.get("updated_at")
if latest_payload.get("temp_symbol"):
next_payload["temp_symbol"] = latest_payload.get("temp_symbol")
return next_payload
async def _overlay_cached_canonical_observation(city: str, payload: Dict[str, Any]) -> Dict[str, Any]:
if not isinstance(payload, dict) or not payload:
return payload
canonical_payload = await _get_canonical_city_payload(city, detail_depth=str(payload.get("detail_depth") or "full"))
if not canonical_payload:
return payload
return _merge_latest_observation_payload(payload, canonical_payload)
def _overlay_cached_runway_history_from_db(city: str, payload: Dict[str, Any]) -> Dict[str, Any]:
if not isinstance(payload, dict) or not payload:
return payload
@@ -405,6 +473,7 @@ async def _get_city_full_data(city: str, *, force_refresh: bool) -> Dict[str, An
async def _get_city_chart_data(city: str, *, force_refresh: bool) -> Dict[str, Any]:
if force_refresh:
payload = await _get_city_full_data(city, force_refresh=True)
payload = await _overlay_cached_canonical_observation(city, payload)
payload = await _run_optional_city_chart_overlay(
city=city,
overlay_name="runway_history",
@@ -433,6 +502,7 @@ async def _get_city_chart_data(city: str, *, force_refresh: bool) -> Dict[str, A
if payload:
if not legacy_routes._city_cache_is_fresh(cached_entry, legacy_routes.CITY_FULL_CACHE_TTL_SEC):
_start_city_full_stale_refresh(city)
payload = await _overlay_cached_canonical_observation(city, payload)
payload = await _run_optional_city_chart_overlay(
city=city,
overlay_name="runway_history",
+41 -1
View File
@@ -36,6 +36,12 @@ _ATTACH_METHODS: dict[str, str] = {
"madis_hfmetar": "_attach_madis_hfmetar_data",
"hko_obs": "_attach_hko_obs_official_nearby",
"cowin_obs": "_attach_cowin_official_nearby",
"mgm": "_attach_turkish_mgm_data",
}
_TURKISH_MGM_STATION_CODES = {
"ankara": "17128",
"istanbul": "17058",
}
@@ -115,6 +121,31 @@ def _record_from_row(
)
def _enrich_mgm_results(results: dict[str, Any], city: str) -> None:
mgm = results.get("mgm")
if not isinstance(mgm, dict):
return
current = mgm.get("current") if isinstance(mgm.get("current"), dict) else {}
station_code = str(
mgm.get("station_code")
or mgm.get("istNo")
or _TURKISH_MGM_STATION_CODES.get(city)
or ""
).strip()
station_name = str(
mgm.get("station_name")
or current.get("station_name")
or current.get("station_label")
or ""
).strip()
mgm.setdefault("source", "mgm")
mgm.setdefault("source_label", "MGM")
if station_code:
mgm.setdefault("station_code", station_code)
if station_name:
mgm.setdefault("station_name", station_name)
def collect_observation_source(
weather: Any,
source: Any,
@@ -144,7 +175,16 @@ def collect_observation_source(
)
results: dict[str, Any] = {}
attach(results, normalized_city, bool(use_fahrenheit))
if normalized_source == "mgm":
attach(
results,
normalized_city,
include_mgm=True,
include_nearby=True,
)
_enrich_mgm_results(results, normalized_city)
else:
attach(results, normalized_city, bool(use_fahrenheit))
if not results:
return ObservationSourceResult(
source=normalized_source,