feat: implement Telegram push utility and corresponding hashtag validation tests

This commit is contained in:
2569718930@qq.com
2026-05-17 13:37:51 +08:00
parent 8e2317592d
commit 07bb8e1d15
2 changed files with 155 additions and 6 deletions
+99 -6
View File
@@ -6,7 +6,7 @@ import threading
import time import time
from concurrent.futures import ThreadPoolExecutor, as_completed from concurrent.futures import ThreadPoolExecutor, as_completed
from datetime import datetime, timezone from datetime import datetime, timezone
from typing import Any, Dict, List, Optional, Tuple from typing import Any, Dict, List, Optional, Set, Tuple
from loguru import logger from loguru import logger
@@ -489,8 +489,17 @@ def _alert_signature(alert_payload: Dict[str, Any]) -> str:
# ── high-freq airport push loop ── # ── high-freq airport push loop ──
HIGH_FREQ_AIRPORT_CITIES = {"seoul", "busan", "tokyo", "ankara", "helsinki", "amsterdam", "istanbul", "paris", "hong kong", "lau fau shan", "taipei", "beijing", "shanghai", "guangzhou", "shenzhen", "qingdao", "chengdu", "chongqing", "wuhan"} HIGH_FREQ_AIRPORT_CITIES = {"seoul", "singapore", "busan", "tokyo", "ankara", "helsinki", "amsterdam", "istanbul", "paris", "hong kong", "lau fau shan", "taipei", "beijing", "shanghai", "guangzhou", "shenzhen", "qingdao", "chengdu", "chongqing", "wuhan"}
HIGH_FREQ_AIRPORT_ICAO = {"seoul": "RKSI", "busan": "RKPK", "tokyo": "44166", "ankara": "17128", "helsinki": "EFHK", "amsterdam": "EHAM", "istanbul": "17058", "paris": "LFPB", "hong kong": "HKO", "lau fau shan": "LFS", "taipei": "466920", "beijing": "ZBAA", "shanghai": "ZSPD", "guangzhou": "ZGGG", "shenzhen": "ZGSZ", "qingdao": "ZSQD", "chengdu": "ZUUU", "chongqing": "ZUCK", "wuhan": "ZHHH"} HIGH_FREQ_AIRPORT_ICAO = {"seoul": "RKSI", "singapore": "WSSS", "busan": "RKPK", "tokyo": "44166", "ankara": "17128", "helsinki": "EFHK", "amsterdam": "EHAM", "istanbul": "17058", "paris": "LFPB", "hong kong": "HKO", "lau fau shan": "LFS", "taipei": "466920", "beijing": "ZBAA", "shanghai": "ZSPD", "guangzhou": "ZGGG", "shenzhen": "ZGSZ", "qingdao": "ZSQD", "chengdu": "ZUUU", "chongqing": "ZUCK", "wuhan": "ZHHH"}
FOCUS_RUNWAY_PAIRS = {
"chongqing": {("02L", "20R")},
"shanghai": {("17L", "35R")},
"wuhan": {("04", "22")},
"beijing": {("01", "19")},
"guangzhou": {("02L", "20R")},
"chengdu": {("02L", "20R")},
"seoul": {("15R", "33L")},
}
MARKET_MONITOR_INTERVAL_SEC = 300 MARKET_MONITOR_INTERVAL_SEC = 300
MARKET_MONITOR_CITIES = [ MARKET_MONITOR_CITIES = [
"seoul", "busan", "tokyo", "helsinki", "amsterdam", "seoul", "busan", "tokyo", "helsinki", "amsterdam",
@@ -561,6 +570,72 @@ def _fmt(value: Any) -> str:
return "--" return "--"
def _normalize_runway_label(value: Any) -> str:
return re.sub(r"[^0-9A-Z]+", "", str(value or "").strip().upper())
def _runway_pair_key(r1: Any, r2: Any) -> Tuple[str, str]:
a = _normalize_runway_label(r1)
b = _normalize_runway_label(r2)
return tuple(sorted((a, b))) # type: ignore[return-value]
def _focus_runway_pairs_for_city(city: str) -> Set[Tuple[str, str]]:
return {_runway_pair_key(a, b) for a, b in FOCUS_RUNWAY_PAIRS.get(city, set())}
def _select_focus_runway_obs(
city: str,
runway_pairs: List[Any],
runway_temps: List[Any],
point_temps: Optional[List[Any]] = None,
) -> Tuple[List[Any], List[Any], List[Any]]:
"""Return only market-relevant runway pairs when configured for the city.
If a configured focus pair is not present in the upstream payload, fall back
to the original lists so the push still carries useful airport evidence.
"""
focus_pairs = _focus_runway_pairs_for_city(city)
if not focus_pairs or not runway_pairs or not runway_temps:
return runway_pairs, runway_temps, point_temps or []
selected_pairs: List[Any] = []
selected_temps: List[Any] = []
selected_points: List[Any] = []
points = point_temps or []
for i, (pair, temp) in enumerate(zip(runway_pairs, runway_temps)):
try:
r1, r2 = pair
except Exception:
continue
if _runway_pair_key(r1, r2) not in focus_pairs:
continue
selected_pairs.append(pair)
selected_temps.append(temp)
if i < len(points):
selected_points.append(points[i])
if selected_pairs:
return selected_pairs, selected_temps, selected_points
return runway_pairs, runway_temps, points
def _focused_runway_max(city: str, city_weather: Dict[str, Any]) -> Optional[float]:
amos = city_weather.get("amos") or {}
runway_obs = (amos.get("runway_obs") or {}) if isinstance(amos, dict) else {}
runway_pairs = runway_obs.get("runway_pairs") or []
runway_temps = runway_obs.get("temperatures") or []
runway_pairs, runway_temps, _points = _select_focus_runway_obs(
city,
runway_pairs,
runway_temps,
runway_obs.get("point_temperatures") or [],
)
del runway_pairs
valid = [float(t) for (t, _d) in runway_temps if t is not None]
return max(valid) if valid else None
def _format_percent(value: Any) -> str: def _format_percent(value: Any) -> str:
try: try:
numeric = float(value) numeric = float(value)
@@ -587,6 +662,9 @@ def _build_market_monitor_message(city: str, city_weather: Dict[str, Any]) -> st
local_time = str(city_weather.get("local_time") or "").strip() or "--" local_time = str(city_weather.get("local_time") or "").strip() or "--"
city_label = str(city or "").strip().title() city_label = str(city or "").strip().title()
current_temp = airport_cur.get("temp") if airport_cur.get("temp") is not None else current.get("temp") current_temp = airport_cur.get("temp") if airport_cur.get("temp") is not None else current.get("temp")
focus_runway_temp = _focused_runway_max(city, city_weather)
if focus_runway_temp is not None:
current_temp = focus_runway_temp
deb_pred = deb.get("prediction") deb_pred = deb.get("prediction")
temp_symbol = str(city_weather.get("temp_symbol") or "°C").strip() temp_symbol = str(city_weather.get("temp_symbol") or "°C").strip()
@@ -730,7 +808,7 @@ def _build_airport_status_message(
local_time: str = "", local_time: str = "",
state: str = "", state: str = "",
) -> str: ) -> str:
_AIRPORT_EN = {"seoul": "Incheon", "busan": "Gimhae", "tokyo": "Haneda", _AIRPORT_EN = {"seoul": "Incheon", "singapore": "Changi", "busan": "Gimhae", "tokyo": "Haneda",
"ankara": "Esenboğa", "helsinki": "Vantaa", "amsterdam": "Schiphol", "ankara": "Esenboğa", "helsinki": "Vantaa", "amsterdam": "Schiphol",
"istanbul": "Airport", "paris": "Le Bourget", "istanbul": "Airport", "paris": "Le Bourget",
"hong kong": "Observatory", "lau fau shan": "Lau Fau Shan", "hong kong": "Observatory", "lau fau shan": "Lau Fau Shan",
@@ -746,6 +824,10 @@ def _build_airport_status_message(
runway_data = (amos.get("runway_obs") or {}) if amos else {} runway_data = (amos.get("runway_obs") or {}) if amos else {}
runway_pairs = runway_data.get("runway_pairs") or [] runway_pairs = runway_data.get("runway_pairs") or []
runway_temps = runway_data.get("temperatures") or [] runway_temps = runway_data.get("temperatures") or []
point_temps = runway_data.get("point_temperatures") or []
runway_pairs, runway_temps, point_temps = _select_focus_runway_obs(
city, runway_pairs, runway_temps, point_temps
)
mgm_nearby = city_weather.get("mgm_nearby") or [] mgm_nearby = city_weather.get("mgm_nearby") or []
airport_icao = HIGH_FREQ_AIRPORT_ICAO.get(city, "") airport_icao = HIGH_FREQ_AIRPORT_ICAO.get(city, "")
airport_row = None airport_row = None
@@ -788,7 +870,6 @@ def _build_airport_status_message(
# Runway detail block # Runway detail block
if is_amsc and runway_pairs and runway_temps and len(runway_pairs) == len(runway_temps): if is_amsc and runway_pairs and runway_temps and len(runway_pairs) == len(runway_temps):
lines.append("") lines.append("")
point_temps = runway_data.get("point_temperatures") or []
for i, ((r1, r2), (t, _d)) in enumerate(zip(runway_pairs, runway_temps)): for i, ((r1, r2), (t, _d)) in enumerate(zip(runway_pairs, runway_temps)):
if t is not None: if t is not None:
pts = point_temps[i] if i < len(point_temps) else {} pts = point_temps[i] if i < len(point_temps) else {}
@@ -874,6 +955,7 @@ _AIRPORT_PUSH_INTERVAL = {
"paris": 60, "paris": 60,
"hong kong": 60, "hong kong": 60,
"lau fau shan": 60, "lau fau shan": 60,
"singapore": 60,
"taipei": 60, "taipei": 60,
"beijing": 60, "beijing": 60,
"shanghai": 60, "shanghai": 60,
@@ -992,7 +1074,15 @@ def _run_high_freq_airport_cycle(
station_temp = airport_row.get("temp") if airport_row else None station_temp = airport_row.get("temp") if airport_row else None
current_obs_time = str(airport_row.get("obs_time") or "") current_obs_time = str(airport_row.get("obs_time") or "")
runway_temps = (amos.get("runway_obs") or {}).get("temperatures") or [] runway_obs = (amos.get("runway_obs") or {})
runway_pairs = runway_obs.get("runway_pairs") or []
runway_temps = runway_obs.get("temperatures") or []
runway_pairs, runway_temps, _point_temps = _select_focus_runway_obs(
city,
runway_pairs,
runway_temps,
runway_obs.get("point_temperatures") or [],
)
if runway_temps: if runway_temps:
valid_temps = [t for (t, _d) in runway_temps if t is not None] valid_temps = [t for (t, _d) in runway_temps if t is not None]
if valid_temps: if valid_temps:
@@ -1175,6 +1265,9 @@ def _run_market_monitor_cycle(bot: Any, chat_ids: List[str]) -> bool:
city_weather["market_scan"] = market city_weather["market_scan"] = market
ac = city_weather.get("airport_current") or {} ac = city_weather.get("airport_current") or {}
current_temp = ac.get("temp") if ac.get("temp") is not None else (city_weather.get("current") or {}).get("temp") current_temp = ac.get("temp") if ac.get("temp") is not None else (city_weather.get("current") or {}).get("temp")
focus_runway_temp = _focused_runway_max(city, city_weather)
if focus_runway_temp is not None:
current_temp = focus_runway_temp
deb_pred = (city_weather.get("deb") or {}).get("prediction") deb_pred = (city_weather.get("deb") or {}).get("prediction")
if current_temp is not None and deb_pred is not None: if current_temp is not None and deb_pred is not None:
delta = float(current_temp) - float(deb_pred) delta = float(current_temp) - float(deb_pred)
+56
View File
@@ -1,4 +1,7 @@
from src.utils.telegram_push import ( from src.utils.telegram_push import (
HIGH_FREQ_AIRPORT_CITIES,
HIGH_FREQ_AIRPORT_ICAO,
MARKET_MONITOR_CITIES,
MARKET_MONITOR_INTERVAL_SEC, MARKET_MONITOR_INTERVAL_SEC,
_build_airport_status_message, _build_airport_status_message,
_build_market_monitor_message, _build_market_monitor_message,
@@ -52,5 +55,58 @@ def test_market_monitor_message_starts_with_market_hashtag_and_city():
assert "当前:29.4°C · DEB31.2°C" in text assert "当前:29.4°C · DEB31.2°C" in text
def test_market_monitor_uses_focus_runway_temperature_for_key_airports():
text = _build_market_monitor_message(
"shanghai",
{
"local_time": "14:01",
"airport_current": {"temp": 29.4},
"deb": {"prediction": 31.2},
"amos": {
"source": "amsc_awos",
"runway_obs": {
"runway_pairs": [("16L", "34R"), ("17L", "35R")],
"temperatures": [(36.0, None), (31.8, None)],
},
},
},
)
assert "当前:31.8°C · DEB31.2°C" in text
def test_airport_status_hides_non_focus_runways_for_key_airports():
text = _build_airport_status_message(
"chongqing",
{
"current": {"temp": 28.0},
"airport_current": {"max_so_far": 30.0, "max_temp_time": "13:00"},
"amos": {
"source": "amsc_awos",
"runway_obs": {
"runway_pairs": [("02L", "20R"), ("02R", "20L")],
"temperatures": [(31.1, None), (34.9, None)],
"point_temperatures": [
{"tdz_temp": 31.0, "mid_temp": 31.1, "end_temp": 31.2},
{"tdz_temp": 34.8, "mid_temp": 34.9, "end_temp": 35.0},
],
},
},
},
32.0,
"13:00",
)
assert "02L/20R" in text
assert "02R/20L" not in text
assert "当前:31.1°C" in text
def test_market_monitor_default_interval_is_five_minutes(): def test_market_monitor_default_interval_is_five_minutes():
assert MARKET_MONITOR_INTERVAL_SEC == 300 assert MARKET_MONITOR_INTERVAL_SEC == 300
def test_singapore_is_in_telegram_push_city_lists():
assert "singapore" not in MARKET_MONITOR_CITIES
assert "singapore" in HIGH_FREQ_AIRPORT_CITIES
assert HIGH_FREQ_AIRPORT_ICAO["singapore"] == "WSSS"