From fb1c88463fe83ddd959cd033317d49a92968318a Mon Sep 17 00:00:00 2001
From: "2569718930@qq.com" <2569718930@qq.com>
Date: Tue, 10 Mar 2026 11:51:43 +0800
Subject: [PATCH] feat: Add multi-source weather data collection
(OpenWeatherMap, Meteoblue, Visual Crossing, METAR) and Polygon wallet
watching functionality, introducing new configuration variables.
---
.env.example | 18 +-
bot_listener.py | 369 +++++++++++++------------
src/data_collection/weather_sources.py | 40 ++-
src/onchain/__init__.py | 1 +
src/onchain/polygon_wallet_watcher.py | 354 ++++++++++++++++++++++++
5 files changed, 597 insertions(+), 185 deletions(-)
create mode 100644 src/onchain/__init__.py
create mode 100644 src/onchain/polygon_wallet_watcher.py
diff --git a/.env.example b/.env.example
index 57fe23d9..de4f2981 100644
--- a/.env.example
+++ b/.env.example
@@ -1,4 +1,4 @@
-# Telegram Bot
+# Telegram Bot
TELEGRAM_BOT_TOKEN=your_bot_token_here
TELEGRAM_CHAT_ID=your_chat_id_here
TELEGRAM_ALERT_PUSH_ENABLED=true
@@ -11,6 +11,10 @@ TELEGRAM_ALERT_CITIES=ankara,london,paris,seoul,toronto,buenos aires,wellington,
# AI
GROQ_API_KEY=your_groq_api_key_here
+# Meteoblue
+METEOBLUE_API_KEY=your_meteoblue_api_key_here
+METEOBLUE_CACHE_TTL_SEC=1800
+
# Proxy Setting (optional)
HTTPS_PROXY=http://127.0.0.1:7890
HTTP_PROXY=http://127.0.0.1:7890
@@ -32,3 +36,15 @@ POLYMARKET_DISCOVERY_PAGES=6
POLYMARKET_DISCOVERY_LIMIT=200
POLYMARKET_SIGNAL_MIN_LIQUIDITY=500
POLYMARKET_SIGNAL_EDGE_PCT=2
+
+# Polygon Wallet Watcher (Single Chain P0)
+POLYGON_WALLET_WATCH_ENABLED=false
+POLYGON_RPC_URL=https://polygon-rpc.com
+POLYGON_WALLET_WATCH_ADDRESSES=0x0000000000000000000000000000000000000000
+POLYGON_WALLET_WATCH_INTERVAL_SEC=8
+POLYGON_WALLET_WATCH_CONFIRMATIONS=2
+POLYGON_WALLET_WATCH_MAX_BLOCKS_PER_CYCLE=30
+POLYGON_WALLET_WATCH_SEEN_TTL_SEC=604800
+POLYGON_WALLET_WATCH_RPC_TIMEOUT_SEC=10
+POLYGON_WALLET_WATCH_TX_BASE=https://polygonscan.com/tx
+POLYGON_WALLET_WATCH_ADDR_BASE=https://polygonscan.com/address
diff --git a/bot_listener.py b/bot_listener.py
index 85375064..0117a08a 100644
--- a/bot_listener.py
+++ b/bot_listener.py
@@ -1,16 +1,17 @@
-import sys
+import sys
import os
from typing import List
import telebot # type: ignore
from loguru import logger # type: ignore
-# 确保项目根目录在 sys.path 中
+# 纭繚椤圭洰鏍圭洰褰曞湪 sys.path 涓?
project_root = os.path.dirname(os.path.abspath(__file__))
if project_root not in sys.path:
sys.path.insert(0, project_root)
from src.utils.config_loader import load_config # type: ignore # noqa: E402
from src.utils.telegram_push import start_trade_alert_push_loop # type: ignore # noqa: E402
+from src.onchain.polygon_wallet_watcher import start_polygon_wallet_watch_loop # type: ignore # noqa: E402
from src.data_collection.weather_sources import WeatherDataCollector # type: ignore # noqa: E402
from src.data_collection.city_risk_profiles import get_city_risk_profile # type: ignore # noqa: E402
from src.analysis.deb_algorithm import calculate_dynamic_weights, update_daily_record # noqa: E402
@@ -25,7 +26,7 @@ DEB_QUERY_COST = 1
def analyze_weather_trend(weather_data, temp_symbol, city_name=None):
- """Thin wrapper — delegates to shared trend_engine module."""
+ """Thin wrapper 鈥?delegates to shared trend_engine module."""
from src.analysis.trend_engine import analyze_weather_trend as _analyze
display_str, ai_context, _structured = _analyze(weather_data, temp_symbol, city_name)
return display_str, ai_context
@@ -35,13 +36,14 @@ def start_bot():
config = load_config()
token = os.getenv("TELEGRAM_BOT_TOKEN")
if not token:
- logger.error("未找到 TELEGRAM_BOT_TOKEN 环境变量")
+ logger.error("鏈壘鍒?TELEGRAM_BOT_TOKEN 鐜鍙橀噺")
return
bot = telebot.TeleBot(token)
db = DBManager()
weather = WeatherDataCollector(config)
start_trade_alert_push_loop(bot, config)
+ start_polygon_wallet_watch_loop(bot)
def _display_name(user) -> str:
return user.username or user.first_name or f"User_{user.id}"
@@ -59,12 +61,12 @@ def start_bot():
bot.reply_to(
message,
(
- f"❌ 积分不足,无法执行 {label}\n"
- f"当前积分: {balance}\n"
- f"需要积分: {required}\n"
- f"还差积分: {missing}\n\n"
- f"积分规则:每日签到(有效发言满 {MESSAGE_MIN_LENGTH} 字)获得 {MESSAGE_POINTS} 积分,"
- f"每日上限 {MESSAGE_DAILY_CAP} 分。"
+ f"鉂?绉垎涓嶈冻锛屾棤娉曟墽琛?{label}\n"
+ f"褰撳墠绉垎: {balance}\n"
+ f"闇€瑕佺Н鍒? {required}\n"
+ f"杩樺樊绉垎: {missing}\n\n"
+ f"绉垎瑙勫垯锛氭瘡鏃ョ鍒?鏈夋晥鍙戣█婊?{MESSAGE_MIN_LENGTH} 瀛?鑾峰緱 {MESSAGE_POINTS} 绉垎锛?
+ f"姣忔棩涓婇檺 {MESSAGE_DAILY_CAP} 鍒嗐€?
),
parse_mode="HTML",
)
@@ -73,15 +75,15 @@ def start_bot():
@bot.message_handler(commands=["start", "help"])
def send_welcome(message):
welcome_text = (
- "🌡️ PolyWeather 天气查询机器人\n\n"
- "可用指令:\n"
- f"/city [城市名] - 查询城市天气预测与实测 (消耗 {CITY_QUERY_COST} 积分)\n"
- f"/deb [城市名] - 查看 DEB 融合预测准确率 (消耗 {DEB_QUERY_COST} 积分)\n"
- "/top - 查看积分排行榜\n"
- "/id - 获取当前聊天的 Chat ID\n\n"
- "示例: /city 伦敦\n"
- f"💡 提示: 每日签到(有效发言满 {MESSAGE_MIN_LENGTH} 字)获得 {MESSAGE_POINTS} 积分,"
- f"每日上限 {MESSAGE_DAILY_CAP} 分。"
+ "馃尅锔?PolyWeather 澶╂皵鏌ヨ鏈哄櫒浜?/b>\n\n"
+ "鍙敤鎸囦护:\n"
+ f"/city [鍩庡競鍚峕 - 鏌ヨ鍩庡競澶╂皵棰勬祴涓庡疄娴?(娑堣€?{CITY_QUERY_COST} 绉垎)\n"
+ f"/deb [鍩庡競鍚峕 - 鏌ョ湅 DEB 铻嶅悎棰勬祴鍑嗙‘鐜?(娑堣€?{DEB_QUERY_COST} 绉垎)\n"
+ "/top - 鏌ョ湅绉垎鎺掕姒淺n"
+ "/id - 鑾峰彇褰撳墠鑱婂ぉ鐨?Chat ID\n\n"
+ "绀轰緥: /city 浼︽暒\n"
+ f"馃挕 鎻愮ず: 姣忔棩绛惧埌(鏈夋晥鍙戣█婊?{MESSAGE_MIN_LENGTH} 瀛?鑾峰緱 {MESSAGE_POINTS} 绉垎锛?
+ f"姣忔棩涓婇檺 {MESSAGE_DAILY_CAP} 鍒嗐€?/i>"
)
bot.reply_to(message, welcome_text, parse_mode="HTML")
@@ -89,44 +91,44 @@ def start_bot():
def get_chat_id(message):
bot.reply_to(
message,
- f"🎯 当前聊天的 Chat ID 是: {message.chat.id}",
+ f"馃幆 褰撳墠鑱婂ぉ鐨?Chat ID 鏄? {message.chat.id}",
parse_mode="HTML",
)
@bot.message_handler(commands=["top"])
def show_points(message):
- """显示当前用户的积分及排行榜"""
+ """鏄剧ず褰撳墠鐢ㄦ埛鐨勭Н鍒嗗強鎺掕姒?""
user = message.from_user
db.upsert_user(user.id, _display_name(user))
user_info = db.get_user(user.id)
leaderboard = db.get_leaderboard(limit=5)
- rank_text = "🏆 PolyWeather 活跃度排行榜\n"
- rank_text += "────────────────────\n"
+ rank_text = "馃弳 PolyWeather 娲昏穬搴︽帓琛屾\n"
+ rank_text += "鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€\n"
for i, entry in enumerate(leaderboard):
- medal = ["🥇", "🥈", "🥉", " ", " "][i] if i < 5 else " "
- rank_text += f"{medal} {entry['username'][:12]}: {entry['points']} 点\n"
+ medal = ["馃", "馃", "馃", " ", " "][i] if i < 5 else " "
+ rank_text += f"{medal} {entry['username'][:12]}: {entry['points']} 鐐筡n"
if user_info:
- rank_text += "────────────────────\n"
+ rank_text += "鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€鈹€\n"
rank_text += (
- f"👤 我的状态:\n"
- f"└ 积分: {user_info['points']}\n"
- f"└ 发言: {user_info['message_count']} 次\n"
- f"└ 今日发言积分: {user_info.get('daily_points') or 0}/{MESSAGE_DAILY_CAP}\n"
- f"└ /city 消耗: {CITY_QUERY_COST} | /deb 消耗: {DEB_QUERY_COST}"
+ f"馃懁 鎴戠殑鐘舵€侊細\n"
+ f"鈹?绉垎: {user_info['points']}\n"
+ f"鈹?鍙戣█: {user_info['message_count']} 娆n"
+ f"鈹?浠婃棩鍙戣█绉垎: {user_info.get('daily_points') or 0}/{MESSAGE_DAILY_CAP}\n"
+ f"鈹?/city 娑堣€? {CITY_QUERY_COST} | /deb 娑堣€? {DEB_QUERY_COST}"
)
bot.send_message(message.chat.id, rank_text, parse_mode="HTML")
@bot.message_handler(commands=["deb"])
def deb_accuracy(message):
- """查询 DEB 融合预测的近 7 天准确率。"""
+ """鏌ヨ DEB 铻嶅悎棰勬祴鐨勮繎 7 澶╁噯纭巼銆?""
try:
parts = message.text.split(maxsplit=1)
if len(parts) < 2:
bot.reply_to(
message,
- "❌ 用法: /deb ankara",
+ "鉂?鐢ㄦ硶: /deb ankara",
parse_mode="HTML",
)
return
@@ -147,7 +149,7 @@ def start_bot():
if city_name not in data or not data[city_name]:
bot.reply_to(
message,
- f"❌ 暂无 {city_name} 的历史数据。",
+ f"鉂?鏆傛棤 {city_name} 鐨勫巻鍙叉暟鎹€?,
parse_mode="HTML",
)
return
@@ -172,9 +174,9 @@ def start_bot():
recent_items.sort(key=lambda item: item[0])
lines = [
- f"📊 DEB 准确率报告 - {city_name.title()}",
+ f"馃搳 DEB 鍑嗙‘鐜囨姤鍛?- {city_name.title()}",
"",
- "📅 近7日记录:",
+ "馃搮 杩?鏃ヨ褰曪細",
]
total_days = 0
hits = 0
@@ -205,7 +207,7 @@ def start_bot():
actual_wu = round(actual)
if date_str == today_str:
- lines.append(f" {date_str}: 📍 今天进行中 (实测暂 {actual:.1f})")
+ lines.append(f" {date_str}: 馃搷 浠婂ぉ杩涜涓?(瀹炴祴鏆?{actual:.1f})")
elif deb_pred is not None:
total_days += 1
deb_wu = round(deb_pred)
@@ -217,18 +219,18 @@ def start_bot():
signed_errors.append(err)
if hit:
- result_icon = "✅"
- err_text = f"偏差{abs(err):.1f}°"
+ result_icon = "鉁?
+ err_text = f"鍋忓樊{abs(err):.1f}掳"
elif err < 0:
- result_icon = "❌"
- err_text = f"低估{abs(err):.1f}°"
+ result_icon = "鉂?
+ err_text = f"浣庝及{abs(err):.1f}掳"
else:
- result_icon = "❌"
- err_text = f"高估{abs(err):.1f}°"
+ result_icon = "鉂?
+ err_text = f"楂樹及{abs(err):.1f}掳"
- retro = "≈" if "deb_prediction" not in record else ""
+ retro = "鈮? if "deb_prediction" not in record else ""
lines.append(
- f" {date_str}: DEB {retro}{deb_pred:.1f}→{deb_wu} vs 实测 {actual:.1f}→{actual_wu} "
+ f" {date_str}: DEB {retro}{deb_pred:.1f}鈫抺deb_wu} vs 瀹炴祴 {actual:.1f}鈫抺actual_wu} "
f"{result_icon} {err_text}"
)
@@ -246,18 +248,18 @@ def start_bot():
deb_mae = sum(deb_errors) / len(deb_errors)
lines.append("")
lines.append(
- f"🎯 DEB 总战绩:WU命中 {hits}/{total_days} ({hit_rate:.0f}%) | MAE: {deb_mae:.1f}°"
+ f"馃幆 DEB 鎬绘垬缁╋細WU鍛戒腑 {hits}/{total_days} ({hit_rate:.0f}%) | MAE: {deb_mae:.1f}掳"
)
if model_errors:
lines.append("")
- lines.append("📈 模型 MAE 对比:")
+ lines.append("馃搱 妯″瀷 MAE 瀵规瘮锛?/b>")
model_maes = {m: sum(e) / len(e) for m, e in model_errors.items() if e}
sorted_models = sorted(model_maes.items(), key=lambda item: item[1])
for model, mae in sorted_models:
- tag = " ⭐" if mae <= deb_mae else ""
- lines.append(f" {model}: {mae:.1f}°{tag}")
- lines.append(f" DEB融合: {deb_mae:.1f}°")
+ tag = " 猸? if mae <= deb_mae else ""
+ lines.append(f" {model}: {mae:.1f}掳{tag}")
+ lines.append(f" DEB铻嶅悎: {deb_mae:.1f}掳")
mean_bias = sum(signed_errors) / len(signed_errors)
underest = sum(1 for e in signed_errors if e < -0.3)
@@ -265,55 +267,55 @@ def start_bot():
accurate = total_days - underest - overest
lines.append("")
- lines.append("🔍 偏差分析:")
+ lines.append("馃攳 鍋忓樊鍒嗘瀽锛?/b>")
if abs(mean_bias) > 0.3:
- bias_label = "系统性低估" if mean_bias < 0 else "系统性高估"
- lines.append(f" ⚠️ {bias_label}:平均偏差 {mean_bias:+.1f}°")
+ bias_label = "绯荤粺鎬т綆浼? if mean_bias < 0 else "绯荤粺鎬ч珮浼?
+ lines.append(f" 鈿狅笍 {bias_label}锛氬钩鍧囧亸宸?{mean_bias:+.1f}掳")
else:
- lines.append(f" ✅ 整体无明显系统偏差:平均偏差 {mean_bias:+.1f}°")
- lines.append(f" 低估 {underest} 次 | 高估 {overest} 次 | 准确 {accurate} 次")
+ lines.append(f" 鉁?鏁翠綋鏃犳槑鏄剧郴缁熷亸宸細骞冲潎鍋忓樊 {mean_bias:+.1f}掳")
+ lines.append(f" 浣庝及 {underest} 娆?| 楂樹及 {overest} 娆?| 鍑嗙‘ {accurate} 娆?)
lines.append("")
- lines.append("💡 建议:")
+ lines.append("馃挕 寤鸿锛?/b>")
if underest > overest and abs(mean_bias) > 0.5:
lines.append(
- f" 该城市模型集体低估趋势明显({mean_bias:+.1f}°),实际最高温可能比 DEB 融合值高 "
- f"{abs(mean_bias):.0f}-{abs(mean_bias) + 0.5:.0f}°。交易时建议适当看高。"
+ f" 璇ュ煄甯傛ā鍨嬮泦浣撲綆浼拌秼鍔挎槑鏄撅紙{mean_bias:+.1f}掳锛夛紝瀹為檯鏈€楂樻俯鍙兘姣?DEB 铻嶅悎鍊奸珮 "
+ f"{abs(mean_bias):.0f}-{abs(mean_bias) + 0.5:.0f}掳銆備氦鏄撴椂寤鸿閫傚綋鐪嬮珮銆?
)
elif overest > underest and abs(mean_bias) > 0.5:
lines.append(
- f" 该城市模型集体高估趋势明显({mean_bias:+.1f}°),实际最高温可能低于 DEB 融合值。交易时注意追高风险。"
+ f" 璇ュ煄甯傛ā鍨嬮泦浣撻珮浼拌秼鍔挎槑鏄撅紙{mean_bias:+.1f}掳锛夛紝瀹為檯鏈€楂樻俯鍙兘浣庝簬 DEB 铻嶅悎鍊笺€備氦鏄撴椂娉ㄦ剰杩介珮椋庨櫓銆?
)
elif deb_mae > 1.5:
- lines.append(f" 近期模型波动较大(MAE {deb_mae:.1f}°),建议降低对单一日预测的信任度。")
+ lines.append(f" 杩戞湡妯″瀷娉㈠姩杈冨ぇ锛圡AE {deb_mae:.1f}掳锛夛紝寤鸿闄嶄綆瀵瑰崟涓€鏃ラ娴嬬殑淇′换搴︺€?)
elif hit_rate >= 60:
- lines.append(" DEB 近期表现稳定,可继续作为主要参考。")
+ lines.append(" DEB 杩戞湡琛ㄧ幇绋冲畾锛屽彲缁х画浣滀负涓昏鍙傝€冦€?)
else:
- lines.append(" 近期准确率一般,建议结合主站实测与周边站点共同判断。")
+ lines.append(" 杩戞湡鍑嗙‘鐜囦竴鑸紝寤鸿缁撳悎涓荤珯瀹炴祴涓庡懆杈圭珯鐐瑰叡鍚屽垽鏂€?)
lines.append("")
- lines.append("📝 MAE = 平均绝对误差,越小越准。⭐ = 优于 DEB 融合。")
- lines.append("🗓 统计窗口:近7天滚动样本。")
+ lines.append("馃摑 MAE = 骞冲潎缁濆璇樊锛岃秺灏忚秺鍑嗐€傗瓙 = 浼樹簬 DEB 铻嶅悎銆?)
+ lines.append("馃棑 缁熻绐楀彛锛氳繎7澶╂粴鍔ㄦ牱鏈€?)
else:
lines.append("")
- lines.append("⏳ 近7天尚无完整的 DEB 预测记录。")
+ lines.append("鈴?杩?澶╁皻鏃犲畬鏁寸殑 DEB 棰勬祴璁板綍銆?)
lines.append("")
- lines.append(f"💳 本次消耗 {DEB_QUERY_COST} 积分。")
+ lines.append(f"馃挸 鏈娑堣€?{DEB_QUERY_COST} 绉垎銆?)
bot.reply_to(message, "\n".join(lines), parse_mode="HTML")
except Exception as e:
- bot.reply_to(message, f"❌ 查询失败: {e}")
+ bot.reply_to(message, f"鉂?鏌ヨ澶辫触: {e}")
@bot.message_handler(commands=["city"])
def get_city_info(message):
- """查询指定城市的天气详情"""
+ """鏌ヨ鎸囧畾鍩庡競鐨勫ぉ姘旇鎯?""
try:
parts = message.text.split(maxsplit=1)
if len(parts) < 2:
bot.reply_to(
message,
- "❓ 请输入城市名称\n\n用法: /city chicago",
+ "鉂?璇疯緭鍏ュ煄甯傚悕绉癨n\n鐢ㄦ硶: /city chicago",
parse_mode="HTML",
)
return
@@ -321,35 +323,35 @@ def start_bot():
from src.data_collection.city_registry import ALIASES, CITY_REGISTRY
city_input = parts[1].strip().lower()
- # --- 使用统一注册表解析城市 ---
+ # --- 浣跨敤缁熶竴娉ㄥ唽琛ㄨВ鏋愬煄甯?---
SUPPORTED_CITIES = list(CITY_REGISTRY.keys())
- # 1. 第一优先级:全称或别名完全匹配
+ # 1. 绗竴浼樺厛绾э細鍏ㄧО鎴栧埆鍚嶅畬鍏ㄥ尮閰?
city_name = ALIASES.get(city_input)
if not city_name and city_input in SUPPORTED_CITIES:
city_name = city_input
- # 2. 第二优先级:前缀模糊匹配
+ # 2. 绗簩浼樺厛绾э細鍓嶇紑妯$硦鍖归厤
if not city_name and len(city_input) >= 2:
- # 搜别名
+ # 鎼滃埆鍚?
for k, v in ALIASES.items():
if k.startswith(city_input):
city_name = v
break
- # 搜城市全名
+ # 鎼滃煄甯傚叏鍚?
if not city_name:
for full_name in SUPPORTED_CITIES:
if full_name.startswith(city_input):
city_name = full_name
break
- # 3. 未找到 → 报错
+ # 3. 鏈壘鍒?鈫?鎶ラ敊
if not city_name:
city_list = ", ".join(sorted(SUPPORTED_CITIES))
bot.reply_to(
message,
- f"❌ 未找到城市: {city_input}\n\n"
- f"支持的城市: {city_list}",
+ f"鉂?鏈壘鍒板煄甯? {city_input}\n\n"
+ f"鏀寔鐨勫煄甯? {city_list}",
parse_mode="HTML",
)
return
@@ -358,12 +360,12 @@ def start_bot():
return
bot.send_message(
- message.chat.id, f"🔍 正在查询 {city_name.title()} 的天气数据..."
+ message.chat.id, f"馃攳 姝e湪鏌ヨ {city_name.title()} 鐨勫ぉ姘旀暟鎹?.."
)
coords = weather.get_coordinates(city_name)
if not coords:
- bot.reply_to(message, f"❌ 未找到城市坐标: {city_name}")
+ bot.reply_to(message, f"鉂?鏈壘鍒板煄甯傚潗鏍? {city_name}")
return
weather_data = weather.fetch_all_sources(
@@ -373,7 +375,7 @@ def start_bot():
metar = weather_data.get("metar", {})
mgm = weather_data.get("mgm") or {}
- # 数值归一化
+ # 鏁板€煎綊涓€鍖?
def _sf(v):
if v is None:
return None
@@ -383,26 +385,26 @@ def start_bot():
return None
temp_unit = open_meteo.get("unit", "celsius")
- temp_symbol = "°F" if temp_unit == "fahrenheit" else "°C"
+ temp_symbol = "掳F" if temp_unit == "fahrenheit" else "掳C"
- # --- 1. 紧凑 Header (城市 + 时间 + 风险状态) ---
+ # --- 1. 绱у噾 Header (鍩庡競 + 鏃堕棿 + 椋庨櫓鐘舵€? ---
local_time = open_meteo.get("current", {}).get("local_time", "")
time_str = local_time.split(" ")[1][:5] if " " in local_time else "N/A"
risk_profile = get_city_risk_profile(city_name)
- risk_emoji = risk_profile.get("risk_level", "⚪") if risk_profile else "⚪"
+ risk_emoji = risk_profile.get("risk_level", "鈿?) if risk_profile else "鈿?
- msg_header = f"📍 {city_name.title()} ({time_str}) {risk_emoji}"
+ msg_header = f"馃搷 {city_name.title()} ({time_str}) {risk_emoji}"
msg_lines = [msg_header]
- # --- 2. 紧凑 风险提示 ---
+ # --- 2. 绱у噾 椋庨櫓鎻愮ず ---
if risk_profile:
- bias = risk_profile.get("bias", "±0.0")
+ bias = risk_profile.get("bias", "卤0.0")
msg_lines.append(
- f"⚠️ {risk_profile.get('airport_name', '')}: {bias}{temp_symbol} | {risk_profile.get('warning', '')}"
+ f"鈿狅笍 {risk_profile.get('airport_name', '')}: {bias}{temp_symbol} | {risk_profile.get('warning', '')}"
)
- # --- 3. 紧凑 预测区 ---
+ # --- 3. 绱у噾 棰勬祴鍖?---
daily = open_meteo.get("daily", {})
dates = daily.get("time", [])[:3]
max_temps = daily.get("temperature_2m_max", [])[:3]
@@ -411,7 +413,7 @@ def start_bot():
mgm_high = _sf(mgm.get("today_high"))
mb_high = _sf(weather_data.get("meteoblue", {}).get("today_high"))
- # 今天对比
+ # 浠婂ぉ瀵规瘮
today_t = max_temps[0] if max_temps else "N/A"
comp_parts = []
sources = ["Open-Meteo"]
@@ -433,45 +435,45 @@ def start_bot():
if mgm_high is not None:
sources.append("MGM")
comp_parts.append(
- f"🇹🇷 MGM: {mgm_high:.1f}{temp_symbol}"
+ f"馃嚬馃嚪 MGM: {mgm_high:.1f}{temp_symbol}"
if isinstance(mgm_high, (int, float))
- else f"🇹🇷 MGM: {mgm_high}"
+ else f"馃嚬馃嚪 MGM: {mgm_high}"
)
- # 检查是否有显著分歧 (超过 5°F 或 2.5°C)
+ # 妫€鏌ユ槸鍚︽湁鏄捐憲鍒嗘 (瓒呰繃 5掳F 鎴?2.5掳C)
divergence_warning = ""
if mb_high is not None and max_temps:
diff = abs(mb_high - (_sf(max_temps[0]) or 0))
threshold = 5.0 if temp_unit == "fahrenheit" else 2.5
if diff > threshold:
divergence_warning = (
- f" ⚠️ 模型显著分歧 ({diff:.1f}{temp_symbol})"
+ f" 鈿狅笍 妯″瀷鏄捐憲鍒嗘 ({diff:.1f}{temp_symbol})"
)
comp_str = f" ({' | '.join(comp_parts)})" if comp_parts else ""
sources_str = " | ".join(sources)
- msg_lines.append(f"\n📊 预报 ({sources_str})")
+ msg_lines.append(f"\n馃搳 棰勬姤 ({sources_str})")
msg_lines.append(
- f"👉 今天: {today_t}{temp_symbol}{comp_str}{divergence_warning}"
+ f"馃憠 浠婂ぉ: {today_t}{temp_symbol}{comp_str}{divergence_warning}"
)
- # 明后天
+ # 鏄庡悗澶?
if len(dates) > 1:
future_forecasts = []
mgm_daily = mgm.get("daily_forecasts", {}) or {}
for d, t in zip(dates[1:], max_temps[1:]):
- # 检查 MGM 是否有该日期的预报
+ # 妫€鏌?MGM 鏄惁鏈夎鏃ユ湡鐨勯鎶?
mgm_f = mgm_daily.get(d)
if mgm_f is not None:
future_forecasts.append(
- f"{d[5:]}: {t}{temp_symbol} | 🇹🇷 MGM: {mgm_f}{temp_symbol}"
+ f"{d[5:]}: {t}{temp_symbol} | 馃嚬馃嚪 MGM: {mgm_f}{temp_symbol}"
)
else:
future_forecasts.append(f"{d[5:]}: {t}{temp_symbol}")
- msg_lines.append("📅 " + " | ".join(future_forecasts))
+ msg_lines.append("馃搮 " + " | ".join(future_forecasts))
- # --- 3.5 日出日落 + 日照时长 ---
+ # --- 3.5 鏃ュ嚭鏃ヨ惤 + 鏃ョ収鏃堕暱 ---
sunrises = daily.get("sunrise", [])
sunsets = daily.get("sunset", [])
sunshine_durations = daily.get("sunshine_duration", [])
@@ -486,14 +488,14 @@ def start_bot():
if "T" in str(sunsets[0])
else sunsets[0]
)
- sun_line = f"🌅 日出 {sunrise_t} | 🌇 日落 {sunset_t}"
+ sun_line = f"馃寘 鏃ュ嚭 {sunrise_t} | 馃寚 鏃ヨ惤 {sunset_t}"
if sunshine_durations:
- sunshine_hours = sunshine_durations[0] / 3600 # 秒 -> 小时
- sun_line += f" | ☀️ 日照 {sunshine_hours:.1f}h"
+ sunshine_hours = sunshine_durations[0] / 3600 # 绉?-> 灏忔椂
+ sun_line += f" | 鈽€锔?鏃ョ収 {sunshine_hours:.1f}h"
msg_lines.append(sun_line)
- # --- 4. 核心 实测区 (合并 METAR 和 MGM) ---
- # 基础数据优先用 METAR
+ # --- 4. 鏍稿績 瀹炴祴鍖?(鍚堝苟 METAR 鍜?MGM) ---
+ # 鍩虹鏁版嵁浼樺厛鐢?METAR
cur_temp = _sf(
metar.get("current", {}).get("temp")
if metar
@@ -506,7 +508,7 @@ def start_bot():
metar.get("current", {}).get("max_temp_time") if metar else None
)
obs_t_str = "N/A"
- metar_age_min = None # METAR 数据年龄(分钟)
+ metar_age_min = None # METAR 鏁版嵁骞撮緞锛堝垎閽燂級
main_source = "METAR" if metar else "MGM"
if metar:
@@ -521,7 +523,7 @@ def start_bot():
timezone(timedelta(seconds=utc_offset))
)
obs_t_str = local_dt.strftime("%H:%M")
- # 计算数据年龄
+ # 璁$畻鏁版嵁骞撮緞
now_utc = datetime.now(timezone.utc)
metar_age_min = int((now_utc - dt).total_seconds() / 60)
elif " " in obs_t:
@@ -543,27 +545,27 @@ def start_bot():
m_time = m_time.split(" ")[1][:5]
obs_t_str = m_time
- # 数据年龄标注
+ # 鏁版嵁骞撮緞鏍囨敞
age_tag = ""
if metar_age_min is not None:
if metar_age_min >= 60:
- age_tag = f" ⚠️{metar_age_min}分钟前"
+ age_tag = f" 鈿狅笍{metar_age_min}鍒嗛挓鍓?
elif metar_age_min >= 30:
- age_tag = f" ⏳{metar_age_min}分钟前"
+ age_tag = f" 鈴硔metar_age_min}鍒嗛挓鍓?
max_str = ""
if max_p is not None:
import math
settled_val = math.floor(max_p + 0.5)
- max_str = f" (最高: {max_p}{temp_symbol}"
+ max_str = f" (鏈€楂? {max_p}{temp_symbol}"
if max_p_time:
max_str += f" @{max_p_time}"
- max_str += f" → WU {settled_val}{temp_symbol})"
+ max_str += f" 鈫?WU {settled_val}{temp_symbol})"
- # --- 天气状况总结 ---
+ # --- 澶╂皵鐘跺喌鎬荤粨 ---
wx_summary = ""
- # 优先使用 METAR 天气现象
+ # 浼樺厛浣跨敤 METAR 澶╂皵鐜拌薄
metar_wx = metar.get("current", {}).get("wx_desc", "") if metar else ""
metar_clouds = metar.get("current", {}).get("clouds", []) if metar else []
mgm_cloud = mgm.get("current", {}).get("cloud_cover") if mgm else None
@@ -586,21 +588,21 @@ def start_bot():
fog_codes = {"FG", "BR", "HZ", "FZFG"}
ts_codes = {"TS", "TSRA"}
if ts_codes & wx_tokens:
- wx_summary = "⛈️ 雷暴"
+ wx_summary = "鉀堬笍 闆锋毚"
elif {"+RA", "+SN"} & wx_tokens:
- wx_summary = "🌧️ 大雨" if "+RA" in wx_tokens else "❄️ 大雪"
+ wx_summary = "馃導锔?澶ч洦" if "+RA" in wx_tokens else "鉂勶笍 澶ч洩"
elif rain_codes & wx_tokens:
wx_summary = (
- "🌧️ 小雨" if {"-RA", "-DZ", "DZ"} & wx_tokens else "🌧️ 下雨"
+ "馃導锔?灏忛洦" if {"-RA", "-DZ", "DZ"} & wx_tokens else "馃導锔?涓嬮洦"
)
elif snow_codes & wx_tokens:
- wx_summary = "❄️ 下雪"
+ wx_summary = "鉂勶笍 涓嬮洩"
elif fog_codes & wx_tokens:
- wx_summary = "🌫️ 雾/霾"
+ wx_summary = "馃尗锔?闆?闇?
- # 如果 METAR 没有特殊现象,用云量推断
+ # 濡傛灉 METAR 娌℃湁鐗规畩鐜拌薄锛岀敤浜戦噺鎺ㄦ柇
if not wx_summary:
- # 优先 METAR 云层,回退 MGM
+ # 浼樺厛 METAR 浜戝眰锛屽洖閫€ MGM
cover_code = ""
if metar_clouds:
cover_code = metar_clouds[-1].get("cover", "")
@@ -608,98 +610,98 @@ def start_bot():
if cover_code in ("SKC", "CLR") or (
cover_code == "" and mgm_cloud is not None and mgm_cloud <= 1
):
- wx_summary = "☀️ 晴"
+ wx_summary = "鈽€锔?鏅?
elif cover_code == "FEW" or (
cover_code == "" and mgm_cloud is not None and mgm_cloud <= 2
):
- wx_summary = "🌤️ 晴间少云"
+ wx_summary = "馃尋锔?鏅撮棿灏戜簯"
elif cover_code == "SCT" or (
cover_code == "" and mgm_cloud is not None and mgm_cloud <= 4
):
- wx_summary = "⛅ 晴间多云"
+ wx_summary = "鉀?鏅撮棿澶氫簯"
elif cover_code == "BKN" or (
cover_code == "" and mgm_cloud is not None and mgm_cloud <= 6
):
- wx_summary = "🌥️ 多云"
+ wx_summary = "馃尌锔?澶氫簯"
elif cover_code == "OVC" or (
cover_code == "" and mgm_cloud is not None and mgm_cloud <= 8
):
- wx_summary = "☁️ 阴天"
+ wx_summary = "鈽侊笍 闃村ぉ"
elif mgm_cloud is not None:
cloud_names = {
- 0: "☀️ 晴",
- 1: "🌤️ 晴",
- 2: "🌤️ 少云",
- 3: "⛅ 散云",
- 4: "⛅ 散云",
- 5: "🌥️ 多云",
- 6: "🌥️ 多云",
- 7: "☁️ 阴",
- 8: "☁️ 阴天",
+ 0: "鈽€锔?鏅?,
+ 1: "馃尋锔?鏅?,
+ 2: "馃尋锔?灏戜簯",
+ 3: "鉀?鏁d簯",
+ 4: "鉀?鏁d簯",
+ 5: "馃尌锔?澶氫簯",
+ 6: "馃尌锔?澶氫簯",
+ 7: "鈽侊笍 闃?,
+ 8: "鈽侊笍 闃村ぉ",
}
wx_summary = cloud_names.get(mgm_cloud, "")
wx_display = f" {wx_summary}" if wx_summary else ""
msg_lines.append(
- f"\n✈️ 实测 ({main_source}): {cur_temp}{temp_symbol}{max_str} |{wx_display} | {obs_t_str}{age_tag}"
+ f"\n鉁堬笍 瀹炴祴 ({main_source}): {cur_temp}{temp_symbol}{max_str} |{wx_display} | {obs_t_str}{age_tag}"
)
if mgm:
m_c = mgm.get("current", {})
- # 翻译风向
+ # 缈昏瘧椋庡悜
wind_dir = m_c.get("wind_dir")
wind_speed_ms = m_c.get("wind_speed_ms")
dir_str = ""
if wind_dir is not None:
- dirs = ["北", "东北", "东", "东南", "南", "西南", "西", "西北"]
- dir_str = dirs[int((float(wind_dir) + 22.5) % 360 / 45)] + "风 "
+ dirs = ["鍖?, "涓滃寳", "涓?, "涓滃崡", "鍗?, "瑗垮崡", "瑗?, "瑗垮寳"]
+ dir_str = dirs[int((float(wind_dir) + 22.5) % 360 / 45)] + "椋?"
- # 体感和湿度(跳过缺失数据)
+ # 浣撴劅鍜屾箍搴︼紙璺宠繃缂哄け鏁版嵁锛?
feels_like = m_c.get("feels_like")
humidity = m_c.get("humidity")
if feels_like is not None or humidity is not None:
parts = []
if feels_like is not None:
- parts.append(f"🌡️ 体感: {feels_like}°C")
+ parts.append(f"馃尅锔?浣撴劅: {feels_like}掳C")
- # 针对安卡拉,补充市区(Center)实测值
- ankara_center = next((s for s in weather_data.get("mgm_nearby", []) if "Bölge/Center" in s.get("name", "")), None)
+ # 閽堝瀹夊崱鎷夛紝琛ュ厖甯傚尯(Center)瀹炴祴鍊?
+ ankara_center = next((s for s in weather_data.get("mgm_nearby", []) if "B枚lge/Center" in s.get("name", "")), None)
if ankara_center:
- parts.append(f"Ankara (Bölge/Center): {ankara_center['temp']}°C")
+ parts.append(f"Ankara (B枚lge/Center): {ankara_center['temp']}掳C")
if humidity is not None:
- parts.append(f"💧 {humidity}%")
+ parts.append(f"馃挧 {humidity}%")
msg_lines.append(f" [MGM] {' | '.join(parts)}")
- # 风况(跳过缺失数据)
+ # 椋庡喌锛堣烦杩囩己澶辨暟鎹級
if wind_dir is not None and wind_speed_ms is not None:
msg_lines.append(
- f" [MGM] 🌬️ {dir_str}{wind_dir}° ({wind_speed_ms} m/s) | 💧 降水: {m_c.get('rain_24h') or 0}mm"
+ f" [MGM] 馃尙锔?{dir_str}{wind_dir}掳 ({wind_speed_ms} m/s) | 馃挧 闄嶆按: {m_c.get('rain_24h') or 0}mm"
)
- # 新增:气压和云量
+ # 鏂板锛氭皵鍘嬪拰浜戦噺
extra_parts = []
pressure = m_c.get("pressure")
if pressure is not None:
- extra_parts.append(f"🌡 气压: {pressure}hPa")
+ extra_parts.append(f"馃尅 姘斿帇: {pressure}hPa")
cloud_cover = m_c.get("cloud_cover")
if cloud_cover is not None:
cloud_desc_map = {
- 0: "晴朗",
- 1: "少云",
- 2: "少云",
- 3: "散云",
- 4: "散云",
- 5: "多云",
- 6: "多云",
- 7: "很多云",
- 8: "阴天",
+ 0: "鏅存湕",
+ 1: "灏戜簯",
+ 2: "灏戜簯",
+ 3: "鏁d簯",
+ 4: "鏁d簯",
+ 5: "澶氫簯",
+ 6: "澶氫簯",
+ 7: "寰堝浜?,
+ 8: "闃村ぉ",
}
cloud_text = cloud_desc_map.get(cloud_cover, f"{cloud_cover}/8")
- extra_parts.append(f"☁️ 云量: {cloud_text}({cloud_cover}/8)")
+ extra_parts.append(f"鈽侊笍 浜戦噺: {cloud_text}({cloud_cover}/8)")
mgm_max = m_c.get("mgm_max_temp")
if mgm_max is not None:
- extra_parts.append(f"🌡️ MGM最高: {mgm_max}°C")
+ extra_parts.append(f"馃尅锔?MGM鏈€楂? {mgm_max}掳C")
if extra_parts:
msg_lines.append(f" [MGM] {' | '.join(extra_parts)}")
@@ -713,45 +715,45 @@ def start_bot():
cloud_desc = ""
if clouds:
c_map = {
- "BKN": "多云",
- "OVC": "阴天",
- "FEW": "少云",
- "SCT": "散云",
- "SKC": "晴",
- "CLR": "晴",
+ "BKN": "澶氫簯",
+ "OVC": "闃村ぉ",
+ "FEW": "灏戜簯",
+ "SCT": "鏁d簯",
+ "SKC": "鏅?,
+ "CLR": "鏅?,
}
main = clouds[-1]
- cloud_desc = f"☁️ {c_map.get(main.get('cover'), main.get('cover'))}"
+ cloud_desc = f"鈽侊笍 {c_map.get(main.get('cover'), main.get('cover'))}"
prefix = "[METAR]" if mgm else " "
if not mgm:
msg_lines.append(
- f" {prefix} 💨 {wind or 0}kt ({wind_dir or 0}°) | 👁️ {vis or 10}mi"
+ f" {prefix} 馃挩 {wind or 0}kt ({wind_dir or 0}掳) | 馃憗锔?{vis or 10}mi"
)
if cloud_desc:
msg_lines.append(
- f" {prefix} {cloud_desc} | 👁️ {vis or 10}mi | 💨 {wind or 0}kt"
+ f" {prefix} {cloud_desc} | 馃憗锔?{vis or 10}mi | 馃挩 {wind or 0}kt"
)
- # --- 5. 态势特征提取 ---
+ # --- 5. 鎬佸娍鐗瑰緛鎻愬彇 ---
feature_str, ai_context = analyze_weather_trend(
weather_data, temp_symbol, city_name
)
if feature_str:
- # 仅将最核心的信息展示给用户作为"态势分析"
- # 但后面会把更全的数据传给 AI
- msg_lines.append("\n💡 分析:")
+ # 浠呭皢鏈€鏍稿績鐨勪俊鎭睍绀虹粰鐢ㄦ埛浣滀负"鎬佸娍鍒嗘瀽"
+ # 浣嗗悗闈細鎶婃洿鍏ㄧ殑鏁版嵁浼犵粰 AI
+ msg_lines.append("\n馃挕 鍒嗘瀽:")
for line in feature_str.split("\n"):
if line.strip():
msg_lines.append(f"- {line.strip()}")
- # --- 6. Groq AI 深度分析 ---
+ # --- 6. Groq AI 娣卞害鍒嗘瀽 ---
try:
from src.analysis.ai_analyzer import get_ai_analysis
- # 构建更全的背景数据给 AI
+ # 鏋勫缓鏇村叏鐨勮儗鏅暟鎹粰 AI
- # 补充多模型分歧
+ # 琛ュ厖澶氭ā鍨嬪垎姝?
mm = weather_data.get("multi_model", {})
if mm.get("forecasts"):
mm_str = " | ".join(
@@ -761,26 +763,26 @@ def start_bot():
if v
]
)
- ai_context += f"\n模型分歧: {mm_str}"
+ ai_context += f"\n妯″瀷鍒嗘: {mm_str}"
ai_result = get_ai_analysis(ai_context, city_name, temp_symbol)
if ai_result:
msg_lines.append(f"\n{ai_result}")
except Exception as e:
- logger.error(f"调用 Groq AI 分析失败: {e}")
+ logger.error(f"璋冪敤 Groq AI 鍒嗘瀽澶辫触: {e}")
- msg_lines.append(f"\n💳 本次消耗 {CITY_QUERY_COST} 积分。")
+ msg_lines.append(f"\n馃挸 鏈娑堣€?{CITY_QUERY_COST} 绉垎銆?)
bot.send_message(message.chat.id, "\n".join(msg_lines), parse_mode="HTML")
except Exception as e:
import traceback
- logger.error(f"查询失败: {e}\n{traceback.format_exc()}")
- bot.reply_to(message, f"❌ 查询失败: {e}")
+ logger.error(f"鏌ヨ澶辫触: {e}\n{traceback.format_exc()}")
+ bot.reply_to(message, f"鉂?鏌ヨ澶辫触: {e}")
@bot.message_handler(func=lambda message: True, content_types=['text'])
def track_activity(message):
- """全量监听消息,用于记录群内发言积分(非指令消息)"""
+ """鍏ㄩ噺鐩戝惉娑堟伅锛岀敤浜庤褰曠兢鍐呭彂瑷€绉垎(闈炴寚浠ゆ秷鎭?"""
if message.text.startswith('/'):
return
if message.chat.type not in ("group", "supergroup"):
@@ -804,9 +806,10 @@ def start_bot():
f"daily_points={result.get('daily_points')}/{MESSAGE_DAILY_CAP}"
)
- logger.info("🤖 Bot 启动中...")
+ logger.info("馃 Bot 鍚姩涓?..")
bot.infinity_polling()
if __name__ == "__main__":
start_bot()
+
diff --git a/src/data_collection/weather_sources.py b/src/data_collection/weather_sources.py
index 0f684047..74969853 100644
--- a/src/data_collection/weather_sources.py
+++ b/src/data_collection/weather_sources.py
@@ -1,7 +1,9 @@
import csv
+import os
import requests
import re
import time
+import threading
from typing import Optional, Dict, List
from datetime import datetime, timedelta
from loguru import logger
@@ -42,6 +44,7 @@ class WeatherDataCollector:
# Meteoblue 仅在增益最大的城市启用(减少配额消耗与冗余请求)
METEOBLUE_PRIORITY_CITIES = {
+ "ankara",
"london",
"paris",
"seoul",
@@ -61,6 +64,11 @@ class WeatherDataCollector:
self.timeout = 30 # 增加超时以支持高延迟 VPS
self.session = requests.Session()
+ self.meteoblue_cache_ttl_sec = int(
+ os.getenv("METEOBLUE_CACHE_TTL_SEC", "1800")
+ )
+ self._meteoblue_cache: Dict[str, Dict] = {}
+ self._meteoblue_cache_lock = threading.Lock()
# 设置代理
proxy = config.get("proxy")
@@ -1194,11 +1202,24 @@ class WeatherDataCollector:
) -> Optional[Dict]:
"""
通过 Meteoblue 官方 API 获取高精度预测数据
+ 带本地缓存,避免频繁请求触发 429。
"""
if not self.meteoblue_key:
logger.warning("Meteoblue API Key 未配置,跳过抓取。")
return None
+ cache_key = f"{round(float(lat), 4)}:{round(float(lon), 4)}:{'f' if use_fahrenheit else 'c'}"
+ now_ts = time.time()
+ with self._meteoblue_cache_lock:
+ cached = self._meteoblue_cache.get(cache_key)
+ if (
+ cached
+ and now_ts - float(cached.get("t", 0)) < self.meteoblue_cache_ttl_sec
+ ):
+ cached_data = cached.get("data")
+ if isinstance(cached_data, dict):
+ return dict(cached_data)
+
try:
# 1. 调用官方 API (使用 basic-day 包,它是多模型 ML 融合结果)
# 格式: https://my.meteoblue.com/packages/basic-day?apikey=KEY&lat=LAT&lon=LON&format=json
@@ -1246,12 +1267,29 @@ class WeatherDataCollector:
else:
result["daily_highs"] = max_temps
+ with self._meteoblue_cache_lock:
+ self._meteoblue_cache[cache_key] = {
+ "t": now_ts,
+ "data": dict(result),
+ }
+
logger.info(
f"✅ Meteoblue API 获取成功 ({lat},{lon}): 今天 {result['today_high']}{result['unit']}"
)
return result
except Exception as e:
- logger.error(f"Meteoblue API fetch failed: {e}")
+ status_code = getattr(getattr(e, "response", None), "status_code", None)
+ if status_code == 429:
+ logger.warning("Meteoblue API 限流(429),尝试使用本地缓存回退。")
+ else:
+ logger.error(f"Meteoblue API fetch failed: {e}")
+
+ with self._meteoblue_cache_lock:
+ stale = self._meteoblue_cache.get(cache_key)
+ if stale and isinstance(stale.get("data"), dict):
+ fallback = dict(stale["data"])
+ fallback["stale_cache"] = True
+ return fallback
return None
def extract_date_from_title(self, title: str) -> Optional[str]:
diff --git a/src/onchain/__init__.py b/src/onchain/__init__.py
new file mode 100644
index 00000000..64a8126f
--- /dev/null
+++ b/src/onchain/__init__.py
@@ -0,0 +1 @@
+"""On-chain monitoring modules."""
diff --git a/src/onchain/polygon_wallet_watcher.py b/src/onchain/polygon_wallet_watcher.py
new file mode 100644
index 00000000..b43b59a8
--- /dev/null
+++ b/src/onchain/polygon_wallet_watcher.py
@@ -0,0 +1,354 @@
+import json
+import os
+import threading
+import time
+from datetime import datetime, timezone
+from decimal import Decimal, InvalidOperation
+from typing import Any, Dict, List, Optional, Set, Tuple
+
+from loguru import logger
+from web3 import Web3
+
+TRANSFER_TOPIC = "0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef"
+ERC20_ABI = [
+ {
+ "constant": True,
+ "inputs": [],
+ "name": "symbol",
+ "outputs": [{"name": "", "type": "string"}],
+ "type": "function",
+ },
+ {
+ "constant": True,
+ "inputs": [],
+ "name": "decimals",
+ "outputs": [{"name": "", "type": "uint8"}],
+ "type": "function",
+ },
+]
+
+
+def _env_bool(name: str, default: bool) -> bool:
+ raw = os.getenv(name)
+ if raw is None:
+ return default
+ return raw.strip().lower() in {"1", "true", "yes", "on"}
+
+
+def _env_int(name: str, default: int) -> int:
+ raw = os.getenv(name)
+ if raw is None:
+ return default
+ try:
+ return int(raw)
+ except Exception:
+ return default
+
+
+def _short(addr: str, left: int = 6, right: int = 4) -> str:
+ if not addr:
+ return "unknown"
+ if len(addr) <= left + right + 2:
+ return addr
+ return f"{addr[:left + 2]}...{addr[-right:]}"
+
+
+def _state_file() -> str:
+ root = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
+ return os.path.join(root, "data", "polygon_wallet_watch_state.json")
+
+
+def _load_state(path: str) -> Dict[str, Any]:
+ if not os.path.exists(path):
+ return {"last_scanned_block": 0, "seen_tx": {}}
+ try:
+ with open(path, "r", encoding="utf-8") as fh:
+ data = json.load(fh)
+ if isinstance(data, dict):
+ data.setdefault("last_scanned_block", 0)
+ data.setdefault("seen_tx", {})
+ return data
+ except Exception as exc:
+ logger.warning(f"failed to load polygon watch state: {exc}")
+ return {"last_scanned_block": 0, "seen_tx": {}}
+
+
+def _save_state(path: str, state: Dict[str, Any]) -> None:
+ os.makedirs(os.path.dirname(path), exist_ok=True)
+ tmp_path = f"{path}.tmp"
+ with open(tmp_path, "w", encoding="utf-8") as fh:
+ json.dump(state, fh, ensure_ascii=False, indent=2)
+ os.replace(tmp_path, path)
+
+
+def _cleanup_seen_tx(state: Dict[str, Any], now_ts: int, keep_sec: int) -> None:
+ seen = state.get("seen_tx", {})
+ if not isinstance(seen, dict):
+ state["seen_tx"] = {}
+ return
+ stale = [key for key, value in seen.items() if now_ts - int(value or 0) > keep_sec]
+ for tx_hash in stale:
+ seen.pop(tx_hash, None)
+
+
+def _parse_addresses(raw: Optional[str]) -> Set[str]:
+ out: Set[str] = set()
+ if not raw:
+ return out
+ for part in raw.split(","):
+ addr = part.strip().lower()
+ if addr and addr.startswith("0x") and len(addr) == 42:
+ out.add(addr)
+ return out
+
+
+def _polygon_scan_tx_url(tx_hash: str) -> str:
+ base = os.getenv("POLYGON_WALLET_WATCH_TX_BASE") or "https://polygonscan.com/tx/"
+ return f"{base.rstrip('/')}/{tx_hash}"
+
+
+def _polygon_scan_addr_url(address: str) -> str:
+ base = os.getenv("POLYGON_WALLET_WATCH_ADDR_BASE") or "https://polygonscan.com/address/"
+ return f"{base.rstrip('/')}/{address}"
+
+
+def _format_matic(wei_value: int) -> str:
+ try:
+ matic = Decimal(wei_value) / Decimal(10**18)
+ except (InvalidOperation, ValueError):
+ return "0"
+
+ if matic == matic.to_integral_value():
+ return f"{int(matic)}"
+ return f"{matic.normalize():f}".rstrip("0").rstrip(".")
+
+
+def _safe_lower(value: Any) -> str:
+ if value is None:
+ return ""
+ return str(value).lower()
+
+
+def _extract_transfer_events(
+ w3: Web3,
+ receipt: Any,
+ watch_set: Set[str],
+ token_meta_cache: Dict[str, Tuple[str, int]],
+) -> List[str]:
+ lines: List[str] = []
+ for log in receipt.logs or []:
+ try:
+ if not log.topics or len(log.topics) < 3:
+ continue
+ topic0 = log.topics[0].hex().lower()
+ if topic0 != TRANSFER_TOPIC:
+ continue
+
+ from_addr = "0x" + log.topics[1].hex()[-40:]
+ to_addr = "0x" + log.topics[2].hex()[-40:]
+ from_addr = from_addr.lower()
+ to_addr = to_addr.lower()
+ if from_addr not in watch_set and to_addr not in watch_set:
+ continue
+
+ token_addr = _safe_lower(log.address)
+ symbol, decimals = token_meta_cache.get(token_addr, ("ERC20", 18))
+ if token_addr not in token_meta_cache:
+ try:
+ token = w3.eth.contract(address=Web3.to_checksum_address(token_addr), abi=ERC20_ABI)
+ symbol_raw = token.functions.symbol().call()
+ decimals_raw = token.functions.decimals().call()
+ symbol = str(symbol_raw or "ERC20")
+ decimals = int(decimals_raw)
+ except Exception:
+ symbol = "ERC20"
+ decimals = 18
+ token_meta_cache[token_addr] = (symbol, decimals)
+
+ amount_int = int(log.data.hex(), 16) if log.data else 0
+ amount = Decimal(amount_int) / (Decimal(10) ** Decimal(max(decimals, 0)))
+ amount_txt = f"{amount.normalize():f}".rstrip("0").rstrip(".") if amount else "0"
+
+ direction = "IN" if to_addr in watch_set and from_addr not in watch_set else "OUT"
+ if from_addr in watch_set and to_addr in watch_set:
+ direction = "SELF"
+
+ lines.append(
+ f"- {direction} {symbol}: {amount_txt} ({_short(from_addr)} -> {_short(to_addr)})"
+ )
+ except Exception:
+ continue
+ return lines
+
+
+def _build_message(
+ tx: Any,
+ block_ts: int,
+ matched_wallet: str,
+ transfer_lines: List[str],
+) -> str:
+ tx_hash = tx["hash"].hex()
+ from_addr = _safe_lower(tx.get("from"))
+ to_addr = _safe_lower(tx.get("to"))
+
+ if from_addr == matched_wallet and to_addr == matched_wallet:
+ direction = "SELF"
+ elif to_addr == matched_wallet:
+ direction = "IN"
+ elif from_addr == matched_wallet:
+ direction = "OUT"
+ else:
+ direction = "RELATED"
+
+ matic_value = int(tx.get("value", 0) or 0)
+ block_time = datetime.fromtimestamp(block_ts, tz=timezone.utc).strftime("%Y-%m-%d %H:%M:%S UTC")
+
+ lines = [
+ "⛓ Polygon 钱包异动",
+ f"钱包: {_short(matched_wallet)}",
+ f"方向: {direction}",
+ f"MATIC: {_format_matic(matic_value)}",
+ f"区块: {tx.get('blockNumber')}",
+ f"时间: {block_time}",
+ f"Tx: {tx_hash}",
+ ]
+
+ if transfer_lines:
+ lines.append("Token 转账:")
+ lines.extend(transfer_lines[:6])
+
+ lines.append(f"交易链接: {_polygon_scan_tx_url(tx_hash)}")
+ lines.append(f"钱包链接: {_polygon_scan_addr_url(matched_wallet)}")
+ return "\n".join(lines)
+
+
+def start_polygon_wallet_watch_loop(bot: Any) -> Optional[threading.Thread]:
+ enabled = _env_bool("POLYGON_WALLET_WATCH_ENABLED", False)
+ chat_id = os.getenv("TELEGRAM_CHAT_ID")
+ rpc_url = os.getenv("POLYGON_RPC_URL")
+ watch_set = _parse_addresses(os.getenv("POLYGON_WALLET_WATCH_ADDRESSES"))
+
+ if not enabled:
+ logger.info("polygon wallet watcher disabled")
+ return None
+ if not chat_id:
+ logger.warning("polygon wallet watcher skipped: TELEGRAM_CHAT_ID is not set")
+ return None
+ if not rpc_url:
+ logger.warning("polygon wallet watcher skipped: POLYGON_RPC_URL is not set")
+ return None
+ if not watch_set:
+ logger.warning("polygon wallet watcher skipped: POLYGON_WALLET_WATCH_ADDRESSES is empty")
+ return None
+
+ poll_sec = max(3, _env_int("POLYGON_WALLET_WATCH_INTERVAL_SEC", 8))
+ confirmations = max(0, _env_int("POLYGON_WALLET_WATCH_CONFIRMATIONS", 2))
+ max_blocks_per_cycle = max(1, _env_int("POLYGON_WALLET_WATCH_MAX_BLOCKS_PER_CYCLE", 30))
+ seen_ttl_sec = max(3600, _env_int("POLYGON_WALLET_WATCH_SEEN_TTL_SEC", 7 * 86400))
+ state_path = _state_file()
+
+ provider_timeout = max(5, _env_int("POLYGON_WALLET_WATCH_RPC_TIMEOUT_SEC", 10))
+ w3 = Web3(Web3.HTTPProvider(rpc_url, request_kwargs={"timeout": provider_timeout}))
+
+ def _runner() -> None:
+ token_meta_cache: Dict[str, Tuple[str, int]] = {}
+ state = _load_state(state_path)
+
+ if not w3.is_connected():
+ logger.error("polygon wallet watcher failed: cannot connect to POLYGON_RPC_URL")
+ return
+
+ try:
+ chain_id = int(w3.eth.chain_id)
+ if chain_id != 137:
+ logger.warning(f"polygon wallet watcher connected to unexpected chain_id={chain_id}")
+ except Exception as exc:
+ logger.warning(f"polygon wallet watcher cannot read chain id: {exc}")
+
+ latest_block = int(w3.eth.block_number)
+ if int(state.get("last_scanned_block") or 0) <= 0:
+ state["last_scanned_block"] = max(0, latest_block - confirmations)
+ _save_state(state_path, state)
+
+ logger.info(
+ f"polygon wallet watcher started wallets={len(watch_set)} "
+ f"poll={poll_sec}s confirmations={confirmations} state_path={state_path}"
+ )
+
+ while True:
+ cycle_ts = int(time.time())
+ try:
+ _cleanup_seen_tx(state, cycle_ts, seen_ttl_sec)
+
+ latest = int(w3.eth.block_number)
+ safe_latest = latest - confirmations
+ last_scanned = int(state.get("last_scanned_block") or 0)
+
+ if safe_latest <= last_scanned:
+ time.sleep(poll_sec)
+ continue
+
+ from_block = last_scanned + 1
+ to_block = min(safe_latest, from_block + max_blocks_per_cycle - 1)
+
+ for block_num in range(from_block, to_block + 1):
+ block = w3.eth.get_block(block_num, full_transactions=True)
+ block_ts = int(block.get("timestamp") or cycle_ts)
+
+ for tx in block.transactions or []:
+ tx_hash = tx["hash"].hex().lower()
+ from_addr = _safe_lower(tx.get("from"))
+ to_addr = _safe_lower(tx.get("to"))
+
+ matched_wallet = ""
+ if from_addr in watch_set:
+ matched_wallet = from_addr
+ elif to_addr in watch_set:
+ matched_wallet = to_addr
+
+ if not matched_wallet:
+ continue
+
+ if tx_hash in state.get("seen_tx", {}):
+ continue
+
+ transfer_lines: List[str] = []
+ try:
+ receipt = w3.eth.get_transaction_receipt(tx["hash"])
+ transfer_lines = _extract_transfer_events(
+ w3=w3,
+ receipt=receipt,
+ watch_set=watch_set,
+ token_meta_cache=token_meta_cache,
+ )
+ except Exception:
+ transfer_lines = []
+
+ message = _build_message(
+ tx=tx,
+ block_ts=block_ts,
+ matched_wallet=matched_wallet,
+ transfer_lines=transfer_lines,
+ )
+ bot.send_message(chat_id, message, disable_web_page_preview=True)
+ state.setdefault("seen_tx", {})[tx_hash] = cycle_ts
+ logger.info(
+ f"polygon wallet alert pushed wallet={matched_wallet} "
+ f"tx={tx_hash} block={block_num}"
+ )
+
+ state["last_scanned_block"] = block_num
+ _save_state(state_path, state)
+
+ time.sleep(1)
+ except Exception:
+ logger.exception("polygon wallet watcher cycle failed")
+ time.sleep(poll_sec)
+
+ thread = threading.Thread(
+ target=_runner,
+ name="polygon-wallet-watcher",
+ daemon=True,
+ )
+ thread.start()
+ return thread