Files
PolyWeather/web/services/latest_observation_overlay.py
2026-06-17 00:19:41 +08:00

1571 lines
54 KiB
Python

from __future__ import annotations
from copy import deepcopy
from datetime import datetime, timedelta, timezone
import time
from typing import Any, Optional
from loguru import logger
from src.database.runtime_state import OfficialIntradayObservationRepository
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
_TAIPEI_TZ = timezone(timedelta(hours=8))
_RAW_AMSC_RUNWAY_HISTORY_MIN_SPAN_SEC = 12 * 60 * 60
def parse_observation_epoch(value: Any) -> Optional[int]:
if value is None or value == "":
return None
if isinstance(value, datetime):
dt = value
else:
text = str(value).strip()
if not text:
return None
if text.endswith("Z"):
text = text[:-1] + "+00:00"
try:
dt = datetime.fromisoformat(text)
except ValueError:
for fmt in ("%Y-%m-%d %H:%M:%S%z", "%Y-%m-%d %H:%M:%S"):
try:
dt = datetime.strptime(text, fmt)
break
except ValueError:
continue
else:
return None
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return int(dt.timestamp())
def _parse_observation_datetime(value: Any) -> Optional[datetime]:
if value is None or value == "":
return None
if isinstance(value, datetime):
dt = value
else:
text = str(value).strip()
if not text:
return None
if text.endswith("Z"):
text = text[:-1] + "+00:00"
try:
dt = datetime.fromisoformat(text)
except ValueError:
for fmt in ("%Y-%m-%d %H:%M:%S%z", "%Y-%m-%d %H:%M:%S"):
try:
dt = datetime.strptime(text, fmt)
break
except ValueError:
continue
else:
return None
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt
def _payload_latest_epoch(payload: dict[str, Any], keys: tuple[str, ...]) -> Optional[int]:
values = [payload.get(key) for key in keys]
parsed = [epoch for epoch in (parse_observation_epoch(value) for value in values) if epoch is not None]
return max(parsed) if parsed else None
def _block_epoch(block: Any) -> Optional[int]:
if not isinstance(block, dict):
return None
canonical_epoch = _payload_latest_epoch(
block,
(
"observed_at",
"observation_time",
"obs_time",
),
)
if canonical_epoch is not None:
return canonical_epoch
return _payload_latest_epoch(
block,
(
"observed_at_local",
"observation_time_local",
),
)
def _raw_amsc_epoch(row: dict[str, Any], raw_payload: dict[str, Any]) -> Optional[int]:
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 _to_float(value: Any) -> Optional[float]:
try:
if value is None or value == "":
return None
return float(value)
except (TypeError, ValueError):
return None
def _to_int(value: Any) -> Optional[int]:
try:
if value is None or value == "":
return None
return int(float(value))
except (TypeError, ValueError):
return None
def _latest_airport_obs_log_row(
db: Any,
*,
station_code: str,
city: str,
source_code: str,
source_label: str,
station_label: str,
use_fahrenheit: bool,
) -> Optional[dict[str, Any]]:
reader = getattr(db, "get_airport_obs_recent", None)
if not callable(reader):
return None
try:
rows = reader(station_code, minutes=180)
except Exception as exc:
logger.debug("latest airport obs log read failed city={} station={}: {}", city, station_code, exc)
return None
latest: Optional[tuple[int, dict[str, Any]]] = None
normalized_city = str(city or "").strip().lower()
for row in rows if isinstance(rows, list) else []:
if not isinstance(row, dict):
continue
row_city = str(row.get("city") or "").strip().lower()
if row_city and normalized_city and row_city != normalized_city:
continue
temp_c = _to_float(row.get("temp_c") if row.get("temp_c") is not None else row.get("temp"))
obs_time = str(row.get("obs_time") or row.get("observed_at") or "").strip()
epoch = parse_observation_epoch(obs_time)
if temp_c is None or not obs_time or epoch is None:
continue
if latest is None or epoch > latest[0]:
latest = (epoch, row)
if latest is None:
return None
row = latest[1]
temp = _to_float(row.get("temp_c") if row.get("temp_c") is not None else row.get("temp"))
obs_time = str(row.get("obs_time") or row.get("observed_at") or "").strip()
if temp is None or not obs_time:
return None
if use_fahrenheit:
temp = temp * 9 / 5 + 32
return {
"station_label": station_label,
"temp": round(float(temp), 1),
"icao": station_code,
"source": source_code,
"source_label": source_label,
"obs_time": obs_time,
}
def _latest_jma_row(
weather: Any,
city: str,
use_fahrenheit: bool,
db: Any = None,
) -> Optional[dict[str, Any]]:
airport_obs_row = _latest_airport_obs_log_row(
db,
station_code="44166",
city=city,
source_code="jma_amedas",
source_label="JMA",
station_label="\u7fbd\u7530 10\u5206\u5b9e\u51b5 (JMA)",
use_fahrenheit=use_fahrenheit,
)
if airport_obs_row:
return airport_obs_row
fetcher = getattr(weather, "fetch_jma_amedas_official_nearby", None)
if callable(fetcher):
try:
rows = fetcher(city, use_fahrenheit=use_fahrenheit)
except Exception as exc:
logger.debug("latest JMA overlay read failed city={}: {}", city, exc)
rows = []
for row in rows if isinstance(rows, list) else []:
if not isinstance(row, dict):
continue
if _to_float(row.get("temp")) is not None and row.get("obs_time"):
return row
current_fetcher = getattr(weather, "fetch_jma_amedas_current", None)
if callable(current_fetcher):
try:
current = current_fetcher(city, use_fahrenheit=use_fahrenheit)
except Exception as exc:
logger.debug("latest JMA current overlay read failed city={}: {}", city, exc)
return None
if isinstance(current, dict):
temp = _to_float((current.get("current") or {}).get("temp"))
obs_time = current.get("obs_time")
if temp is not None and obs_time:
return {
"station_label": current.get("station_name"),
"temp": temp,
"icao": current.get("station_code"),
"source": "jma",
"source_label": "JMA",
"obs_time": obs_time,
}
return None
def _jma_observation_update(
city: str,
row: dict[str, Any],
obs_time: str,
temp: float,
) -> dict[str, Any]:
source_label = str(row.get("source_label") or "JMA").strip() or "JMA"
station_code = str(row.get("icao") or row.get("istNo") or "").strip() or None
station_name = str(row.get("station_label") or row.get("name") or source_label).strip()
freshness = {
"freshness_status": "fresh",
"observed_at": obs_time,
"source_code": "jma_amedas",
"source_label": source_label,
}
return {
"temp": round(float(temp), 1),
"source_code": "jma_amedas",
"source_label": source_label,
"station_code": station_code,
"station_name": station_name,
"station_label": station_name,
"observed_at": obs_time,
"observation_time": obs_time,
"obs_time": obs_time,
"freshness": freshness,
"observation_status": "live",
"city": city,
}
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": 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],
*,
replace_all: bool,
) -> list[dict[str, Any]]:
if replace_all:
return [point]
next_rows: list[dict[str, Any]] = []
replaced = False
for row in rows if isinstance(rows, list) else []:
if not isinstance(row, dict):
continue
current_time = str(row.get("time") or "").strip()
current_obs_time = str(row.get("obs_time") or row.get("observed_at") or "").strip()
if current_time == point["time"] or current_obs_time == point["obs_time"]:
next_rows.append(point)
replaced = True
else:
next_rows.append(dict(row))
if not replaced:
next_rows.append(point)
return next_rows
def _sync_jma_today_series(
payload: dict[str, Any],
point: dict[str, Any],
*,
replace_all: bool,
) -> None:
_sync_today_series_points(payload, [point], replace_all=replace_all)
def _sync_today_series_points(
payload: dict[str, Any],
points: list[dict[str, Any]],
*,
replace_all: bool,
) -> None:
clean_points = [dict(point) for point in points if isinstance(point, dict)]
if not clean_points:
return
if replace_all:
rows = clean_points
else:
rows = payload.get("metar_today_obs")
for point in clean_points:
rows = _replace_or_append_today_point(
rows,
point,
replace_all=False,
)
payload["metar_today_obs"] = rows
payload["airport_primary_today_obs"] = rows
official = payload.get("official")
if not isinstance(official, dict):
official = {}
official["airport_primary_today_obs"] = rows
payload["official"] = official
timeseries = payload.get("timeseries")
if not isinstance(timeseries, dict):
timeseries = {}
timeseries["metar_today_obs"] = rows
for key in ("metar_recent_obs", "settlement_today_obs"):
if replace_all and key in timeseries:
timeseries[key] = []
payload["timeseries"] = timeseries
if replace_all:
for key in ("metar_recent_obs", "settlement_today_obs"):
if key in payload:
payload[key] = []
def overlay_latest_jma_amedas_observation(
weather: Any,
city: str,
payload: dict[str, Any],
db: Any = None,
) -> dict[str, Any]:
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
if normalized_city != "tokyo":
return payload
use_fahrenheit = "F" in str(payload.get("temp_symbol") or "").upper()
row = _latest_jma_row(weather, normalized_city, use_fahrenheit, db=db)
if not isinstance(row, dict):
return payload
temp = _to_float(row.get("temp"))
obs_time = str(row.get("obs_time") or "").strip()
if temp is None or not obs_time:
return payload
raw_epoch = parse_observation_epoch(obs_time)
local_dt = _parse_observation_datetime(obs_time)
if raw_epoch is None or local_dt is None:
return payload
existing_epochs = [
epoch
for epoch in (
_block_epoch(payload.get("current")),
_block_epoch(payload.get("airport_primary")),
_block_epoch(payload.get("airport_current")),
_block_epoch(payload.get("canonical_temperature")),
)
if epoch is not None
]
if existing_epochs and max(existing_epochs) >= raw_epoch:
return payload
update = _jma_observation_update(normalized_city, row, obs_time, temp)
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": obs_time,
"current": update,
},
fetched_at=obs_time,
)
if canonical_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,
_jma_today_point(local_time, obs_time, temp),
replace_all=replace_today_series,
)
changed = True
return next_payload if changed else payload
def _amsc_payload_has_observation(raw_payload: dict[str, Any]) -> bool:
temp = raw_payload.get("temp_c") if raw_payload.get("temp_c") is not None else raw_payload.get("temp")
return _to_float(temp) is not None
def _latest_amsc_row(db: Any, city: str) -> tuple[Optional[dict[str, Any]], Optional[dict[str, Any]]]:
getter = getattr(db, "get_latest_raw_observation", None)
if not callable(getter):
return None, None
try:
row = getter("amsc_awos", city)
except Exception as exc:
logger.debug("latest AMSC 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 _amsc_payload_has_observation(raw_payload):
return row, raw_payload
lister = getattr(db, "list_latest_raw_observations_for_city", None)
if callable(lister):
try:
candidates = lister(city, limit=20)
except Exception as exc:
logger.debug("latest AMSC raw fallback list failed city={}: {}", city, exc)
candidates = []
for candidate in candidates if isinstance(candidates, list) else []:
if not isinstance(candidate, dict):
continue
if str(candidate.get("source") or "").strip().lower() != "amsc_awos":
continue
candidate_payload = candidate.get("payload")
if isinstance(candidate_payload, dict) and _amsc_payload_has_observation(candidate_payload):
return candidate, candidate_payload
return row, None
def _raw_observation_update(
city: str,
row: dict[str, Any],
raw_payload: dict[str, Any],
) -> Optional[dict[str, Any]]:
temp = _to_float(raw_payload.get("temp_c") if raw_payload.get("temp_c") is not None else raw_payload.get("temp"))
if temp is None:
return None
observed_at = str(
raw_payload.get("observation_time")
or raw_payload.get("observed_at")
or row.get("observed_at")
or ""
).strip()
observed_at_local = str(
raw_payload.get("observation_time_local")
or raw_payload.get("observed_at_local")
or ""
).strip()
source_label = str(raw_payload.get("source_label") or "AMSC AWOS").strip()
station_code = str(raw_payload.get("icao") or row.get("station_code") or "").strip().upper() or None
station_name = str(raw_payload.get("station_label") or row.get("station_name") or source_label).strip()
freshness = {
"freshness_status": "fresh",
"observed_at": observed_at or None,
"observed_at_local": observed_at_local or None,
"source_code": "amsc_awos",
"source_label": source_label,
}
return {
"temp": round(temp, 1),
"source_code": "amsc_awos",
"source_label": source_label,
"settlement_source": "amsc_awos",
"settlement_source_label": source_label,
"station_code": station_code,
"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,
}
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,
update: dict[str, Any],
raw_epoch: int,
) -> bool:
current = payload.get(key)
if not isinstance(current, dict):
current = {}
current_epoch = _block_epoch(current)
if current_epoch is not None and current_epoch >= raw_epoch:
return False
merged = dict(current)
merged.update(update)
payload[key] = merged
return True
def _runway_history_temp(point: dict[str, Any]) -> Optional[float]:
for key in ("target_runway_max", "temp", "tdz_temp", "end_temp", "mid_temp"):
temp = _to_float(point.get(key))
if temp is not None:
return temp
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,
payload: dict[str, Any],
row: dict[str, Any],
raw_payload: dict[str, Any],
*,
persist: bool = True,
copy_history: bool = True,
) -> bool:
runway_obs = raw_payload.get("runway_obs")
if not isinstance(runway_obs, dict):
return False
point_temperatures = runway_obs.get("point_temperatures")
if not isinstance(point_temperatures, list) or not point_temperatures:
return False
observed_at = str(
raw_payload.get("observation_time")
or raw_payload.get("observed_at")
or row.get("observed_at")
or ""
).strip()
if not observed_at:
return False
history = payload.get("runway_plate_history")
if not isinstance(history, dict):
history = {}
elif copy_history:
history = deepcopy(history)
use_fahrenheit = "F" in str(payload.get("temp_symbol") or "").upper()
appender = getattr(db, "append_runway_obs", None) if persist else None
icao = str(raw_payload.get("icao") or row.get("station_code") or "").strip().upper()
changed = False
for point in point_temperatures:
if not isinstance(point, dict):
continue
runway = str(point.get("runway") or "").strip().upper()
temp = _runway_history_temp(point)
if not runway or temp is None:
continue
if use_fahrenheit:
temp = temp * 9.0 / 5.0 + 32.0
entry = {
"time": observed_at,
"temp": round(float(temp), 1),
}
existing = history.get(runway)
runway_history = list(existing) if isinstance(existing, list) else []
replaced = False
for index, current in enumerate(runway_history):
if not isinstance(current, dict):
continue
current_time = str(
current.get("time")
or current.get("timestamp")
or current.get("observed_at")
or ""
).strip()
if current_time == observed_at:
if current != entry:
runway_history[index] = entry
changed = True
replaced = True
break
if not replaced:
runway_history.append(entry)
changed = True
history[runway] = runway_history
if callable(appender) and icao:
try:
appender(
icao=icao,
city=city,
runway=runway,
tdz_temp=_to_float(point.get("tdz_temp")),
mid_temp=_to_float(point.get("mid_temp")),
end_temp=_to_float(point.get("end_temp")),
target_runway_max=_to_float(point.get("target_runway_max")),
wind_dir=_to_int(point.get("wind_dir")),
wind_speed=_to_float(point.get("wind_speed")),
rvr=_to_int(point.get("rvr")),
mor=_to_float(point.get("mor")),
humidity=_to_float(point.get("humidity")),
otime_utc=observed_at,
)
except Exception as exc:
logger.debug(
"latest AMSC runway history persist skipped city={} runway={}: {}",
city,
runway,
exc,
)
if changed:
payload["runway_plate_history"] = history
return changed
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 _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):
return False
try:
rows = lister("amsc_awos", city, minutes=24 * 60, limit=1500)
except Exception as exc:
logger.debug("latest AMSC raw history overlay skipped city={}: {}", city, exc)
return False
payload["runway_plate_history"] = {}
changed = False
for row in rows if isinstance(rows, list) else []:
if not isinstance(row, dict):
continue
raw_payload = row.get("payload")
if not isinstance(raw_payload, dict) or not _amsc_payload_has_observation(raw_payload):
continue
changed = _append_latest_amsc_runway_history(
db,
city,
payload,
row,
raw_payload,
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
def overlay_latest_amsc_observation(
db: Any,
city: str,
payload: dict[str, Any],
) -> dict[str, Any]:
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_amsc_row(db, normalized_city)
if not isinstance(row, dict) or not isinstance(raw_payload, dict):
return payload
raw_epoch = _raw_amsc_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
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")
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
# ═══════════════════════════════════════════════════════════════════════════════
# CWA (Central Weather Administration — Taipei)
# ═══════════════════════════════════════════════════════════════════════════════
def _latest_cwa_data_from_airport_obs_log(db: Any, city: str, use_fahrenheit: bool) -> Optional[dict[str, Any]]:
row = _latest_airport_obs_log_row(
db,
station_code="466920",
city=city,
source_code="cwa",
source_label="CWA",
station_label="\u81fa\u5317",
use_fahrenheit=use_fahrenheit,
)
if not row:
return None
return {
"source": "cwa",
"source_label": "CWA",
"station_code": row.get("icao") or "466920",
"station_name": row.get("station_label") or "\u81fa\u5317",
"observation_time": row.get("obs_time"),
"current": {
"temp": row.get("temp"),
},
"unit": "fahrenheit" if use_fahrenheit else "celsius",
}
def _latest_cwa_data_from_official_intraday_history(use_fahrenheit: bool) -> Optional[dict[str, Any]]:
try:
local_now = datetime.now(timezone.utc).astimezone(timezone(timedelta(hours=8)))
target_date = local_now.strftime("%Y-%m-%d")
points = OfficialIntradayObservationRepository().load_points(
source_code="cwa",
station_code="466920",
target_date=target_date,
)
except Exception as exc:
logger.debug("latest CWA official intraday history read failed: {}", exc)
return None
latest: Optional[dict[str, Any]] = None
latest_epoch: Optional[int] = None
for point in points if isinstance(points, list) else []:
if not isinstance(point, dict):
continue
time_text = str(point.get("time") or "").strip()
temp = _to_float(point.get("temp"))
if len(time_text) != 5 or ":" not in time_text or temp is None:
continue
obs_time = f"{target_date}T{time_text}:00+08:00"
epoch = parse_observation_epoch(obs_time)
if epoch is None:
continue
if latest_epoch is None or epoch > latest_epoch:
latest_epoch = epoch
latest = {
"source": "cwa",
"source_label": "CWA",
"station_code": "466920",
"station_name": "\u81fa\u5317",
"observation_time": obs_time,
"current": {
"temp": round(float(temp) * 9 / 5 + 32, 1) if use_fahrenheit else round(float(temp), 1),
},
"unit": "fahrenheit" if use_fahrenheit else "celsius",
}
return latest
def _newer_cwa_payload(left: Optional[dict[str, Any]], right: Optional[dict[str, Any]]) -> Optional[dict[str, Any]]:
if not isinstance(left, dict):
return right if isinstance(right, dict) else None
if not isinstance(right, dict):
return left
left_epoch = parse_observation_epoch(left.get("observation_time"))
right_epoch = parse_observation_epoch(right.get("observation_time"))
if right_epoch is not None and (left_epoch is None or right_epoch > left_epoch):
return right
return left
def _taipei_local_today() -> str:
return datetime.now(timezone.utc).astimezone(_TAIPEI_TZ).strftime("%Y-%m-%d")
def _taipei_local_datetime(value: Any) -> Optional[datetime]:
dt = _parse_observation_datetime(value)
if dt is None:
return None
return dt.astimezone(_TAIPEI_TZ)
def _is_taipei_today_observation(value: Any, target_date: str) -> bool:
local_dt = _taipei_local_datetime(value)
return bool(local_dt and local_dt.date().isoformat() == target_date)
def _latest_taipei_metar_from_airport_obs_log(
db: Any,
use_fahrenheit: bool,
*,
target_date: str,
) -> Optional[dict[str, Any]]:
latest: Optional[dict[str, Any]] = None
for station_code, station_name in (
("RCSS", "Taipei Songshan METAR"),
("RCTP", "Taipei Taoyuan METAR"),
):
row = _latest_airport_obs_log_row(
db,
station_code=station_code,
city="taipei",
source_code="metar",
source_label="METAR",
station_label=station_name,
use_fahrenheit=use_fahrenheit,
)
if not row or not _is_taipei_today_observation(row.get("obs_time"), target_date):
continue
candidate = {
"source": "metar",
"source_label": "METAR",
"station_code": row.get("icao") or station_code,
"station_name": row.get("station_label") or station_name,
"observation_time": row.get("obs_time"),
"current": {
"temp": row.get("temp"),
},
"unit": "fahrenheit" if use_fahrenheit else "celsius",
}
latest = _newer_cwa_payload(latest, candidate)
return latest
def _call_taipei_metar_fetcher(weather: Any, use_fahrenheit: bool) -> Optional[dict[str, Any]]:
fetcher = getattr(weather, "fetch_metar", None)
if not callable(fetcher):
return None
try:
return fetcher("taipei", use_fahrenheit=use_fahrenheit, utc_offset=28800)
except TypeError:
try:
return fetcher("taipei", use_fahrenheit=use_fahrenheit)
except TypeError:
return fetcher("taipei")
except Exception as exc:
logger.debug("latest Taipei METAR fallback fetch failed: {}", exc)
return None
def _latest_taipei_metar_from_fetcher(
weather: Any,
use_fahrenheit: bool,
*,
target_date: str,
) -> Optional[dict[str, Any]]:
metar = _call_taipei_metar_fetcher(weather, use_fahrenheit)
if not isinstance(metar, dict) or metar.get("stale_for_today") is True:
return None
obs_time = str(metar.get("observation_time") or "").strip()
if not obs_time or not _is_taipei_today_observation(obs_time, target_date):
return None
current = metar.get("current") if isinstance(metar.get("current"), dict) else {}
temp = _to_float(current.get("temp"))
if temp is None:
return None
return {
"source": "metar",
"source_label": "METAR",
"station_code": str(metar.get("icao") or metar.get("station_code") or "RCSS").strip(),
"station_name": str(metar.get("station_name") or "Taipei Songshan METAR").strip(),
"observation_time": obs_time,
"current": {
"temp": round(float(temp), 1),
"max_temp_so_far": _to_float(current.get("max_temp_so_far")),
"max_temp_time": current.get("max_temp_time"),
"humidity": current.get("humidity"),
"wind_speed_kt": current.get("wind_speed_kt"),
"wind_dir": current.get("wind_dir"),
"raw_metar": current.get("raw_metar"),
"visibility_mi": current.get("visibility_mi"),
"wx_desc": current.get("wx_desc"),
},
"today_obs": metar.get("today_obs") or [],
"unit": "fahrenheit" if use_fahrenheit else "celsius",
}
def _latest_taipei_metar_fallback(
weather: Any,
db: Any,
use_fahrenheit: bool,
*,
target_date: str,
) -> Optional[dict[str, Any]]:
latest = _latest_taipei_metar_from_airport_obs_log(
db,
use_fahrenheit,
target_date=target_date,
)
latest = _newer_cwa_payload(
latest,
_latest_taipei_metar_from_fetcher(weather, use_fahrenheit, target_date=target_date),
)
return latest
def _normalized_today_points(
rows: Any,
latest_point: dict[str, Any],
*,
source_code: str,
source_label: str,
) -> list[dict[str, Any]]:
points: list[dict[str, Any]] = []
seen: set[tuple[str, str]] = set()
def add(point: dict[str, Any]) -> None:
time_key = str(point.get("time") or "")
obs_key = str(point.get("obs_time") or "")
if not time_key and not obs_key:
return
for index, existing in enumerate(points):
existing_time = str(existing.get("time") or "")
existing_obs = str(existing.get("obs_time") or "")
if existing_time == time_key and (not existing_obs or not obs_key or existing_obs == obs_key):
if obs_key and not existing_obs:
points[index] = point
return
key = (time_key, obs_key)
if key in seen:
return
seen.add(key)
points.append(point)
for row in rows if isinstance(rows, list) else []:
if isinstance(row, dict):
time_text = str(row.get("time") or "").strip()
temp = _to_float(row.get("temp"))
obs_time = str(row.get("obs_time") or row.get("observed_at") or "").strip()
elif isinstance(row, (list, tuple)) and len(row) >= 2:
time_text = str(row[0] or "").strip()
temp = _to_float(row[1])
obs_time = ""
else:
continue
if not time_text or temp is None:
continue
point = {
"time": time_text,
"temp": round(float(temp), 1),
"source_code": source_code,
"source_label": source_label,
}
if obs_time:
point["obs_time"] = obs_time
add(point)
add(dict(latest_point))
return points
def overlay_latest_cwa_observation(weather, city, payload, db=None):
normalized_city = str(city or payload.get("name") or payload.get("city") or "").strip().lower()
if normalized_city != "taipei" or not isinstance(payload, dict) or not payload:
return payload
use_fahrenheit = "F" in str(payload.get("temp_symbol") or "").upper()
target_date = _taipei_local_today()
cwa_data = _latest_cwa_data_from_airport_obs_log(db, normalized_city, use_fahrenheit)
cwa_data = _newer_cwa_payload(cwa_data, _latest_cwa_data_from_official_intraday_history(use_fahrenheit))
fetcher = getattr(weather, "fetch_cwa_taipei_settlement_current", None)
if not callable(fetcher) and not cwa_data:
return payload
if callable(fetcher):
try:
cwa_data = _newer_cwa_payload(cwa_data, fetcher())
except Exception as exc:
logger.debug("latest CWA overlay fetch failed: {}", exc)
if isinstance(cwa_data, dict) and not _is_taipei_today_observation(
cwa_data.get("observation_time"),
target_date,
):
cwa_epoch = parse_observation_epoch(cwa_data.get("observation_time"))
existing_epochs = [
epoch
for epoch in (
_block_epoch(payload.get("current")),
_block_epoch(payload.get("airport_primary")),
_block_epoch(payload.get("airport_current")),
)
if epoch is not None
]
latest_existing_epoch = max(existing_epochs) if existing_epochs else None
if cwa_epoch is None or (latest_existing_epoch is not None and cwa_epoch <= latest_existing_epoch):
cwa_data = None
if not isinstance(cwa_data, dict):
cwa_data = _latest_taipei_metar_fallback(
weather,
db,
use_fahrenheit,
target_date=target_date,
)
if not isinstance(cwa_data, dict):
return payload
temp = _to_float((cwa_data.get("current") or {}).get("temp"))
obs_time = str(cwa_data.get("observation_time") or "").strip()
if temp is None or not obs_time:
return payload
raw_epoch = parse_observation_epoch(obs_time)
local_dt = _taipei_local_datetime(obs_time)
if raw_epoch is None or local_dt is None:
return payload
existing_epochs = [
epoch
for epoch in (
_block_epoch(payload.get("current")),
_block_epoch(payload.get("airport_primary")),
_block_epoch(payload.get("airport_current")),
)
if epoch is not None
]
if existing_epochs and max(existing_epochs) > raw_epoch:
return payload
source_code = str(cwa_data.get("source") or "cwa").strip().lower() or "cwa"
source_label = str(cwa_data.get("source_label") or ("METAR" if source_code == "metar" else "CWA")).strip()
source_label = source_label or ("METAR" if source_code == "metar" else "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": source_code,
"source_label": source_label,
"settlement_source": source_code,
"settlement_source_label": source_label,
"station_code": str(cwa_data.get("station_code") or ("RCSS" if source_code == "metar" else "466920")).strip(),
"station_name": str(cwa_data.get("station_name") or ("Taipei Songshan METAR" if source_code == "metar" else "\u81fa\u5317")).strip(),
"observed_at": obs_time,
"observation_time": obs_time,
"obs_time": obs_time,
"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
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": obs_time,
"current": update,
},
fetched_at=obs_time,
)
if canonical_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,
"settlement_source": source_code,
"settlement_source_label": source_label,
}
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)
latest_point = _observation_today_point(
local_time,
obs_time,
temp,
source_code=source_code,
source_label=source_label,
)
today_points = _normalized_today_points(
cwa_data.get("today_obs"),
latest_point,
source_code=source_code,
source_label=source_label,
)
_sync_today_series_points(
next_payload,
today_points,
replace_all=replace_today_series,
)
changed = True
return next_payload if changed else payload
# ═══════════════════════════════════════════════════════════════════════════════
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_source_observation_update(
normalized_city,
row,
raw_payload,
source_code="amos",
default_label="AMOS",
)
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
# ═══════════════════════════════════════════════════════════════════════════════
# HKO (Hong Kong Observatory — Hong Kong, Shenzhen)
# ═══════════════════════════════════════════════════════════════════════════════
def _latest_hko_row(db, city):
getter = getattr(db, "get_latest_raw_observation", None)
if not callable(getter):
return None, None
try:
row = getter("hko_obs", city)
except Exception as exc:
logger.debug("latest HKO 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_hko_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_hko_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_hko_row(db, normalized_city)
if not isinstance(row, dict) or not isinstance(raw_payload, dict):
return payload
raw_epoch = _raw_hko_epoch(row, raw_payload)
if raw_epoch is None:
return 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
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
# ═══════════════════════════════════════════════════════════════════════════════
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_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
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