From 6943f637ad78643870b0244d593b6a1fcae9e44e Mon Sep 17 00:00:00 2001 From: "2569718930@qq.com" <2569718930@qq.com> Date: Tue, 16 Jun 2026 16:29:09 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=E5=9B=BE=E8=A1=A8=E5=AE=9E?= =?UTF-8?q?=E6=97=B6=E8=A7=82=E6=B5=8B=E5=8F=A0=E5=8A=A0=E7=BC=93=E5=AD=98?= =?UTF-8?q?=E9=99=8D=E7=BA=A7?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tests/test_latest_observation_overlay.py | 168 +++++++++++++++++++ tests/test_web_observability.py | 70 ++++++++ web/services/city_api.py | 44 +++-- web/services/latest_observation_overlay.py | 180 +++++++++++++++++++-- 4 files changed, 441 insertions(+), 21 deletions(-) diff --git a/tests/test_latest_observation_overlay.py b/tests/test_latest_observation_overlay.py index 38e2fd28..7af78b58 100644 --- a/tests/test_latest_observation_overlay.py +++ b/tests/test_latest_observation_overlay.py @@ -385,3 +385,171 @@ def test_overlay_latest_jma_does_not_downgrade_newer_payload_current(): assert result is payload assert result["current"]["temp"] == 25.0 assert result["local_time"] == "07:00" + + +def test_non_amsc_raw_overlays_keep_their_source_codes(): + cases = [ + ( + "amos", + observation_overlay.overlay_latest_amos_observation, + "seoul", + { + "observed_at": "2026-06-16T06:03:00+09:00", + "fetched_at": "2026-06-15T21:03:30+00:00", + "station_code": "RKSS", + "station_name": "Gimpo", + "payload": { + "source": "amos", + "source_label": "AMOS", + "icao": "RKSS", + "temp": 22.4, + "observation_time": "2026-06-16T06:03:00+09:00", + }, + }, + ), + ( + "hko_obs", + observation_overlay.overlay_latest_hko_observation, + "hong kong", + { + "observed_at": "2026-06-16T05:55:00+08:00", + "fetched_at": "2026-06-15T21:55:30+00:00", + "station_code": "HKO", + "station_name": "Hong Kong Observatory", + "payload": { + "source": "hko_obs", + "source_label": "HKO", + "station_code": "HKO", + "station_label": "Hong Kong Observatory", + "temp": 28.6, + "obs_time": "2026-06-16T05:55:00+08:00", + }, + }, + ), + ( + "mgm", + observation_overlay.overlay_latest_mgm_observation, + "ankara", + { + "observed_at": "2026-06-16T00:10:00+03:00", + "fetched_at": "2026-06-15T21:10:30+00:00", + "station_code": "17130", + "station_name": "Ankara", + "payload": { + "source": "mgm", + "source_label": "MGM", + "icao": "17130", + "station_label": "Ankara (Bölge/Center)", + "temp": 25.3, + "obs_time": "2026-06-16T00:10:00+03:00", + }, + }, + ), + ] + + for source_code, overlay_fn, city, row in cases: + class FakeDB: + def get_latest_raw_observation(self, source, requested_city): + assert (source, requested_city) == (source_code, city) + return row + + result = overlay_fn( + FakeDB(), + city, + { + "name": city, + "temp_symbol": "°C", + "current": { + "temp": 19.0, + "source_code": "metar", + "obs_time": "2026-06-15T00:00:00+00:00", + }, + "airport_current": { + "temp": 19.0, + "source_code": "metar", + "obs_time": "2026-06-15T00:00:00+00:00", + }, + }, + ) + + assert result["current"]["source_code"] == source_code + assert result["airport_current"]["source_code"] == source_code + assert result["current"]["source_label"] == row["payload"]["source_label"] + assert result["current"].get("settlement_source") != "amsc_awos" + assert result["canonical_temperature"]["source"] == source_code + + +def test_overlay_latest_cwa_resets_stale_taipei_detail_to_latest_local_day(): + class FakeWeather: + def fetch_cwa_taipei_settlement_current(self): + return { + "source": "cwa", + "source_label": "CWA", + "station_code": "466920", + "station_name": "臺北", + "observation_time": "2026-06-16T15:30:00+08:00", + "current": { + "temp": 29.4, + "max_temp_so_far": 31.0, + }, + } + + stale_payload = { + "name": "taipei", + "display_name": "Taipei", + "temp_symbol": "°C", + "local_date": "2026-06-14", + "local_time": "18:00", + "current": { + "temp": 24.0, + "source_code": "metar", + "obs_time": "2026-06-14T10:00:00+00:00", + }, + "airport_current": { + "temp": 24.0, + "source_code": "metar", + "obs_time": "2026-06-14T10:00:00+00:00", + }, + "airport_primary": { + "temp": 24.0, + "source_code": "metar", + "obs_time": "2026-06-14T10:00:00+00:00", + }, + "overview": { + "local_date": "2026-06-14", + "local_time": "18:00", + "current_temp": 24.0, + }, + "metar_today_obs": [ + {"time": "18:00", "temp": 24.0}, + ], + "timeseries": { + "metar_today_obs": [ + {"time": "18:00", "temp": 24.0}, + ], + }, + } + + result = observation_overlay.overlay_latest_cwa_observation( + FakeWeather(), + "taipei", + stale_payload, + ) + + assert result["local_date"] == "2026-06-16" + assert result["local_time"] == "15:30" + assert result["current"]["temp"] == 29.4 + assert result["current"]["source_code"] == "cwa" + assert result["airport_current"]["obs_time"] == "2026-06-16T15:30:00+08:00" + assert result["overview"]["local_date"] == "2026-06-16" + assert result["overview"]["current_temp"] == 29.4 + assert result["metar_today_obs"] == [ + { + "time": "15:30", + "temp": 29.4, + "obs_time": "2026-06-16T15:30:00+08:00", + "source_code": "cwa", + "source_label": "CWA", + } + ] + assert result["timeseries"]["metar_today_obs"] == result["metar_today_obs"] diff --git a/tests/test_web_observability.py b/tests/test_web_observability.py index 3d210153..faf8ada1 100644 --- a/tests/test_web_observability.py +++ b/tests/test_web_observability.py @@ -1431,6 +1431,76 @@ def test_chart_data_returns_cached_payload_when_optional_overlay_times_out(monke ] +def test_chart_data_waits_for_latest_observation_overlay_when_optional_timeout_is_short(monkeypatch): + import asyncio + + class FakeCache: + def get_city_cache(self, kind, city): + assert kind == "full" + return { + "payload": { + "name": city, + "display_name": city.title(), + "temp_symbol": "°C", + "risk": {"icao": "ZUUU"}, + "current": { + "temp": 21.0, + "source_code": "metar", + "obs_time": "2026-06-14T16:00:00+00:00", + }, + "airport_current": { + "temp": 21.0, + "source_code": "metar", + "obs_time": "2026-06-14T16:00:00+00:00", + }, + "hourly": {"times": ["13:00"], "temps": [21.0]}, + }, + } + + def get_runway_obs_recent(self, icao, minutes=60): + return [] + + def get_latest_raw_observation(self, source, city): + if (source, city) != ("amsc_awos", "chengdu"): + return None + return { + "source": "amsc_awos", + "city": "chengdu", + "station_code": "ZUUU", + "station_name": "Chengdu Shuangliu", + "status": "ok", + "observed_at": "2026-06-14T17:00:00+00:00", + "fetched_at": "2026-06-14T17:00:30+00:00", + "payload": { + "source": "amsc_awos", + "source_label": "AMSC AWOS Chengdu Shuangliu (ZUUU)", + "icao": "ZUUU", + "temp_c": 25.8, + "observation_time": "2026-06-14T17:00:00+00:00", + "observation_time_local": "2026-06-15 01:00:00", + }, + } + + async def fake_run_in_threadpool(fn, *args, **kwargs): + if fn is city_api.overlay_latest_amsc_observation: + await asyncio.sleep(0.05) + return fn(*args, **kwargs) + + monkeypatch.setenv("POLYWEATHER_CITY_CHART_OPTIONAL_OVERLAY_TIMEOUT_MS", "1") + monkeypatch.setattr(city_api, "run_in_threadpool", fake_run_in_threadpool) + monkeypatch.setattr(city_api.legacy_routes, "_CACHE_DB", FakeCache()) + monkeypatch.setattr( + city_api.legacy_routes, + "_overlay_latest_wunderground_current", + lambda city, payload: payload, + ) + + payload = asyncio.run(city_api._get_city_chart_data("chengdu", force_refresh=False)) + + assert payload["current"]["temp"] == 25.8 + assert payload["airport_current"]["source_code"] == "amsc_awos" + + def test_chart_detail_payload_uses_threadpool_and_reuses_short_cache(monkeypatch): import asyncio diff --git a/web/services/city_api.py b/web/services/city_api.py index 144937a4..adb3b5c2 100644 --- a/web/services/city_api.py +++ b/web/services/city_api.py @@ -118,6 +118,26 @@ async def _run_optional_city_chart_overlay( return payload +async def _run_latest_observation_city_chart_overlay( + *, + city: str, + overlay_name: str, + payload: Dict[str, Any], + fn: Callable[..., Dict[str, Any]], + args: Tuple[Any, ...], +) -> Dict[str, Any]: + try: + return await run_in_threadpool(fn, *args) + except Exception as exc: + logger.debug( + "city chart latest observation overlay skipped city={} overlay={}: {}", + city, + overlay_name, + exc, + ) + return payload + + async def _get_cached_city_payload(city: str, kind: str) -> Dict[str, Any]: cached_entry = await run_in_threadpool(legacy_routes._CACHE_DB.get_city_cache, kind, city) if not isinstance(cached_entry, dict): @@ -661,42 +681,42 @@ async def _get_city_chart_data(city: str, *, force_refresh: bool) -> Dict[str, A fn=_overlay_cached_runway_history_from_db, args=(city, payload), ) - payload = await _run_optional_city_chart_overlay( + payload = await _run_latest_observation_city_chart_overlay( city=city, overlay_name="amsc_latest_raw", payload=payload, fn=overlay_latest_amsc_observation, args=(legacy_routes._CACHE_DB, city, payload), ) - payload = await _run_optional_city_chart_overlay( + payload = await _run_latest_observation_city_chart_overlay( city=city, overlay_name="jma_amedas_latest", payload=payload, fn=overlay_latest_jma_amedas_observation, args=(legacy_routes._weather, city, payload), ) - payload = await _run_optional_city_chart_overlay( + payload = await _run_latest_observation_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( + payload = await _run_latest_observation_city_chart_overlay( city=city, overlay_name="mgm_latest_raw", payload=payload, fn=overlay_latest_mgm_observation, args=(legacy_routes._CACHE_DB, city, payload), ) - payload = await _run_optional_city_chart_overlay( + payload = await _run_latest_observation_city_chart_overlay( city=city, overlay_name="hko_latest_raw", payload=payload, fn=overlay_latest_hko_observation, args=(legacy_routes._CACHE_DB, city, payload), ) - payload = await _run_optional_city_chart_overlay( + payload = await _run_latest_observation_city_chart_overlay( city=city, overlay_name="cwa_taipei", payload=payload, @@ -725,42 +745,42 @@ async def _get_city_chart_data(city: str, *, force_refresh: bool) -> Dict[str, A fn=_overlay_cached_runway_history_from_db, args=(city, payload), ) - payload = await _run_optional_city_chart_overlay( + payload = await _run_latest_observation_city_chart_overlay( city=city, overlay_name="amsc_latest_raw", payload=payload, fn=overlay_latest_amsc_observation, args=(legacy_routes._CACHE_DB, city, payload), ) - payload = await _run_optional_city_chart_overlay( + payload = await _run_latest_observation_city_chart_overlay( city=city, overlay_name="jma_amedas_latest", payload=payload, fn=overlay_latest_jma_amedas_observation, args=(legacy_routes._weather, city, payload), ) - payload = await _run_optional_city_chart_overlay( + payload = await _run_latest_observation_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( + payload = await _run_latest_observation_city_chart_overlay( city=city, overlay_name="mgm_latest_raw", payload=payload, fn=overlay_latest_mgm_observation, args=(legacy_routes._CACHE_DB, city, payload), ) - payload = await _run_optional_city_chart_overlay( + payload = await _run_latest_observation_city_chart_overlay( city=city, overlay_name="hko_latest_raw", payload=payload, fn=overlay_latest_hko_observation, args=(legacy_routes._CACHE_DB, city, payload), ) - payload = await _run_optional_city_chart_overlay( + payload = await _run_latest_observation_city_chart_overlay( city=city, overlay_name="cwa_taipei", payload=payload, diff --git a/web/services/latest_observation_overlay.py b/web/services/latest_observation_overlay.py index 118c1c95..e5e4721b 100644 --- a/web/services/latest_observation_overlay.py +++ b/web/services/latest_observation_overlay.py @@ -198,16 +198,33 @@ def _jma_observation_update( } -def _jma_today_point(local_time: str, obs_time: str, temp: float) -> dict[str, Any]: +def _observation_today_point( + local_time: str, + obs_time: str, + temp: float, + *, + source_code: str, + source_label: str, +) -> dict[str, Any]: return { "time": local_time, "temp": round(float(temp), 1), "obs_time": obs_time, - "source_code": "jma_amedas", - "source_label": "JMA", + "source_code": source_code, + "source_label": source_label, } +def _jma_today_point(local_time: str, obs_time: str, temp: float) -> dict[str, Any]: + return _observation_today_point( + local_time, + obs_time, + temp, + source_code="jma_amedas", + source_label="JMA", + ) + + def _replace_or_append_today_point( rows: Any, point: dict[str, Any], @@ -449,6 +466,82 @@ def _raw_observation_update( } +def _first_text(*values: Any) -> str: + for value in values: + text = str(value or "").strip() + if text: + return text + return "" + + +def _raw_source_observation_update( + city: str, + row: dict[str, Any], + raw_payload: dict[str, Any], + *, + source_code: str, + default_label: str, + temp_keys: tuple[str, ...] = ("temp", "temp_c"), + observed_at_keys: tuple[str, ...] = ("observation_time", "obs_time", "observed_at"), +) -> Optional[dict[str, Any]]: + temp = None + for key in temp_keys: + temp = _to_float(raw_payload.get(key)) + if temp is not None: + break + if temp is None: + return None + + observed_at = _first_text( + *(raw_payload.get(key) for key in observed_at_keys), + row.get("observed_at"), + ) + observed_at_local = _first_text( + raw_payload.get("observation_time_local"), + raw_payload.get("observed_at_local"), + row.get("observed_at_local"), + ) + normalized_source = str(source_code or raw_payload.get("source") or "").strip().lower() + source_label = _first_text(raw_payload.get("source_label"), default_label, normalized_source.upper()) + station_code = _first_text( + raw_payload.get("icao"), + raw_payload.get("station_code"), + raw_payload.get("istNo"), + row.get("station_code"), + ) + station_name = _first_text( + raw_payload.get("station_label"), + raw_payload.get("station_name"), + raw_payload.get("name"), + row.get("station_name"), + source_label, + ) + freshness = { + "freshness_status": "fresh", + "observed_at": observed_at or None, + "observed_at_local": observed_at_local or None, + "source_code": normalized_source, + "source_label": source_label, + } + update = { + "temp": round(temp, 1), + "source_code": normalized_source, + "source_label": source_label, + "station_code": station_code or None, + "station_name": station_name, + "observed_at": observed_at or None, + "observed_at_local": observed_at_local or None, + "obs_time": observed_at or observed_at_local, + "freshness": freshness, + "observation_status": "live", + "city": city, + } + if normalized_source: + update["settlement_source"] = normalized_source + update["settlement_source_label"] = source_label + return update + + def _merge_observation_block( payload: dict[str, Any], key: str, @@ -737,7 +830,8 @@ def overlay_latest_cwa_observation(weather, city, payload): return payload raw_epoch = parse_observation_epoch(obs_time) - if raw_epoch is None: + local_dt = _parse_observation_datetime(obs_time) + if raw_epoch is None or local_dt is None: return payload existing_epochs = [ epoch @@ -748,13 +842,18 @@ def overlay_latest_cwa_observation(weather, city, payload): ) if epoch is not None ] - if existing_epochs and max(existing_epochs) >= raw_epoch: + if existing_epochs and max(existing_epochs) > raw_epoch: return payload + source_label = str(cwa_data.get("source_label") or "CWA").strip() or "CWA" + current = cwa_data.get("current") if isinstance(cwa_data.get("current"), dict) else {} + max_so_far = _to_float(current.get("max_temp_so_far") if current else None) update = { "temp": round(float(temp), 1), "source_code": "cwa", - "source_label": str(cwa_data.get("source_label") or "CWA").strip(), + "source_label": source_label, + "settlement_source": "cwa", + "settlement_source_label": source_label, "station_code": str(cwa_data.get("station_code") or "466920").strip(), "station_name": str(cwa_data.get("station_name") or "\u81fa\u5317").strip(), "observed_at": obs_time, @@ -762,6 +861,9 @@ def overlay_latest_cwa_observation(weather, city, payload): "observation_status": "live", "city": normalized_city, } + if max_so_far is not None: + update["max_so_far"] = round(float(max_so_far), 1) + update["max_temp_so_far"] = round(float(max_so_far), 1) next_payload = deepcopy(payload) changed = False @@ -786,6 +888,46 @@ def overlay_latest_cwa_observation(weather, city, payload): next_payload["canonical_temperature"] = canonical_payload changed = True + local_date = local_dt.date().isoformat() + local_time = local_dt.strftime("%H:%M") + previous_local_date = str(next_payload.get("local_date") or "") + if next_payload.get("local_date") != local_date: + next_payload["local_date"] = local_date + changed = True + if next_payload.get("local_time") != local_time: + next_payload["local_time"] = local_time + changed = True + + overview = next_payload.get("overview") + if not isinstance(overview, dict): + overview = {} + next_overview = dict(overview) + overview_updates = { + "local_date": local_date, + "local_time": local_time, + "current_temp": round(float(temp), 1), + "airport_primary": update, + } + for key, value in overview_updates.items(): + if next_overview.get(key) != value: + next_overview[key] = value + changed = True + next_payload["overview"] = next_overview + + replace_today_series = bool(previous_local_date and previous_local_date != local_date) + _sync_jma_today_series( + next_payload, + _observation_today_point( + local_time, + obs_time, + temp, + source_code="cwa", + source_label=source_label, + ), + replace_all=replace_today_series, + ) + changed = True + return next_payload if changed else payload # ═══════════════════════════════════════════════════════════════════════════════ @@ -832,7 +974,13 @@ def overlay_latest_amos_observation(db, city, 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) + update = _raw_source_observation_update( + normalized_city, + row, + raw_payload, + source_code="amos", + default_label="AMOS", + ) if not update: return payload @@ -913,7 +1061,14 @@ def overlay_latest_hko_observation(db, city, payload): raw_epoch = _raw_hko_epoch(row, raw_payload) if raw_epoch is None: return payload - update = _raw_observation_update(normalized_city, row, raw_payload) + update = _raw_source_observation_update( + normalized_city, + row, + raw_payload, + source_code="hko_obs", + default_label="HKO", + observed_at_keys=("obs_time", "observation_time", "observed_at"), + ) if not update: return payload @@ -987,7 +1142,14 @@ def overlay_latest_mgm_observation(db, city, 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) + update = _raw_source_observation_update( + normalized_city, + row, + raw_payload, + source_code="mgm", + default_label="MGM", + observed_at_keys=("obs_time", "observation_time", "observed_at"), + ) if not update: return payload