feat: implement PolyWeather Telegram bot with weather and DEB analysis, a point system, and configuration.

This commit is contained in:
2569718930@qq.com
2026-03-14 11:39:30 +08:00
parent 64f8b37014
commit 4bd531bbc9
9 changed files with 151 additions and 19 deletions
+8 -1
View File
@@ -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"}
+5
View File
@@ -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"
+18 -8
View File
@@ -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用法: <code>/city chicago</code>",
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"❌ 未找到城市: <b>{city_input}</b>\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()
+15 -6
View File
@@ -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,
"❌ 用法: <code>/deb ankara</code>",
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()
+53 -1
View File
@@ -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"❌ 积分不足,无法执行 <b>{label}</b>\n"
+12 -2
View File
@@ -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()
+16
View File
@@ -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",
+3
View File
@@ -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",
@@ -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: