From ce2a553e08c1359a5c6a8e3d2b1f52d3726950cb Mon Sep 17 00:00:00 2001 From: "2569718930@qq.com" <2569718930@qq.com> Date: Fri, 6 Mar 2026 18:50:11 +0800 Subject: [PATCH] feat: implement a rule-based market weather alert engine, including a Telegram push utility and tests. --- src/analysis/market_alert_engine.py | 14 +++------ src/utils/telegram_push.py | 49 +++++++++++++++++++++++------ tests/test_market_alert_engine.py | 31 ++++++++++++++++++ 3 files changed, 75 insertions(+), 19 deletions(-) diff --git a/src/analysis/market_alert_engine.py b/src/analysis/market_alert_engine.py index ec5052ae..16d06f77 100644 --- a/src/analysis/market_alert_engine.py +++ b/src/analysis/market_alert_engine.py @@ -212,17 +212,13 @@ def _pick_ankara_center_station(nearby: List[Dict[str, Any]]) -> Optional[Dict[s if not nearby: return None - def _temp(row: Dict[str, Any]) -> float: - return _sf(row.get("temp")) or -999.0 - - priority_rows = [] for row in nearby: - name = str(row.get("name") or "").lower() + name = str(row.get("name") or "").strip().lower() sid = str(row.get("istNo") or "").strip() - if sid == "17130" or "center" in name or "b枚lge" in name or "etimesgut" in name: - priority_rows.append(row) - if priority_rows: - return max(priority_rows, key=_temp) + if sid == "17130": + return row + if name in {"ankara (bölge/center)", "ankara (bolge/center)"}: + return row return None diff --git a/src/utils/telegram_push.py b/src/utils/telegram_push.py index 9c35ff39..5b1c2f8d 100644 --- a/src/utils/telegram_push.py +++ b/src/utils/telegram_push.py @@ -111,6 +111,15 @@ def _severity_ok(alert_payload: Dict[str, Any], min_severity: str, min_trigger_c return SEVERITY_RANK.get(severity, 0) >= SEVERITY_RANK.get(min_severity, 0) +def _trigger_type_key(alert_payload: Dict[str, Any]) -> str: + trigger_types = sorted( + str(alert.get("type") or "").strip() + for alert in (alert_payload.get("triggered_alerts") or []) + if alert.get("type") + ) + return "|".join(trigger_types) + + def _alert_signature(alert_payload: Dict[str, Any]) -> str: rules = alert_payload.get("rules") or {} center_deb = rules.get("ankara_center_deb_hit") or {} @@ -178,19 +187,33 @@ def _maybe_send_alert( min_severity: str, min_trigger_count: int, ) -> bool: - if not _severity_ok(alert_payload, min_severity, min_trigger_count): - return False - - message = ((alert_payload.get("telegram") or {}).get("zh") or "").strip() - if not message: - return False - now_ts = int(time.time()) + last_by_city = state.setdefault("last_by_city", {}) + last_city = last_by_city.get(city) or {} + is_active = _severity_ok(alert_payload, min_severity, min_trigger_count) + message = ((alert_payload.get("telegram") or {}).get("zh") or "").strip() + + if not is_active or not message: + if last_city.get("active"): + last_by_city[city] = { + **last_city, + "active": False, + "cleared_ts": now_ts, + } + logger.info(f"trade alert disarmed city={city}") + return True + return False + signature = _alert_signature(alert_payload) - last_city = (state.get("last_by_city") or {}).get(city) or {} + trigger_key = _trigger_type_key(alert_payload) last_city_sig = last_city.get("signature") + last_city_key = str(last_city.get("trigger_key") or "") last_city_ts = int(last_city.get("ts") or 0) last_sig_ts = int((state.get("by_signature") or {}).get(signature) or 0) + last_city_active = bool(last_city.get("active")) + + if last_city_active and last_city_key == trigger_key: + return False if last_city_ts and now_ts - last_city_ts < cooldown_sec: return False @@ -198,11 +221,17 @@ def _maybe_send_alert( return False bot.send_message(chat_id, message) - state.setdefault("last_by_city", {})[city] = {"signature": signature, "ts": now_ts} + last_by_city[city] = { + "signature": signature, + "trigger_key": trigger_key, + "severity": alert_payload.get("severity"), + "ts": now_ts, + "active": True, + } state.setdefault("by_signature", {})[signature] = now_ts logger.info( f"trade alert pushed city={city} severity={alert_payload.get('severity')} " - f"trigger_count={alert_payload.get('trigger_count')}" + f"trigger_count={alert_payload.get('trigger_count')} trigger_key={trigger_key}" ) return True diff --git a/tests/test_market_alert_engine.py b/tests/test_market_alert_engine.py index dd5f016c..accd4c02 100644 --- a/tests/test_market_alert_engine.py +++ b/tests/test_market_alert_engine.py @@ -100,6 +100,37 @@ def test_ankara_center_hits_deb_triggers_force_push(): assert "Center信号" in out["telegram"]["zh"] +def test_ankara_center_signal_only_uses_official_center_station(): + city_weather = _sample_weather_payload() + city_weather["deb"]["prediction"] = 11.2 + city_weather["current"]["temp"] = 10.7 + city_weather["mgm_nearby"] = [ + { + "name": "Etimesgut", + "istNo": "17069", + "lat": 39.95, + "lon": 32.68, + "temp": 12.6, + }, + { + "name": "Ankara (Bölge/Center)", + "istNo": "17130", + "lat": 39.95, + "lon": 32.97, + "temp": 11.3, + }, + ] + + out = build_trading_alerts(city_weather=city_weather) + + center_rule = out["rules"]["ankara_center_deb_hit"] + assert center_rule["triggered"] is True + assert center_rule["center_station"]["istNo"] == "17130" + assert center_rule["center_station"]["name"] == "Ankara (Bölge/Center)" + assert "Ankara (Bölge/Center)" in out["telegram"]["zh"] + assert "Etimesgut" not in out["telegram"]["zh"] + + def test_peak_passed_guard_suppresses_late_day_cooldown_alerts(): city_weather = { "name": "wellington",