From 4bd531bbc9fdec12f6c2ddd6558dd8d3060663d0 Mon Sep 17 00:00:00 2001
From: "2569718930@qq.com" <2569718930@qq.com>
Date: Sat, 14 Mar 2026 11:39:30 +0800
Subject: [PATCH] feat: implement PolyWeather Telegram bot with weather and DEB
analysis, a point system, and configuration.
---
.env.example | 9 +++-
config/config.yaml | 5 ++
src/bot/handlers/city.py | 26 ++++++---
src/bot/handlers/deb.py | 21 +++++---
src/bot/io_layer.py | 54 ++++++++++++++++++-
src/bot/orchestrator.py | 14 ++++-
src/data_collection/city_registry.py | 16 ++++++
src/data_collection/weather_sources.py | 3 ++
.../polymarket_wallet_activity_watcher.py | 22 +++++++-
9 files changed, 151 insertions(+), 19 deletions(-)
diff --git a/.env.example b/.env.example
index 7d11cfdc..f46ffe74 100644
--- a/.env.example
+++ b/.env.example
@@ -4,6 +4,10 @@ TELEGRAM_CHAT_ID=your_chat_id_here
# Optional multi-chat target (comma-separated). If set, it will be merged with TELEGRAM_CHAT_ID.
# Example: TELEGRAM_CHAT_IDS=-1003586303099,-1003539418691
TELEGRAM_CHAT_IDS=
+# Optional: route /city and /deb outputs to a fixed forum topic.
+# If TELEGRAM_QUERY_TOPIC_CHAT_ID is empty, fallback to command source chat.
+TELEGRAM_QUERY_TOPIC_CHAT_ID=
+TELEGRAM_QUERY_TOPIC_ID=
TELEGRAM_ALERT_PUSH_ENABLED=true
TELEGRAM_ALERT_PUSH_INTERVAL_SEC=300
TELEGRAM_ALERT_PUSH_COOLDOWN_SEC=1800
@@ -11,7 +15,7 @@ TELEGRAM_ALERT_MIN_TRIGGER_COUNT=2
TELEGRAM_ALERT_MIN_SEVERITY=medium
# Mispricing radar: skip push when YES buy price is above this cap (10c = 0.10)
TELEGRAM_ALERT_MISPRICING_MAX_YES_BUY=0.10
-TELEGRAM_ALERT_CITIES=ankara,london,paris,seoul,hong kong,shanghai,singapore,tokyo,toronto,buenos aires,wellington,new york,chicago,dallas,miami,atlanta,seattle,lucknow,sao paulo,munich
+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
# AI
GROQ_API_KEY=your_groq_api_key_here
@@ -145,6 +149,9 @@ POLYMARKET_WALLET_ACTIVITY_USERS=0x0000000000000000000000000000000000000000
# If unset, fallback to TELEGRAM_CHAT_IDS / TELEGRAM_CHAT_ID.
POLYMARKET_WALLET_ACTIVITY_CHAT_ID=
POLYMARKET_WALLET_ACTIVITY_CHAT_IDS=
+# Optional: mirror wallet activity push to a forum topic, while keeping existing chat targets unchanged.
+POLYMARKET_WALLET_ACTIVITY_TOPIC_CHAT_ID=
+POLYMARKET_WALLET_ACTIVITY_TOPIC_ID=
# Optional wallet nicknames:
# - CSV: 0xabc...=Whale_A,0xdef...=Main_Account
# - JSON: {"0xabc...":"Whale A","0xdef...":"Main Account"}
diff --git a/config/config.yaml b/config/config.yaml
index 27b73352..6e3df6ef 100644
--- a/config/config.yaml
+++ b/config/config.yaml
@@ -64,6 +64,11 @@ cities:
country: "Japan"
latitude: 35.5523
longitude: 139.7798
+ - id: "tel_aviv"
+ city: "Tel Aviv"
+ country: "Israel"
+ latitude: 32.0114
+ longitude: 34.8867
# Logging
logging:
level: "INFO"
diff --git a/src/bot/handlers/city.py b/src/bot/handlers/city.py
index 1e429b48..f1c4d99a 100644
--- a/src/bot/handlers/city.py
+++ b/src/bot/handlers/city.py
@@ -5,6 +5,7 @@ from typing import Any
from loguru import logger
from src.bot.command_guard import CommandGuard
+from src.bot.io_layer import BotIOLayer
from src.bot.observability import CommandTrace
from src.bot.services.city_command_service import CityCommandService
from src.bot.settings import CITY_QUERY_COST
@@ -16,10 +17,12 @@ class CityCommandHandler:
bot: Any,
guard: CommandGuard,
city_service: CityCommandService,
+ io_layer: BotIOLayer,
):
self.bot = bot
self.guard = guard
self.city_service = city_service
+ self.io_layer = io_layer
def register(self) -> None:
@self.bot.message_handler(commands=["city"])
@@ -32,7 +35,7 @@ class CityCommandHandler:
parts = (message.text or "").split(maxsplit=1)
if len(parts) < 2:
trace.set_status("bad_request", "missing_city")
- self.bot.reply_to(
+ self.io_layer.send_query_message(
message,
"❌ 请输入城市名称\n\n用法: /city chicago",
parse_mode="HTML",
@@ -44,7 +47,7 @@ class CityCommandHandler:
if not resolved.ok:
city_list = ", ".join(resolved.supported_cities or [])
trace.set_status("bad_request", "city_not_supported")
- self.bot.reply_to(
+ self.io_layer.send_query_message(
message,
f"❌ 未找到城市: {city_input}\n\n支持的城市: {city_list}",
parse_mode="HTML",
@@ -56,21 +59,28 @@ class CityCommandHandler:
trace.set_status("blocked", "guard_rejected")
return
- self.bot.send_message(
- message.chat.id, f"🔍 正在查询 {city_name.title()} 的天气数据..."
+ self.io_layer.send_query_message(
+ message,
+ f"🔍 正在查询 {city_name.title()} 的天气数据...",
)
report_result = self.city_service.build_report(city_name, CITY_QUERY_COST)
if not report_result.ok:
trace.set_status("failed", report_result.error or "city_report_failed")
- self.bot.reply_to(message, f"❌ 查询失败: {report_result.error}")
+ self.io_layer.send_query_message(
+ message,
+ f"❌ 查询失败: {report_result.error}",
+ )
return
- self.bot.send_message(message.chat.id, str(report_result.report), parse_mode="HTML")
+ self.io_layer.send_query_message(
+ message,
+ str(report_result.report),
+ parse_mode="HTML",
+ )
trace.set_status("ok", city_name)
except Exception as exc:
trace.set_status("failed", "unexpected_error")
logger.exception("查询 /city 失败")
- self.bot.reply_to(message, f"❌ 查询失败: {exc}")
+ self.io_layer.send_query_message(message, f"❌ 查询失败: {exc}")
finally:
trace.emit()
-
diff --git a/src/bot/handlers/deb.py b/src/bot/handlers/deb.py
index 0056d3aa..fe087111 100644
--- a/src/bot/handlers/deb.py
+++ b/src/bot/handlers/deb.py
@@ -5,6 +5,7 @@ from typing import Any
from loguru import logger
from src.bot.command_guard import CommandGuard
+from src.bot.io_layer import BotIOLayer
from src.bot.observability import CommandTrace
from src.bot.services.deb_command_service import DebCommandService
from src.bot.settings import DEB_QUERY_COST
@@ -16,10 +17,12 @@ class DebCommandHandler:
bot: Any,
guard: CommandGuard,
deb_service: DebCommandService,
+ io_layer: BotIOLayer,
):
self.bot = bot
self.guard = guard
self.deb_service = deb_service
+ self.io_layer = io_layer
def register(self) -> None:
@self.bot.message_handler(commands=["deb"])
@@ -32,7 +35,7 @@ class DebCommandHandler:
parts = (message.text or "").split(maxsplit=1)
if len(parts) < 2:
trace.set_status("bad_request", "missing_city")
- self.bot.reply_to(
+ self.io_layer.send_query_message(
message,
"❌ 用法: /deb ankara",
parse_mode="HTML",
@@ -43,7 +46,7 @@ class DebCommandHandler:
city_name = self.deb_service.resolve_city(city_input)
if not self.deb_service.has_history(city_name):
trace.set_status("bad_request", "history_missing")
- self.bot.reply_to(
+ self.io_layer.send_query_message(
message,
f"❌ 暂无 {city_name} 的历史数据。",
parse_mode="HTML",
@@ -57,15 +60,21 @@ class DebCommandHandler:
report_result = self.deb_service.build_report(city_name, DEB_QUERY_COST)
if not report_result.ok:
trace.set_status("failed", report_result.error or "deb_report_failed")
- self.bot.reply_to(message, f"❌ 查询失败: {report_result.error}")
+ self.io_layer.send_query_message(
+ message,
+ f"❌ 查询失败: {report_result.error}",
+ )
return
- self.bot.reply_to(message, str(report_result.report), parse_mode="HTML")
+ self.io_layer.send_query_message(
+ message,
+ str(report_result.report),
+ parse_mode="HTML",
+ )
trace.set_status("ok", city_name)
except Exception as exc:
trace.set_status("failed", "unexpected_error")
logger.exception("查询 /deb 失败")
- self.bot.reply_to(message, f"❌ 查询失败: {exc}")
+ self.io_layer.send_query_message(message, f"❌ 查询失败: {exc}")
finally:
trace.emit()
-
diff --git a/src/bot/io_layer.py b/src/bot/io_layer.py
index feb61eb1..4c898e8c 100644
--- a/src/bot/io_layer.py
+++ b/src/bot/io_layer.py
@@ -1,5 +1,6 @@
from __future__ import annotations
+import os
from datetime import datetime
from typing import Any
@@ -22,11 +23,62 @@ class BotIOLayer:
def __init__(self, bot: Any, db: DBManager):
self.bot = bot
self.db = db
+ self.query_topic_chat_id = str(
+ os.getenv("TELEGRAM_QUERY_TOPIC_CHAT_ID") or ""
+ ).strip()
+ self.query_topic_id = self._safe_int(
+ os.getenv("TELEGRAM_QUERY_TOPIC_ID"),
+ default=0,
+ )
@staticmethod
def display_name(user: Any) -> str:
return user.username or user.first_name or f"User_{user.id}"
+ @staticmethod
+ def _safe_int(raw: Any, default: int = 0) -> int:
+ try:
+ return int(raw)
+ except Exception:
+ return default
+
+ def send_query_message(
+ self,
+ message: Any,
+ text: str,
+ *,
+ parse_mode: str | None = None,
+ ) -> None:
+ chat = getattr(message, "chat", None)
+ fallback_chat_id = getattr(chat, "id", None)
+ target_chat_id = self.query_topic_chat_id or fallback_chat_id
+ if target_chat_id is None:
+ self.bot.send_message(message.chat.id, text, parse_mode=parse_mode)
+ return
+
+ kwargs = {}
+ if parse_mode:
+ kwargs["parse_mode"] = parse_mode
+ if self.query_topic_id > 0:
+ kwargs["message_thread_id"] = self.query_topic_id
+
+ try:
+ self.bot.send_message(target_chat_id, text, **kwargs)
+ return
+ except Exception as exc:
+ logger.warning(
+ "query topic send failed chat_id={} topic_id={} error={}",
+ target_chat_id,
+ self.query_topic_id,
+ exc,
+ )
+
+ # Fallback: drop topic and send to source chat.
+ if fallback_chat_id is not None:
+ fallback_kwargs = dict(kwargs)
+ fallback_kwargs.pop("message_thread_id", None)
+ self.bot.send_message(fallback_chat_id, text, **fallback_kwargs)
+
def ensure_query_points(self, message: Any, cost: int, label: str) -> bool:
user = message.from_user
self.db.upsert_user(user.id, self.display_name(user))
@@ -37,7 +89,7 @@ class BotIOLayer:
balance = int(result.get("balance") or 0)
required = int(result.get("required") or cost)
missing = max(0, required - balance)
- self.bot.reply_to(
+ self.send_query_message(
message,
(
f"❌ 积分不足,无法执行 {label}\n"
diff --git a/src/bot/orchestrator.py b/src/bot/orchestrator.py
index 44866553..9507e0a7 100644
--- a/src/bot/orchestrator.py
+++ b/src/bot/orchestrator.py
@@ -36,8 +36,18 @@ def _register_handlers(
io_layer=io_layer,
runtime_status_provider=startup_coordinator.get_runtime_status,
).register()
- CityCommandHandler(bot=bot, guard=guard, city_service=city_service).register()
- DebCommandHandler(bot=bot, guard=guard, deb_service=deb_service).register()
+ CityCommandHandler(
+ bot=bot,
+ guard=guard,
+ city_service=city_service,
+ io_layer=io_layer,
+ ).register()
+ DebCommandHandler(
+ bot=bot,
+ guard=guard,
+ deb_service=deb_service,
+ io_layer=io_layer,
+ ).register()
ActivityHandler(bot=bot, io_layer=io_layer).register()
diff --git a/src/data_collection/city_registry.py b/src/data_collection/city_registry.py
index 55ea9f9c..7d908b00 100644
--- a/src/data_collection/city_registry.py
+++ b/src/data_collection/city_registry.py
@@ -114,6 +114,20 @@ CITY_REGISTRY = {
"distance_km": 15.0,
"warning": "湾风与热岛效应并存,日最高温时点有时会后移。",
},
+ "tel aviv": {
+ "name": "Tel Aviv",
+ "lat": 32.0114,
+ "lon": 34.8867,
+ "icao": "LLBG",
+ "tz_offset": 7200,
+ "use_fahrenheit": False,
+ "is_major": True,
+ "risk_level": "medium",
+ "risk_emoji": "🟡",
+ "airport_name": "Ben Gurion 机场",
+ "distance_km": 14.8,
+ "warning": "海陆风转换明显,午后海风可压制升温,最高温时点可能后移。",
+ },
"toronto": {
"name": "Toronto",
"lat": 43.6777,
@@ -293,6 +307,7 @@ ALIASES = {
"seo": "seoul", "hkg": "hong kong", "hk": "hong kong",
"sha": "shanghai", "sh": "shanghai", "sin": "singapore",
"sg": "singapore", "tok": "tokyo", "tyo": "tokyo",
+ "tlv": "tel aviv", "telaviv": "tel aviv",
"ba": "buenos aires", "wel": "wellington",
"luc": "lucknow", "sp": "sao paulo", "mun": "munich",
@@ -313,6 +328,7 @@ ALIASES = {
"新加坡": "singapore",
"东京": "tokyo",
"東京": "tokyo",
+ "特拉维夫": "tel aviv",
"布宜诺斯艾利斯": "buenos aires",
"惠灵顿": "wellington",
"勒克瑙": "lucknow",
diff --git a/src/data_collection/weather_sources.py b/src/data_collection/weather_sources.py
index 624d5e7e..ef72b92a 100644
--- a/src/data_collection/weather_sources.py
+++ b/src/data_collection/weather_sources.py
@@ -36,6 +36,7 @@ class WeatherDataCollector:
"shanghai": ["ZSPD", "ZSSS", "ZSNB", "ZSHC"],
"singapore": ["WSSS", "WSAP", "WMKK"],
"tokyo": ["RJTT", "RJAA", "RJAH", "RJTJ"],
+ "tel aviv": ["LLBG"],
"toronto": ["CYYZ", "CYTZ", "CYKF"],
"chicago": ["KORD", "KMDW", "KPWK", "KDPA"],
"dallas": ["KDAL", "KDFW", "KADS", "KGKY"],
@@ -1670,6 +1671,8 @@ class WeatherDataCollector:
"tokyo": "Tokyo",
"东京": "Tokyo",
"東京": "Tokyo",
+ "tel aviv": "Tel Aviv",
+ "特拉维夫": "Tel Aviv",
"toronto": "Toronto",
"多伦多": "Toronto",
"ankara": "Ankara",
diff --git a/src/onchain/polymarket_wallet_activity_watcher.py b/src/onchain/polymarket_wallet_activity_watcher.py
index 74b75feb..9aae0908 100644
--- a/src/onchain/polymarket_wallet_activity_watcher.py
+++ b/src/onchain/polymarket_wallet_activity_watcher.py
@@ -580,6 +580,16 @@ def _filter_changes_by_position_value(
def start_polymarket_wallet_activity_loop(bot: Any) -> Optional[threading.Thread]:
enabled = _env_bool("POLYMARKET_WALLET_ACTIVITY_ENABLED", False)
chat_ids = get_polymarket_wallet_activity_chat_ids_from_env()
+ topic_chat_id = str(
+ os.getenv("POLYMARKET_WALLET_ACTIVITY_TOPIC_CHAT_ID") or ""
+ ).strip()
+ topic_thread_id = max(
+ 0,
+ _env_int("POLYMARKET_WALLET_ACTIVITY_TOPIC_ID", 0),
+ )
+ if topic_chat_id and topic_chat_id not in chat_ids:
+ # Mirror wallet activity push to a forum topic without replacing existing channels.
+ chat_ids = [*chat_ids, topic_chat_id]
users = _parse_addresses(os.getenv("POLYMARKET_WALLET_ACTIVITY_USERS"))
user_aliases = _parse_address_aliases(
os.getenv("POLYMARKET_WALLET_ACTIVITY_USER_ALIASES")
@@ -670,6 +680,7 @@ def start_polymarket_wallet_activity_loop(bot: Any) -> Optional[threading.Thread
f"min_position_value_usd={min_position_value_usd} "
f"min_value_exempt_users={len(exempt_wallets)} "
f"chat_targets={len(chat_ids)} "
+ f"topic_chat_id={topic_chat_id or '-'} topic_id={topic_thread_id} "
f"aliases={len(user_aliases)} link_preview={link_preview} "
f"min_avg_price_delta={min_avg_price_delta} "
f"immediate_on_size_delta={immediate_on_size_delta} "
@@ -841,10 +852,19 @@ def start_polymarket_wallet_activity_loop(bot: Any) -> Optional[threading.Thread
sent_count = 0
for chat_id in chat_ids:
try:
+ send_kwargs = {
+ "disable_web_page_preview": not link_preview,
+ }
+ if (
+ topic_thread_id > 0
+ and topic_chat_id
+ and str(chat_id).strip() == topic_chat_id
+ ):
+ send_kwargs["message_thread_id"] = topic_thread_id
bot.send_message(
chat_id,
msg,
- disable_web_page_preview=not link_preview,
+ **send_kwargs,
)
sent_count += 1
except Exception as exc: