From c9719cf577a97563f06e74aa4056e0ca0734c162 Mon Sep 17 00:00:00 2001 From: "2569718930@qq.com" <2569718930@qq.com> Date: Wed, 1 Apr 2026 18:24:50 +0800 Subject: [PATCH] Add local peak timing and /markets bot digest --- .env.example | 1 - docs/CONFIGURATION_ZH.md | 7 +- src/bot/handlers/basic.py | 25 +++++ src/bot/io_layer.py | 1 + src/bot/orchestrator.py | 1 + src/bot/runtime_coordinator.py | 2 +- src/utils/telegram_push.py | 170 +++++++++++++++++++++++++------- tests/test_bot_basic_handler.py | 30 ++++++ 8 files changed, 192 insertions(+), 45 deletions(-) diff --git a/.env.example b/.env.example index cad17fc8..0b11d25b 100644 --- a/.env.example +++ b/.env.example @@ -83,7 +83,6 @@ TELEGRAM_MARKET_FOCUS_DIGEST_ENABLED=true TELEGRAM_MARKET_FOCUS_DIGEST_HOURS=11,18 TELEGRAM_MARKET_FOCUS_DIGEST_TOP_N=5 TELEGRAM_MARKET_FOCUS_DIGEST_GRACE_MINUTES=180 -TELEGRAM_MARKET_TIMEZONE=Asia/Shanghai TELEGRAM_ALERT_CITIES=ankara,london,paris,seoul,hong kong,shanghai,singapore,tokyo,tel aviv,toronto,buenos aires,wellington,new york,chicago,dallas,miami,atlanta,seattle,lucknow,sao paulo,munich POLYWEATHER_MONITORING_ALERT_CHAT_IDS= diff --git a/docs/CONFIGURATION_ZH.md b/docs/CONFIGURATION_ZH.md index 3662b5e2..f12605ed 100644 --- a/docs/CONFIGURATION_ZH.md +++ b/docs/CONFIGURATION_ZH.md @@ -129,7 +129,6 @@ PolyWeather 的环境变量很多,但不是所有变量都属于同一层级 - `TELEGRAM_MARKET_FOCUS_DIGEST_HOURS` - `TELEGRAM_MARKET_FOCUS_DIGEST_TOP_N` - `TELEGRAM_MARKET_FOCUS_DIGEST_GRACE_MINUTES` -- `TELEGRAM_MARKET_TIMEZONE` - `POLYWEATHER_PAYMENT_RPC_URLS` - `TAF_CACHE_TTL_SEC` @@ -238,7 +237,6 @@ TELEGRAM_MARKET_FOCUS_DIGEST_ENABLED=true TELEGRAM_MARKET_FOCUS_DIGEST_HOURS=11,18 TELEGRAM_MARKET_FOCUS_DIGEST_TOP_N=5 TELEGRAM_MARKET_FOCUS_DIGEST_GRACE_MINUTES=180 -TELEGRAM_MARKET_TIMEZONE=Asia/Shanghai POLYMARKET_WALLET_ACTIVITY_ENABLED=false ``` @@ -251,7 +249,7 @@ POLYMARKET_WALLET_ACTIVITY_ENABLED=false - `POLYWEATHER_STATE_STORAGE_MODE` 当前线上推荐直接使用 `sqlite`。 - `POLYWEATHER_PAYMENT_RPC_URLS` 支持逗号分隔多个 RPC;如果暂时只用单 RPC,也可以继续只配 `POLYWEATHER_PAYMENT_RPC_URL`。 - 机器人市场监控当前分成两类消息:`关键提醒` 与 `关注清单`。 -- `TELEGRAM_MARKET_TIMEZONE` 建议保持 `Asia/Shanghai`,更适合亚洲时区盯盘节奏。 +- `TELEGRAM_MARKET_FOCUS_DIGEST_HOURS` 直接按服务器本地时间解释,不再单独配置时区。 - `POLYMARKET_WALLET_ACTIVITY_ENABLED` 已退役,保留为 `false` 即可,不建议再启用钱包异动监听。 ### 6.3 机器人市场监控建议配置 @@ -276,14 +274,13 @@ TELEGRAM_MARKET_FOCUS_DIGEST_ENABLED=true TELEGRAM_MARKET_FOCUS_DIGEST_HOURS=11,18 TELEGRAM_MARKET_FOCUS_DIGEST_TOP_N=5 TELEGRAM_MARKET_FOCUS_DIGEST_GRACE_MINUTES=180 -TELEGRAM_MARKET_TIMEZONE=Asia/Shanghai POLYMARKET_WALLET_ACTIVITY_ENABLED=false ``` 说明: - `TELEGRAM_ALERT_MISPRICING_ONLY=true` 表示关键提醒优先围绕错价/市场触发,不把机器人做成泛通知器。 -- `TELEGRAM_MARKET_FOCUS_DIGEST_HOURS=11,18` 对应亚洲白天与傍晚两个摘要窗口,适合提前筛出今晚值得关注的市场。 +- `TELEGRAM_MARKET_FOCUS_DIGEST_HOURS=11,18` 按服务器本地时间解释,通常可以对应白天与傍晚两个摘要窗口。 - `TELEGRAM_MARKET_FOCUS_DIGEST_TOP_N=5` 建议先保持较小,避免机器人一次推太多城市。 - `TELEGRAM_MARKET_FOCUS_DIGEST_GRACE_MINUTES=180` 用于容忍进程重启/部署后的补发窗口。 - `POLYMARKET_WALLET_ACTIVITY_ENABLED=false` 表示停用旧的钱包异动监听,统一收敛到市场监控。 diff --git a/src/bot/handlers/basic.py b/src/bot/handlers/basic.py index f1034f07..cf49d6e7 100644 --- a/src/bot/handlers/basic.py +++ b/src/bot/handlers/basic.py @@ -9,6 +9,7 @@ from src.bot.observability import CommandTrace from src.bot.runtime_coordinator import RuntimeStatus, render_runtime_status_html _BASIC_COMMANDS = {"start", "help", "id", "top", "diag", "bind", "unbind"} +_BASIC_COMMANDS = {"start", "help", "id", "top", "diag", "bind", "unbind", "markets"} class BasicCommandHandler: @@ -17,10 +18,12 @@ class BasicCommandHandler: bot: Any, io_layer: BotIOLayer, runtime_status_provider: Callable[[], RuntimeStatus], + config: dict | None = None, ): self.bot = bot self.io_layer = io_layer self.runtime_status_provider = runtime_status_provider + self.config = config or {} def register(self) -> None: @self.bot.message_handler(commands=["start", "help"]) @@ -47,6 +50,10 @@ class BasicCommandHandler: def _unbind(message): self._dispatch(message) + @self.bot.message_handler(commands=["markets"]) + def _markets(message): + self._dispatch(message) + @self.bot.message_handler( content_types=["text"], func=lambda message: extract_command_name( @@ -87,6 +94,9 @@ class BasicCommandHandler: if command == "unbind": self.handle_unbind(message) return + if command == "markets": + self.handle_markets(message) + return def handle_start_help(self, message: Any) -> None: trace = CommandTrace("/start", message) @@ -233,3 +243,18 @@ class BasicCommandHandler: trace.set_status("ok", "unbound") finally: trace.emit() + + def handle_markets(self, message: Any) -> None: + trace = CommandTrace("/markets", message) + try: + from src.utils.telegram_push import build_market_monitor_digest + + summary = build_market_monitor_digest( + self.config, + slot_label="当前市场概览", + force_refresh=False, + ) + self.bot.reply_to(message, summary) + trace.set_status("ok") + finally: + trace.emit() diff --git a/src/bot/io_layer.py b/src/bot/io_layer.py index 0cbe1463..2cc9c106 100644 --- a/src/bot/io_layer.py +++ b/src/bot/io_layer.py @@ -178,6 +178,7 @@ class BotIOLayer: "可用指令:\n" f"/city [城市名] 或 /pwcity [城市名] - 查询城市天气预测与实测 (消耗 {CITY_QUERY_COST} 积分)\n" f"/deb [城市名] 或 /pwdeb [城市名] - 查看 DEB 融合预测准确率 (消耗 {DEB_QUERY_COST} 积分)\n" + "/markets - 查看当前市场监控摘要\n" "/top - 查看积分排行榜\n" "/id - 获取当前聊天的 Chat ID\n\n" "/diag - 查看 Bot 启动诊断\n\n" diff --git a/src/bot/orchestrator.py b/src/bot/orchestrator.py index f8c78f8a..cc62038e 100644 --- a/src/bot/orchestrator.py +++ b/src/bot/orchestrator.py @@ -36,6 +36,7 @@ def _register_handlers( bot=bot, io_layer=io_layer, runtime_status_provider=startup_coordinator.get_runtime_status, + config=config, ).register() CityCommandHandler( bot=bot, diff --git a/src/bot/runtime_coordinator.py b/src/bot/runtime_coordinator.py index a19e42b1..4205e6d9 100644 --- a/src/bot/runtime_coordinator.py +++ b/src/bot/runtime_coordinator.py @@ -169,7 +169,7 @@ class StartupCoordinator: "focus_digest_hours": str( os.getenv("TELEGRAM_MARKET_FOCUS_DIGEST_HOURS") or "11,18" ).strip(), - "timezone": str(os.getenv("TELEGRAM_MARKET_TIMEZONE") or "Asia/Shanghai").strip(), + "schedule_basis": "server_local_time", } validation_error = None if chat_ids else "missing_TELEGRAM_CHAT_IDS" return self._start_with_validation( diff --git a/src/utils/telegram_push.py b/src/utils/telegram_push.py index 03bd4012..21e55929 100644 --- a/src/utils/telegram_push.py +++ b/src/utils/telegram_push.py @@ -4,7 +4,7 @@ import os import re import threading import time -from datetime import datetime, timedelta, timezone +from datetime import datetime, timezone from typing import Any, Dict, List, Optional, Tuple from loguru import logger @@ -282,16 +282,75 @@ def _parse_hour_list(raw: Optional[str], default: List[int]) -> List[int]: return sorted(out) or default -def _resolve_market_timezone(timezone_name: str): - normalized = str(timezone_name or "").strip() or "Asia/Shanghai" - if normalized == "Asia/Shanghai": - return timezone(timedelta(hours=8), name="Asia/Shanghai") +def _minute_of_day(text: Optional[str]) -> Optional[int]: + raw = str(text or "").strip() + if not raw or ":" not in raw: + return None try: - from zoneinfo import ZoneInfo - - return ZoneInfo(normalized) + hour_s, minute_s = raw.split(":", 1) + hour = int(hour_s) + minute = int(minute_s[:2]) except Exception: - return timezone.utc + return None + if not (0 <= hour <= 23 and 0 <= minute <= 59): + return None + return hour * 60 + minute + + +def _format_minutes_window(delta_minutes: int) -> str: + total = abs(int(delta_minutes)) + hours = total // 60 + minutes = total % 60 + if hours > 0 and minutes > 0: + text = f"{hours}h{minutes:02d}m" + elif hours > 0: + text = f"{hours}h" + else: + text = f"{minutes}m" + return text + + +def _local_peak_context(alert_payload: Dict[str, Any]) -> Dict[str, Any]: + evidence = alert_payload.get("evidence") or {} + generated_local_time = str(evidence.get("generated_local_time") or "").strip() + trigger_summary = evidence.get("trigger_summary") or {} + suppression_snapshot = trigger_summary.get("suppression_snapshot") or {} + peak_time = str(suppression_snapshot.get("max_temp_time") or "").strip() + + local_min = _minute_of_day(generated_local_time) + peak_min = _minute_of_day(peak_time) + if local_min is None or peak_min is None: + return { + "local_time": generated_local_time, + "peak_time": peak_time, + "minutes_to_peak": None, + "score_adjustment": 0.0, + "window_label": "", + } + + delta = peak_min - local_min + score = 0.0 + window_label = "" + if 0 <= delta <= 120: + score = 18.0 + window_label = f"峰值前 {_format_minutes_window(delta)}" + elif 120 < delta <= 360: + score = 10.0 + window_label = f"距峰值 {_format_minutes_window(delta)}" + elif -90 <= delta < 0: + score = 6.0 + window_label = f"峰值后 {_format_minutes_window(delta)}" + elif delta < -90: + score = -8.0 + window_label = "峰值已过较久" + + return { + "local_time": generated_local_time, + "peak_time": peak_time, + "minutes_to_peak": delta, + "score_adjustment": score, + "window_label": window_label, + } def _market_monitor_score(alert_payload: Dict[str, Any]) -> float: @@ -328,8 +387,18 @@ def _market_monitor_score(alert_payload: Dict[str, Any]) -> float: suppression = alert_payload.get("suppression") or {} suppressed_penalty = -20.0 if bool(suppression.get("suppressed")) else 0.0 + peak_context = _local_peak_context(alert_payload) - return max(0.0, severity_score + trigger_score + edge_score + pricing_score + confidence_score + suppressed_penalty) + return max( + 0.0, + severity_score + + trigger_score + + edge_score + + pricing_score + + confidence_score + + suppressed_penalty + + float(peak_context.get("score_adjustment") or 0.0), + ) def _priority_label(score: float) -> str: @@ -367,7 +436,6 @@ def _build_focus_digest_message( *, slot_label: str, top_n: int, - timezone_name: str, ) -> str: ranked = sorted( payloads, @@ -383,7 +451,6 @@ def _build_focus_digest_message( lines = [ f"🌐 PolyWeather 市场监控 · {slot_label}", - f"时区:{timezone_name}", "", ] @@ -399,30 +466,36 @@ def _build_focus_digest_message( or snapshot.get("top_bucket") or "--" ).strip() - yes_buy = _norm_prob(snapshot.get("yes_buy")) - if yes_buy is None: - forecast_bucket = snapshot.get("forecast_bucket") or {} - if isinstance(forecast_bucket, dict): - yes_buy = _norm_prob(forecast_bucket.get("yes_buy")) - yes_buy_text = _fmt_cents(yes_buy) or "--" - edge_percent = _safe_float(snapshot.get("edge_percent")) current_temp = _safe_float(inputs.get("current_temp")) deb_prediction = _safe_float(inputs.get("deb_prediction")) market_url = str(snapshot.get("market_url") or snapshot.get("primary_market_url") or "").strip() + peak_context = _local_peak_context(payload) score = _market_monitor_score(payload) lines.append(f"{idx}. {city_name} | {_priority_label(score)}") - lines.append( - " " - + f"桶 {bucket} | Yes {yes_buy_text} | 偏差 " - + (f"{edge_percent:+.1f}%" if edge_percent is not None else "--") - ) + lines.append(" " + f"关注桶 {bucket}") + local_time = str(peak_context.get("local_time") or "").strip() + peak_time = str(peak_context.get("peak_time") or "").strip() + window_label = str(peak_context.get("window_label") or "").strip() + if local_time or peak_time or window_label: + context_parts: List[str] = [] + if local_time: + context_parts.append(f"当地 {local_time}") + if peak_time: + context_parts.append(f"峰值参考 {peak_time}") + if window_label: + context_parts.append(window_label) + lines.append(" " + " | ".join(context_parts)) if current_temp is not None or deb_prediction is not None: lines.append( " " + (f"实测 {current_temp:.1f}°C" if current_temp is not None else "实测 --") + " | " - + (f"DEB {deb_prediction:.1f}°C" if deb_prediction is not None else "DEB --") + + ( + f"DEB 预报 {deb_prediction:.1f}°C" + if deb_prediction is not None + else "DEB 预报 --" + ) ) lines.append(f" 触发:{_focus_trigger_summary(payload)}") if market_url: @@ -439,7 +512,6 @@ def _maybe_send_focus_digest( payloads: List[Dict[str, Any]], state: Dict[str, Any], *, - timezone_name: str, digest_hours: List[int], top_n: int, grace_minutes: int, @@ -447,11 +519,7 @@ def _maybe_send_focus_digest( if not chat_ids or not payloads or not digest_hours: return False - try: - local_now = datetime.now(_resolve_market_timezone(timezone_name)) - except Exception: - local_now = datetime.now() - timezone_name = "local" + local_now = datetime.now().astimezone() eligible_slot: Optional[datetime] = None for hour in sorted(digest_hours, reverse=True): @@ -474,7 +542,6 @@ def _maybe_send_focus_digest( payloads, slot_label=slot_label, top_n=top_n, - timezone_name=timezone_name, ) if not message: return False @@ -496,15 +563,45 @@ def _maybe_send_focus_digest( digest_slots[slot_key] = int(time.time()) logger.info( - "market focus digest pushed slot={} timezone={} items={} chat_targets={}", + "market focus digest pushed slot={} items={} chat_targets={}", slot_key, - timezone_name, min(top_n, len(payloads)), sent_count, ) return True +def build_market_monitor_digest( + config: Dict[str, Any], + *, + slot_label: str = "当前概览", + top_n: Optional[int] = None, + force_refresh: bool = False, +) -> str: + cities = _parse_city_list(os.getenv("TELEGRAM_ALERT_CITIES")) + if not cities: + return "⚠️ 当前未配置 TELEGRAM_ALERT_CITIES,无法生成市场监控摘要。" + + digest_top_n = top_n if top_n is not None else max( + 3, + min(8, _env_int("TELEGRAM_MARKET_FOCUS_DIGEST_TOP_N", 5)), + ) + payloads: List[Dict[str, Any]] = [] + for city in cities: + try: + payloads.append(build_trade_alert_for_city(city, config, force_refresh=force_refresh)) + except Exception as exc: + logger.warning("market monitor digest build skipped city={} error={}", city, exc) + message = _build_focus_digest_message( + payloads, + slot_label=slot_label, + top_n=digest_top_n, + ) + if message: + return message + return "ℹ️ 当前没有可用的市场监控摘要。" + + def _severity_ok(alert_payload: Dict[str, Any], min_severity: str, min_trigger_count: int) -> bool: triggered_alerts = alert_payload.get("triggered_alerts") or [] if any(alert.get("force_push") for alert in triggered_alerts): @@ -907,8 +1004,6 @@ def start_trade_alert_push_loop(bot: Any, config: Dict[str, Any]) -> Optional[th 30, min(240, _env_int("TELEGRAM_MARKET_FOCUS_DIGEST_GRACE_MINUTES", 180)), ) - market_timezone = str(os.getenv("TELEGRAM_MARKET_TIMEZONE") or "Asia/Shanghai").strip() or "Asia/Shanghai" - def _runner() -> None: try: _save_state(state_path, _load_state(state_path)) @@ -919,7 +1014,7 @@ def start_trade_alert_push_loop(bot: Any, config: Dict[str, Any]) -> Optional[th f"cities={len(cities)} interval={interval_sec}s chat_targets={len(chat_ids)} " f"cooldown={cooldown_sec}s min_triggers={min_trigger_count} min_severity={min_severity} " f"focus_digest_enabled={focus_digest_enabled} focus_hours={focus_digest_hours} " - f"timezone={market_timezone} state_path={state_path}" + f"state_path={state_path}" ) while True: cycle_started = time.time() @@ -957,7 +1052,6 @@ def start_trade_alert_push_loop(bot: Any, config: Dict[str, Any]) -> Optional[th chat_ids=chat_ids, payloads=cycle_payloads, state=state, - timezone_name=market_timezone, digest_hours=focus_digest_hours, top_n=focus_digest_top_n, grace_minutes=focus_digest_grace_minutes, diff --git a/tests/test_bot_basic_handler.py b/tests/test_bot_basic_handler.py index db8d7248..10d88b2e 100644 --- a/tests/test_bot_basic_handler.py +++ b/tests/test_bot_basic_handler.py @@ -53,3 +53,33 @@ def test_basic_handler_diag_returns_html(): assert len(bot.replies) == 1 assert bot.replies[0]["parse_mode"] == "HTML" assert "Bot 启动诊断" in bot.replies[0]["text"] + + +def test_basic_handler_markets_returns_summary(monkeypatch): + bot = DummyBot() + io_layer = SimpleNamespace( + build_welcome_text=lambda: "WELCOME", + build_points_rank_text=lambda _user: "TOP", + ) + handler = BasicCommandHandler( + bot=bot, + io_layer=io_layer, + runtime_status_provider=lambda: RuntimeStatus( + started_at="2026-03-12 00:00:00 UTC", + loops=[], + command_access_mode="group_member", + protected_commands=["/city", "/deb"], + required_group_chat_id="-1001234567890", + ), + config={}, + ) + + monkeypatch.setattr( + "src.utils.telegram_push.build_market_monitor_digest", + lambda config, slot_label="当前概览", top_n=None, force_refresh=False: "MARKET DIGEST", + ) + + handler.handle_markets(_message("/markets")) + + assert len(bot.replies) == 1 + assert bot.replies[0]["text"] == "MARKET DIGEST"