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