From 40dc5062dd4c9bc05a11781ac4c546955d72124a Mon Sep 17 00:00:00 2001 From: "2569718930@qq.com" <2569718930@qq.com> Date: Sat, 7 Feb 2026 00:46:15 +0800 Subject: [PATCH] feat: Implement risk management, order book analysis, and a new alerting system with enhanced Polymarket batch price fetching and market data processing. --- bot_listener.py | 109 +---- main.py | 577 ++++++++++++++++++++----- src/analysis/orderbook_analyzer.py | 102 +++-- src/data_collection/polymarket_api.py | 137 ++++-- src/data_collection/weather_sources.py | 10 +- src/models/statistical_model.py | 4 +- src/strategy/risk_manager.py | 206 ++++----- src/utils/notifier.py | 44 +- test_template.py | 65 +++ 9 files changed, 794 insertions(+), 460 deletions(-) create mode 100644 test_template.py diff --git a/bot_listener.py b/bot_listener.py index 2cbf4e3c..bf8eeabf 100644 --- a/bot_listener.py +++ b/bot_listener.py @@ -66,7 +66,7 @@ def start_bot(): # 过滤掉已结束的市场(价格接近0或100)和无日期的 active_signals = [] - for s in signals.values(): + for s in signals: price = s.get("price", 50) if 5 <= price <= 95 and s.get("target_date"): active_signals.append(s) @@ -370,113 +370,6 @@ def start_bot(): return html_path - @bot.message_handler(func=lambda m: True) - def handle_city_query(message): - """输入城市名直查当日天气市场""" - import re - from datetime import datetime - - query = message.text.strip() - if len(query) < 2 or query.startswith("/"): - return - - bot.send_chat_action(message.chat.id, "typing") - - try: - # 1. 优先从本地全量市场缓存读取 (速度快,不依赖实时全量扫描) - cache_path = "data/all_markets.json" - if not os.path.exists(cache_path): - # 扫码还没完成的情形 - bot.reply_to(message, "⏳ 系统正在进行首次数据同步(约需1分钟),请稍后再试。") - return - - with open(cache_path, "r", encoding="utf-8") as f: - cached_data = json.load(f) - - pm = PolymarketClient(config["polymarket"]) - - # 2. 筛选匹配城市及日期的市场 - today_str = datetime.now().strftime("%Y-%m-%d") - city_markets = [] - - for m_id, m in cached_data.items(): - title = m.get("event_title", "") + m.get("question", "") + m.get("full_title", "") - if query.lower() in title.lower(): - # 提取并验证日期 - target_date = m.get("target_date") - if not target_date: - date_match = re.search(r'(\d{4}-\d{2}-\d{2})', title) - target_date = date_match.group(1) if date_match else "Unknown" - - if target_date != "Unknown" and target_date < today_str: - continue - - m["target_date"] = target_date - city_markets.append(m) - - if not city_markets: - if message.chat.type == "private": - bot.reply_to(message, f"❌ 未找到相关的活跃天气市场。\n提示:请确保输入的是城市常用名(如 Seattle, London)。") - return - - # 获取最早日期 - valid_dates = [m["target_date"] for m in city_markets if m["target_date"] != "Unknown"] - if not valid_dates: - bot.reply_to(message, "❌ 该城市目前没有已标明结算日期的活跃市场。") - return - - earliest_date = min(valid_dates) - target_markets = [m for m in city_markets if m["target_date"] == earliest_date] - - # 3. 构建报告 - msg_lines = [ - f"🌡️ {query.upper()} 概率报告 ({earliest_date})\n", - "隐含概率 (Midpoint) 及买入报价:\n" - ] - - # 批量获取实时价格 (确保报价最新) - price_reqs = [] - for m in target_markets: - t_ids = m.get("tokens", []) - if len(t_ids) >= 1: - price_reqs.append({"token_id": t_ids[0], "side": "ask"}) - price_reqs.append({"token_id": t_ids[0], "side": "bid"}) - - price_map = pm.get_multiple_prices(price_reqs) - - for m in target_markets: - tid = m.get("active_token_id") or (m.get("tokens", [])[0] if m.get("tokens") else None) - if not tid: continue - - # 获取中点价 (概率) - mid = pm.get_midpoint(tid) - prob = f"{mid*100:.1f}%" if mid is not None else "N/A" - - # 获取报价 - buy_yes = price_map.get(f"{tid}:ask") - bid_yes = price_map.get(f"{tid}:bid") - buy_no = (1.0 - bid_yes) if bid_yes is not None else None - - yes_str = f"{int(buy_yes*100)}¢" if buy_yes else "??¢" - no_str = f"{int(buy_no*100)}¢" if buy_no else "??¢" - - opt = m.get("option") or m.get("question") or "" - # 简化选项显示 - opt = re.sub(r'.*temperature in.*be ', '', opt, flags=re.I) - - msg_lines.append( - f"🔹 {opt}\n" - f" └ 隐含概率: {prob}\n" - f" └ 买入 是:{yes_str} | 买入 否:{no_str}\n" - ) - - msg_lines.append(f"\n🔗 在 Polymarket 查看") - bot.send_message(message.chat.id, "\n".join(msg_lines), parse_mode="HTML", disable_web_page_preview=True) - - except Exception as e: - logger.error(f"城市直查失败: {e}") - if message.chat.type == "private": - bot.reply_to(message, "❌ 抱歉,数据处理出现异常。") @bot.message_handler(commands=["status"]) def get_status(message): diff --git a/main.py b/main.py index f2425665..c89ed59e 100644 --- a/main.py +++ b/main.py @@ -37,6 +37,8 @@ def main(): # 3. 初始化分析与交易组件 predictor = TemperaturePredictor() + risk_manager = RiskManager(config_data.get("config", {})) + orderbook_analyzer = OrderbookAnalyzer(config_data.get("config", {})) decision_engine = DecisionEngine(config_data.get("config", {})) whale_tracker = WhaleTracker(config_data.get("config", {}), onchain) paper_trader = PaperTrader() @@ -130,14 +132,16 @@ def main(): for m in all_weather_markets: ts = m.get("tokens", []) if isinstance(ts, str): - try: ts = json.loads(ts) - except: ts = [] - + try: + ts = json.loads(ts) + except: + ts = [] + active_tid = m.get("active_token_id") - + # 如果是多选一市场(比如 Dallas 76-77°F) if len(ts) > 2 and active_tid: - # 获取该档位的买入价 (Ask) + # 获取该档位的买入价 (Ask) price_requests.append({"token_id": active_tid, "side": "ask"}) # 获取该档位的买入“否”价所需的 Bid 价 price_requests.append({"token_id": active_tid, "side": "bid"}) @@ -164,11 +168,13 @@ def main(): # 注入实时批量价格 ts = m.get("tokens", []) if isinstance(ts, str): - try: ts = json.loads(ts) - except: ts = [] - + try: + ts = json.loads(ts) + except: + ts = [] + active_tid = m.get("active_token_id") - + # 多选一市场逻辑 if len(ts) > 2 and active_tid: m["buy_yes_live"] = token_price_map.get(f"{active_tid}:ask") @@ -229,13 +235,21 @@ def main(): if not consensus.get("consensus"): continue + temp_unit = weather_data.get("open-meteo", {}).get( + "unit", "celsius" + ) + temp_symbol = "°F" if temp_unit == "fahrenheit" else "°C" logger.info( - f"☁️ {city} 当前气温: {consensus['average_temp']}°C | 监控合约: {len(city_markets)}" + f"☁️ {city} 当前气温: {consensus['average_temp']}{temp_symbol} | 监控合约: {len(city_markets)}" ) # --- 本城市汇总预警缓存 --- city_alerts = [] city_local_time = None + city_total_vol = 0 + city_pred_high = None + city_target_date = None + city_strategy_tips = [] # B. 遍历该城市所有合约 for market in city_markets: @@ -243,6 +257,17 @@ def main(): question = market.get("question", "未知市场") event_title = market.get("event_title", "") + # 累计城市总成交量 + vol_raw = market.get("volume", 0) + if isinstance(vol_raw, str): + try: + vol_raw = float( + vol_raw.replace("$", "").replace(",", "") + ) + except: + vol_raw = 0 + city_total_vol += vol_raw + # 识别该合约的目标日期 target_date = weather.extract_date_from_title( event_title @@ -261,60 +286,168 @@ def main(): break # --- 价格获取逻辑 (增强版) --- - buy_yes_price = market.get("buy_yes_live") - buy_no_price = market.get("buy_no_live") - ts = market.get("tokens", []) - if isinstance(ts, str): ts = json.loads(ts) + # 使用 token_price_map 获取实时数据 active_tid = market.get("active_token_id") + ts = market.get("tokens", []) + if isinstance(ts, str): + ts = json.loads(ts) - # 特殊逻辑:针对多选一型市场(Dallas/Chicago 等温度段) - if len(ts) > 2 and active_tid: - # 1. 尝试从批量映射中重新获取该档位的精确买入价(Ask) - buy_yes_price = polymarket.get_price(active_tid, side="ask") if buy_yes_price is None else buy_yes_price - # 2. 计算“否”的价格(在多选一里是 1 - 最高买入意愿价Bid) - bid_price = polymarket.get_price(active_tid, side="bid") - if bid_price: - buy_no_price = 1.0 - bid_price - - if i < 3: # 仅对前几个档位打印调试 - logger.debug(f"Categorical价格匹配 [{market.get('city')}-{market.get('question')}]: Yes={buy_yes_price}, No={buy_no_price} (from Bid={bid_price})") + buy_yes_price = None + buy_no_price = None + bid_yes_price = None - # 通用回退逻辑 (二选一) - if buy_yes_price is None or buy_no_price is None: - if ts and len(ts) >= 2: - prices = polymarket.get_buy_prices(ts[0], ts[1]) - if prices: - buy_yes_price = prices.get("buy_yes") - buy_no_price = prices.get("buy_no") - - # 最后的回退:使用原始 Gamma 概率 (利用 outcome_index 获取正确的那一个) - if buy_yes_price is None or buy_no_price is None: - gamma_prices = market.get("prices", []) - if isinstance(gamma_prices, str): gamma_prices = json.loads(gamma_prices) - - # 尝试获取该档位相对应的概率索引 - idx = market.get("outcome_index", 0) - prob = float(gamma_prices[idx]) if (gamma_prices and idx < len(gamma_prices)) else 0.5 - - if buy_yes_price is None: buy_yes_price = prob - if buy_no_price is None: buy_no_price = 1.0 - prob + if len(ts) == 2: + # 传统二选一市场 (Yes/No Token 独立) + buy_yes_price = token_price_map.get(f"{ts[0]}:ask") + buy_no_price = token_price_map.get(f"{ts[1]}:ask") + bid_yes_price = token_price_map.get(f"{ts[0]}:bid") + elif active_tid: + # 多选一市场 (单 Token 对应一个档位) + buy_yes_price = token_price_map.get(f"{active_tid}:ask") + bid_yes_price = token_price_map.get(f"{active_tid}:bid") + if bid_yes_price is not None: + buy_no_price = 1.0 - bid_yes_price + + # 兜底概率计算 + current_prob = ( + (buy_yes_price + bid_yes_price) / 2 + if (buy_yes_price and bid_yes_price) + else (buy_yes_price or 0.5) + ) + if buy_no_price is None: + buy_no_price = 1.0 - current_prob + + # 计算价格趋势 + prev_data = price_history.get(market_id, {}) + prev_prob = prev_data.get("price", current_prob) + prob_change = (current_prob - prev_prob) * 100 + trend_str = ( + f"▲{abs(prob_change):.0f}%" + if prob_change > 0.5 + else ( + f"▼{abs(prob_change):.0f}%" + if prob_change < -0.5 + else "" + ) + ) + + # 更新历史缓存 + price_history[market_id] = { + "price": current_prob, + "timestamp": datetime.now().isoformat(), + } + + # --- 预警收集 (自动推送逻辑) --- + # 触发阈值: 价格处于 85-95 锁死区间,或者概率异动 > 10% + is_price_locked = ( + current_prob >= 0.85 or (1 - current_prob) >= 0.85 + ) + is_big_move = abs(prob_change) >= 10 + + if is_price_locked or is_big_move: + alert_key = f"alert_{market_id}_{int(current_prob * 100)}" + if alert_key not in pushed_signals: + # 深度分析订单簿 + ob_data = ( + polymarket.get_orderbook(active_tid) + if active_tid + else None + ) + ob_analysis = ( + orderbook_analyzer.analyze(ob_data) + if ob_data + else { + "tradeable": False, + "liquidity": "枯竭", + "spread": 0, + "mid_price": current_prob, + } + ) + + # 预测偏差分析 + if ref_temp: + city_pred_high = ref_temp # 记录到城市概览 + temp_match = re.search( + r"(\d+)(?:-(\d+))?°[FC]", question + ) + if temp_match: + low_b = int(temp_match.group(1)) + high_b = ( + int(temp_match.group(2)) + if temp_match.group(2) + else low_b + ) + diff = ref_temp - ((low_b + high_b) / 2) + msg += f"\n📐 预测偏差: {diff:+.1f}{temp_symbol} (预测 {ref_temp}{temp_symbol})" + + # 生成策略建议:仅保留模型一致提示 + if abs(diff) < 2 and current_prob > 0.7: + city_strategy_tips.append( + f"预测温度{ref_temp}{temp_symbol}落在{question}区间,市场与模型一致" + ) + + # 模拟下单 - 使用 Ask 价格(实际可成交价格) + if buy_yes_price and buy_yes_price > 0.5: + trigger_side = "Buy Yes" + trigger_price = int(buy_yes_price * 100) + else: + trigger_side = "Buy No" + trigger_price = ( + int(buy_no_price * 100) + if buy_no_price + else int((1 - current_prob) * 100) + ) + + success = paper_trader.open_position( + market_id=market_id, + city=city, + option=question, + price=trigger_price, + side="YES" if trigger_side == "Buy Yes" else "NO", + amount_usd=5.0, + target_date=target_date, + predicted_temp=ref_temp, + ) + + city_alerts.append( + { + "market": target_date or "今日", + "msg": msg, + "bought": success, + "amount": 5.0, + "confidence": "动态", + } + ) + pushed_signals[alert_key] = time.time() + if target_date: + city_target_date = target_date # C. 准备缓存数据 - temp_unit = weather_data.get("open-meteo", {}).get("unit", "celsius") + temp_unit = weather_data.get("open-meteo", {}).get( + "unit", "celsius" + ) temp_symbol = "°F" if temp_unit == "fahrenheit" else "°C" - city_local_time = weather_data.get("open-meteo", {}).get("current", {}).get("local_time") - + city_local_time = ( + weather_data.get("open-meteo", {}) + .get("current", {}) + .get("local_time") + ) + current_price = buy_yes_price if buy_yes_price else 0.5 - + # 计算价格趋势 prev_data = price_history.get(market_id, {}) prev_price = prev_data.get("price", current_price) - price_change_pct = ((current_price - prev_price) / prev_price * 100) if prev_price > 0 else 0 - + price_change_pct = ( + ((current_price - prev_price) / prev_price * 100) + if prev_price > 0 + else 0 + ) + # 更新价格历史缓存 price_history[market_id] = { "price": current_price, - "timestamp": datetime.now().isoformat() + "timestamp": datetime.now().isoformat(), } cache_entry = { @@ -334,9 +467,11 @@ def main(): } # --- 最终过滤器 (拦截垃圾信号) --- - + # 1. 过滤已锁定价格 (>= 98.5c) - if (buy_yes_price and buy_yes_price >= 0.985) or (buy_no_price and buy_no_price >= 0.985): + if (buy_yes_price and buy_yes_price >= 0.985) or ( + buy_no_price and buy_no_price >= 0.985 + ): cache_entry["rationale"] = "ENDED" all_markets_cache[market_id] = cache_entry continue @@ -361,7 +496,9 @@ def main(): whale_activity=None, ) cache_entry["score"] = signal.get("final_score", 0) - cache_entry["rationale"] = signal.get("recommendation", "ACTIVE") + cache_entry["rationale"] = signal.get( + "recommendation", "ACTIVE" + ) except Exception as e: logger.error(f"计算信号失败 [{market_id}]: {e}") cache_entry["score"] = 0 @@ -370,102 +507,293 @@ def main(): all_markets_cache[market_id] = cache_entry # --- 预警收集 (自动推送逻辑) --- - if (buy_yes_price and 0.85 <= buy_yes_price <= 0.95) or (buy_no_price and 0.85 <= buy_no_price <= 0.95): + if (buy_yes_price and 0.85 <= buy_yes_price <= 0.95) or ( + buy_no_price and 0.85 <= buy_no_price <= 0.95 + ): alert_key = f"alert_{market_id}_range_85_95" if alert_key not in pushed_signals: - trigger_side = ( - "Buy Yes" if buy_yes_price >= 0.85 else "Buy No" + # --- 基础参数识别 --- + is_categorical = len(ts) > 2 and active_tid + if is_categorical: + # 语义转换逻辑保持一致 + if buy_no_price and buy_no_price >= 0.85: + trigger_side = "Sell Yes" + trigger_price = int( + buy_no_price * 100 + ) # 预估价 + else: + trigger_side = "Buy Yes" + trigger_price = int(buy_yes_price * 100) + else: + trigger_side = ( + "Buy Yes" if buy_yes_price >= 0.85 else "Buy No" + ) + trigger_price = ( + int(buy_yes_price * 100) + if trigger_side == "Buy Yes" + else int(buy_no_price * 100) + ) + + # --- 深度流动性与 Spread 检查 --- + target_tid = ( + active_tid + if is_categorical + else (ts[0] if trigger_side == "Buy Yes" else ts[1]) ) - trigger_price = ( - int(buy_yes_price * 100) - if trigger_side == "Buy Yes" - else int(buy_no_price * 100) + ob_data = ( + polymarket.get_orderbook(target_tid) + if target_tid + else None ) + ob_analysis = { + "tradeable": True, + "liquidity": "未知", + "spread": 0, + "mid_price": trigger_price / 100, + } + if ob_data: + ob_analysis = orderbook_analyzer.analyze(ob_data) + + if not ob_analysis.get("tradeable", True): + confidence_tag = ( + f"🔴不可交易 ({ob_analysis.get('liquidity')})" + ) + if not is_categorical: + logger.warning( + f"跳过不可交易信号 (Spread {ob_analysis.get('spread')}): {city} {question}" + ) + continue + + # 更新实时数据显示 + mid_c = round(ob_analysis.get("mid_price", 0) * 100, 1) + spr_c = round(ob_analysis.get("spread", 0) * 100, 1) + depth = ob_analysis.get( + "ask_depth" + if trigger_side.startswith("Buy") + else "bid_depth", + 0, + ) + + # 流动性图标 + liq_map = { + "充裕": "✅ 充裕", + "正常": "🟡 正常", + "稀薄": "🟠 稀薄", + "枯竭": "🔴 枯竭", + } + liq_status = liq_map.get( + ob_analysis.get("liquidity", "未知"), "❓ 未知" + ) + + if is_categorical: + ask_str = ( + "--" + if trigger_side == "Sell Yes" + else f"{trigger_price}¢" + ) + bid_str = ( + f"{trigger_price}¢" + if trigger_side == "Sell Yes" + else "--" + ) + + display_side = ( + f"📊 {question}\n" + f"Ask: {ask_str} | Bid: {bid_str} | Mid: {mid_c}¢\n" + f"Spread: {spr_c}¢ | 深度: ${depth}\n" + f"流动性: {liq_status}" + ) + else: + display_side = ( + f"📊 {question}\n" + f"报价: {trigger_side} {trigger_price}¢ | Mid: {mid_c}¢\n" + f"Spread: {spr_c}¢ | 深度: ${depth}\n" + f"流动性: {liq_status}" + ) + # --- 智能动态仓位计算 --- # 1. 获取 Open-Meteo 对目标日期的最高温预测 predicted_high = None weather_supports = False - daily_data = weather_data.get("open-meteo", {}).get("daily", {}) + daily_data = weather_data.get("open-meteo", {}).get( + "daily", {} + ) if daily_data and target_date: dates = daily_data.get("time", []) max_temps = daily_data.get("temperature_2m_max", []) for idx, d_str in enumerate(dates): - if target_date == d_str and idx < len(max_temps): + if target_date == d_str and idx < len( + max_temps + ): predicted_high = max_temps[idx] break - + # 2. 判断天气预测是否支持当前方向 if predicted_high is not None: # 解析选项的温度范围 (例如 "40-41°F" 或 "32°F or below") - temp_match = re.search(r'(\d+)(?:-(\d+))?°[FC]', question) + temp_match = re.search( + r"(\d+)(?:-(\d+))?°[FC]", question + ) if temp_match: low_bound = int(temp_match.group(1)) - high_bound = int(temp_match.group(2)) if temp_match.group(2) else low_bound - + high_bound = ( + int(temp_match.group(2)) + if temp_match.group(2) + else low_bound + ) + # 如果买 NO,天气预测应该在这个区间之外 if trigger_side == "Buy No": - weather_supports = (predicted_high < low_bound - 2) or (predicted_high > high_bound + 2) + weather_supports = ( + predicted_high < low_bound - 2 + ) or (predicted_high > high_bound + 2) else: # 买 YES - weather_supports = (low_bound - 2 <= predicted_high <= high_bound + 2) - + weather_supports = ( + low_bound - 2 + <= predicted_high + <= high_bound + 2 + ) + # 3. 获取成交量信息 market_volume = market.get("volume", 0) if isinstance(market_volume, str): try: - market_volume = float(market_volume.replace("$", "").replace(",", "")) + market_volume = float( + market_volume.replace("$", "").replace( + ",", "" + ) + ) except: market_volume = 0 high_volume = market_volume >= 5000 # $5000+ 算高成交量 - - # 4. 动态仓位决策 - # 条件: 价格锁定程度 + 天气支持 + 成交量 - if trigger_price >= 90 and weather_supports and high_volume: - # 三重确认:重注 - amount_usd = 10.0 - confidence_tag = "🔥高置信" + + # --- Pro 级仓位决策系统 --- + # 1. 计算离结算剩余小时数 (假设气温市场在目标日期晚上 23:59 结算) + hours_to_settle = 24.0 + if target_date: + try: + settle_dt = datetime.strptime( + f"{target_date} 23:59:59", + "%Y-%m-%d %H:%M:%S", + ) + now_utc = datetime.utcnow() + diff = settle_dt - now_utc + hours_to_settle = diff.total_seconds() / 3600.0 + except: + pass + + # 2. 计算相对成交量比例 + total_daily_vol = sum( + [ + float( + str(m.get("volume", 0)) + .replace("$", "") + .replace(",", "") + ) + for m in city_markets + if ( + weather.extract_date_from_title( + m.get("event_title", "") + ) + or weather.extract_date_from_title( + m.get("question", "") + ) + ) + == target_date + ] + ) + market_vol = float( + str(market.get("volume", 0)) + .replace("$", "") + .replace(",", "") + ) + is_rel_high_vol = ( + (market_vol / total_daily_vol > 0.3) + if total_daily_vol > 0 + else False + ) + + # 3. 基础意向仓位 (基于置信度) + base_pos = 3.0 # 默认探路 + confidence_tag = "💡试探" + if ( + trigger_price >= 90 + and weather_supports + and high_volume + ): + base_pos, confidence_tag = 10.0, "🔥高置信" elif trigger_price >= 90 and weather_supports: - # 双重确认:中等仓位 - amount_usd = 7.0 - confidence_tag = "⭐中置信" + base_pos, confidence_tag = 7.0, "⭐中置信" elif trigger_price >= 92: - # 价格接近锁定,即使其他条件不满足也小额参与 - amount_usd = 5.0 - confidence_tag = "📌价格锁定" - else: - # 普通信号:最小仓位 - amount_usd = 3.0 - confidence_tag = "💡试探" - + base_pos, confidence_tag = 5.0, "📌价格锁定" + + # 4. 四层过滤决策 + amount_usd, risk_reason = ( + risk_manager.calculate_position_size( + base_confidence_usd=base_pos, + depth=depth, + hours_to_settle=hours_to_settle, + is_high_relative_volume=is_rel_high_vol, + ) + ) + logger.info( - f"【仓位决策】{city} {question} | " - f"价格:{trigger_price}¢ | 预测:{predicted_high} | 天气支持:{weather_supports} | " - f"高量:{high_volume} | 仓位:${amount_usd} ({confidence_tag})" + f"【Pro仓位】{city} {question} | " + f"基础:{base_pos}$ -> 最终:{amount_usd}$ | 原因:{risk_reason} | " + f"深度:${depth} | 剩:{hours_to_settle:.1f}h" ) # --- 模拟交易触发逻辑 --- - side = "YES" if trigger_side == "Buy Yes" else "NO" - success = paper_trader.open_position( - market_id=market_id, - city=city, - option=question, - price=trigger_price, - side=side, - amount_usd=amount_usd, - target_date=target_date, - predicted_temp=predicted_high, - ) + if amount_usd > 0: + side = "YES" if trigger_side == "Buy Yes" else "NO" + success = paper_trader.open_position( + market_id=market_id, + city=city, + option=question, + price=trigger_price, + side=side, + amount_usd=amount_usd, + target_date=target_date, + predicted_temp=predicted_high, + ) + if success: + risk_manager.record_trade(amount_usd) + else: + # 如果被风控拦截(金额为0),则不进行任何推送,避免刷屏 + success = False + logger.info( + f"Skipping alert for {question}: {risk_reason}" + ) + continue # 构建预测温度显示文本 - temp_unit = weather_data.get("open-meteo", {}).get("unit", "celsius") - temp_symbol = "°F" if temp_unit == "fahrenheit" else "°C" - forecast_text = f"预测:{predicted_high}{temp_symbol}" if predicted_high else "预测:N/A" + temp_unit = weather_data.get("open-meteo", {}).get( + "unit", "celsius" + ) + temp_symbol = ( + "°F" if temp_unit == "fahrenheit" else "°C" + ) + forecast_text = ( + f"{predicted_high}{temp_symbol}" + if predicted_high + else "N/A" + ) + + # 构建简约版消息: ⚡ {question} ({date}): {side} {price}¢ | 预测:{forecast} [🛒 ${amount} {tag}] + side_display = ( + "Buy No" if trigger_side == "Buy No" else "Buy Yes" + ) + msg = ( + f"⚡ {question} ({target_date}): {side_display} {trigger_price}¢ | " + f"预测:{forecast_text} [🛒 ${amount_usd} {confidence_tag}]" + ) city_alerts.append( { "type": "price", - "market": f"{question} ({target_date or '今日'})", - "msg": f"{trigger_side} {trigger_price}¢ | {forecast_text}", + "market": f"{target_date or '今日'}", + "msg": msg, "bought": success, "amount": amount_usd, "confidence": confidence_tag, @@ -476,15 +804,24 @@ def main(): # 3. 信号暂存 cached_signals[market_id] = cache_entry - # --- 循环结束后统一推送本城市汇总 --- - notifier.send_combined_alert( - city, city_alerts, local_time=city_local_time - ) + # E. 统一发送城市汇总通知 (使用新 Pro 模板) + if city_alerts: + # 去重策略建议 + unique_tips = list(dict.fromkeys(city_strategy_tips)) + notifier.send_combined_alert( + city=city, + alerts=city_alerts, + local_time=city_local_time, + forecast_temp=f"{city_pred_high}{temp_symbol}" + if city_pred_high + else "N/A", + total_volume=city_total_vol, + brackets_count=len(city_markets), + strategy_tips=unique_tips, + ) except Exception as e: logger.error(f"分析城市 {city} 时出错: {e}") - continue - # --- 每处理完一个城市,立即更新 JSON 文件 --- try: # --- 周期性结算:保存高价值信号 --- @@ -497,15 +834,17 @@ def main(): if target_dt and target_dt < "2026-02-06": continue active_signals.append(entry) - + # 按分数排序 active_signals.sort(key=lambda x: x.get("score", 0), reverse=True) - + with open("data/active_signals.json", "w", encoding="utf-8") as f: json.dump(active_signals, f, ensure_ascii=False, indent=4) - - logger.info(f"已更新活跃信号库,包含 {len(active_signals)} 个有效信号。") - + + logger.info( + f"已更新活跃信号库,包含 {len(active_signals)} 个有效信号。" + ) + # 2. 更新全量市场缓存 try: with open("data/all_markets.json", "r", encoding="utf-8") as f: diff --git a/src/analysis/orderbook_analyzer.py b/src/analysis/orderbook_analyzer.py index d0abf2b2..4783d1d5 100644 --- a/src/analysis/orderbook_analyzer.py +++ b/src/analysis/orderbook_analyzer.py @@ -9,65 +9,83 @@ class OrderbookAnalyzer: self.wall_threshold = self.config.get("wall_threshold", 500) # 单笔订单超过此值为墙 logger.info("Initializing Orderbook Analyzer...") + def assess_liquidity(self, orderbook, side="ask"): + """ + 分析流动性深度 (基于前 3 档) + """ + orders = orderbook.get('asks' if side == "ask" else 'bids', []) + if not orders: + return "枯竭", 0 + + # 前 3 档总量 (Polymarket 通常返回价格字符串) + depth = sum(float(o.get("size", 0)) for o in orders[:3]) + + if depth < 50: + return "稀薄", depth + elif depth < 500: + return "正常", depth + else: + return "充裕", depth + def analyze(self, orderbook): """ - 订单簿分析决策 - - Args: - orderbook: dict 包含 'bids' 和 'asks' 列表 + 增强版订单簿分析:集成深度与 Spread 评估 """ bids = orderbook.get('bids', []) asks = orderbook.get('asks', []) if not bids or not asks: - return {"signal": "NEUTRAL", "confidence": 0.5, "reason": "Empty orderbook"} + return { + "signal": "NEUTRAL", + "confidence": 0.0, + "tradeable": False, + "reason": "缺乏双边报价", + "liquidity": "枯竭", + "spread": 1.0 + } - # 1. 计算买卖力量对比 (Imbalance) - # Polymarket API 返回的通常是 [{"price": "0.90", "size": "100"}, ...] - bid_volume = sum([float(b.get('size', 0)) for b in bids]) - ask_volume = sum([float(a.get('size', 0)) for a in asks]) - - imbalance = bid_volume / ask_volume if ask_volume > 0 else 0 - - # 2. 识别墙单 - max_bid = max([float(b.get('size', 0)) for b in bids]) if bids else 0 - max_ask = max([float(a.get('size', 0)) for a in asks]) if asks else 0 - - # 3. 计算价差 (Spread) + # 1. 计算核心指标 best_bid = float(bids[0].get('price', 0)) best_ask = float(asks[0].get('price', 0)) - spread = (best_ask - best_bid) / best_ask if best_ask > 0 else 0 + spread = abs(best_ask - best_bid) + mid_price = (best_ask + best_bid) / 2 + + # 2. 评估流动性 + ask_liq, ask_depth = self.assess_liquidity(orderbook, "ask") + bid_liq, bid_depth = self.assess_liquidity(orderbook, "bid") + + # 3. 交易可行性判定 (Spread <= 10c 且 深度 >= $50) + is_tradeable = (spread <= 0.10) and (ask_depth >= 50 or bid_depth >= 50) + + # 4. Imbalance 计算 + bid_volume = sum([float(b.get('size', 0)) for b in bids]) + ask_volume = sum([float(a.get('size', 0)) for a in asks]) + imbalance = bid_volume / ask_volume if ask_volume > 0 else 0 result = { + "best_bid": best_bid, + "best_ask": best_ask, + "mid_price": mid_price, + "spread": round(spread, 4), + "ask_depth": round(ask_depth, 2), + "bid_depth": round(bid_depth, 2), + "liquidity": ask_liq if ask_depth < bid_depth else bid_liq, + "tradeable": is_tradeable, "imbalance": imbalance, - "bid_volume": bid_volume, - "ask_volume": ask_volume, - "max_bid_wall": max_bid, - "max_ask_wall": max_ask, - "spread": spread, "signal": "NEUTRAL", "confidence": 0.5 } - # 4. 决策逻辑 - if imbalance > 2.0: - result["signal"] = "BULLISH" - result["confidence"] = min(0.9, 0.5 + (imbalance - 1) / 4) - elif imbalance < 0.5: - result["signal"] = "BEARISH" - result["confidence"] = min(0.9, 0.5 + (1 / imbalance - 1) / 4) - - if max_bid > self.wall_threshold and bid_volume > ask_volume: - result["signal"] = "STRONG_BUY" - result["confidence"] = 0.85 - elif max_ask > self.wall_threshold and ask_volume > bid_volume: - result["signal"] = "STRONG_SELL" - result["confidence"] = 0.85 - - # 5. 流动性警告 - if spread > 0.05: # 价差超过5% - result["warning"] = "LOW_LIQUIDITY" - result["confidence"] *= 0.8 # 降低置信度 + # 5. 信号修正 + if is_tradeable: + if imbalance > 2.5: + result["signal"] = "BULLISH" + result["confidence"] = 0.75 + elif imbalance < 0.4: + result["signal"] = "BEARISH" + result["confidence"] = 0.75 + else: + result["confidence"] = 0.1 # 不建议交易 return result diff --git a/src/data_collection/polymarket_api.py b/src/data_collection/polymarket_api.py index 4c53effb..136a2527 100644 --- a/src/data_collection/polymarket_api.py +++ b/src/data_collection/polymarket_api.py @@ -5,9 +5,11 @@ import re from typing import Dict, List, Optional from loguru import logger from datetime import datetime +from concurrent.futures import ThreadPoolExecutor from py_clob_client.client import ClobClient from py_clob_client.constants import POLYGON +from py_clob_client.clob_types import ApiCreds, BookParams, OpenOrderParams class PolymarketClient: @@ -19,6 +21,11 @@ class PolymarketClient: self.base_url = config.get("base_url", "https://clob.polymarket.com") self.timeout = config.get("timeout", 20) self.session = requests.Session() + + # 缓存机制 + self._weather_markets_cache = [] + self._last_discovery_time = 0 + self._cache_ttl = 300 # 5 分钟缓存 # 统一代理设置 proxy = os.getenv("HTTPS_PROXY") or os.getenv("HTTP_PROXY") @@ -40,7 +47,6 @@ class PolymarketClient: self.api_passphrase = config.get("api_passphrase") try: - from py_clob_client.clob_types import ApiCreds # 组装凭据对象 (如果提供) creds = None @@ -102,9 +108,9 @@ class PolymarketClient: 获取 Token 的实时盘口价格 (纯官方库实现) """ try: - # 官方语义:BUY 对应的是我们的买入成本 (Ask) - side_val = "BUY" if side == "ask" else "SELL" - price_str = self.clob_client.get_price(token_id=token_id, side=side_val) + # 用户侧语义:BUY 代表我要买 (Ask),SELL 代表我要卖 (Bid) + sdk_side = "BUY" if side == "ask" else "SELL" + price_str = self.clob_client.get_price(token_id=token_id, side=sdk_side) if price_str: return float(price_str) except Exception as e: @@ -165,45 +171,59 @@ class PolymarketClient: def get_multiple_prices(self, token_requests: List[Dict]) -> Dict[str, float]: """ - 批量获取多个 token 的价格 (官方接口实现) + 批量获取多个 token 的价格 (官方接口 + 线程池并行实现) """ if not token_requests: return {} all_prices = {} - try: - # 准备官方批量请求格式 - # 我们映射 ask->BUY, bid->SELL - batch_req = [] - for r in token_requests: - side_val = "BUY" if r.get("side") == "ask" else "SELL" - batch_req.append({"token_id": r["token_id"], "side": side_val}) + batch_size = 20 + + def robust_float(val): + if isinstance(val, (int, float)): return float(val) + if isinstance(val, str): + try: return float(val) + except: return 0.0 + return 0.0 - # 使用官方批量获取接口 - results = self.clob_client.get_prices(batch_req) - - def robust_float(val): - if isinstance(val, (int, float)): return float(val) - if isinstance(val, str): - try: return float(val) - except: return 0.0 - return 0.0 + chunks = [token_requests[i : i + batch_size] for i in range(0, len(token_requests), batch_size)] + + def fetch_chunk(chunk): + for attempt in range(3): + try: + batch_req = [] + for r in chunk: + sdk_side = "BUY" if r.get("side") == "ask" else "SELL" + batch_req.append(BookParams(token_id=r["token_id"], side=sdk_side)) + + results = self.clob_client.get_prices(batch_req) + + chunk_prices = {} + if isinstance(results, list): + for item in results: + tid = item.get("token_id") + price_raw = item.get("price") + res_side = item.get("side") + if tid and price_raw: + val = robust_float(price_raw) + key_side = "ask" if res_side == "BUY" else "bid" + chunk_prices[f"{tid}:{key_side}"] = val + return chunk_prices + except Exception as e: + if attempt < 2: + time.sleep(0.5 * (attempt + 1)) + continue + logger.warning(f"Batch fetch failed after 3 attempts: {e}") + return {} - if isinstance(results, list): - for item in results: - tid = item.get("token_id") - price_raw = item.get("price") - side = item.get("side") - if tid and price_raw: - val = robust_float(price_raw) - key_side = "ask" if side == "BUY" else "bid" - all_prices[f"{tid}:{key_side}"] = val - all_prices[tid] = val + # 使用更保守的线程池并发抓取 + with ThreadPoolExecutor(max_workers=3) as executor: + future_results = list(executor.map(fetch_chunk, chunks)) - return all_prices - except Exception as e: - logger.warning(f"官方库批量获取报价失败: {e}") - return {} + for chunk_result in future_results: + all_prices.update(chunk_result) + + return all_prices def get_midpoint(self, token_id: str) -> Optional[float]: """ @@ -267,10 +287,28 @@ class PolymarketClient: logger.error(f"取消订单失败: {e}") return None + def get_orders(self, market_id: str = None) -> Optional[Dict]: + """ + 获取当前活跃挂单 + """ + try: + params = OpenOrderParams(market=market_id) if market_id else None + return self.clob_client.get_orders(params=params) + except Exception as e: + logger.error(f"获取挂单失败: {e}") + return None + def discover_weather_markets(self) -> list: """ - 通过全量扫描活跃事件发现最高温天气市场。 + 通过全量扫描活跃事件发现最高温天气市场 (支持缓存机制) """ + # 缓存检查 + current_time = time.time() + if self._weather_markets_cache and (current_time - self._last_discovery_time < self._cache_ttl): + logger.debug(f"使用缓存的市场列表 (剩余寿命: {int(self._cache_ttl - (current_time - self._last_discovery_time))}s)") + return self._weather_markets_cache + + logger.info("📡 正在全量扫描 Polymarket 发现天气市场...") gamma_url = "https://gamma-api.polymarket.com/events" all_weather_markets = [] seen_condition_ids = set() @@ -290,13 +328,23 @@ class PolymarketClient: for m in event.get("markets", []): question = m.get("groupItemTitle") or m.get("question") or "" - # 关键词匹配 - if not ( - is_weather_event - or "Highest temperature" in question - or "temperature in" in question.lower() - ): + # 强化过滤:必须在标题中包含明确的气温气象词,且排除非气温市场 + t_lower = title.lower() + q_lower = question.lower() + + # 1. 标题必须像个气温市场 + if not any(k in t_lower for k in ["highest temperature", "high temperature", "will temperature", "daily temperature"]): continue + + # 2. 排除干扰项 + if "climate" in t_lower or "rain" in t_lower or "snow" in t_lower: + continue + + # 3. 确保这个具体的 market (bracket) 是我们想要的 + if not any(k in q_lower for k in ["temperature", "be", "highest", "range"]): + # 补充:如果是多选一市场的子项,question 可能只是一个数字或范围,此时看 title + if not any(k in t_lower for k in ["temperature", "highest"]): + continue c_id = m.get("conditionId") # 识别 outcome_index @@ -429,6 +477,11 @@ class PolymarketClient: logger.info( f"全量发现结束,共获取 {len(all_weather_markets)} 个天气档位合约" ) + + # 更新缓存 + self._weather_markets_cache = all_weather_markets + self._last_discovery_time = current_time + return all_weather_markets except Exception as e: diff --git a/src/data_collection/weather_sources.py b/src/data_collection/weather_sources.py index 748a5ee4..47d01798 100644 --- a/src/data_collection/weather_sources.py +++ b/src/data_collection/weather_sources.py @@ -270,6 +270,7 @@ class WeatherDataCollector: static_coords = { "london": {"lat": 51.5074, "lon": -0.1278}, "new york": {"lat": 40.7128, "lon": -74.0060}, + "new york's central park": {"lat": 40.7812, "lon": -73.9665}, "nyc": {"lat": 40.7128, "lon": -74.0060}, "seattle": {"lat": 47.6062, "lon": -122.3321}, "chicago": {"lat": 41.8781, "lon": -87.6298}, @@ -286,6 +287,12 @@ class WeatherDataCollector: normalized_city = city.lower().strip() if normalized_city in static_coords: return static_coords[normalized_city] + + # 模糊匹配映射 (针对包含城市名的情况) + for key in static_coords: + if key in normalized_city: + logger.debug(f"地理编码命中模糊映射: {city} -> {key}") + return static_coords[key] try: url = "https://geocoding-api.open-meteo.com/v1/search" @@ -377,7 +384,8 @@ class WeatherDataCollector: "phoenix", "philadelphia", ] - use_fahrenheit = city.lower() in us_cities + city_lower = city.lower() + use_fahrenheit = any(uc in city_lower for uc in us_cities) # Open-Meteo (Primary Free Source - No Key) if lat and lon: diff --git a/src/models/statistical_model.py b/src/models/statistical_model.py index 56d325df..df31e956 100644 --- a/src/models/statistical_model.py +++ b/src/models/statistical_model.py @@ -258,10 +258,10 @@ class TemperaturePredictor: weights.append(rf_weight) if not predictions: - logger.warning("No predictions available") + logger.debug("No predictions available (Model not trained)") return { "predicted_temp": None, - "confidence": 0.0, + "confidence": 0.5, "error": "No models available for prediction" } diff --git a/src/strategy/risk_manager.py b/src/strategy/risk_manager.py index 29169817..967960e3 100644 --- a/src/strategy/risk_manager.py +++ b/src/strategy/risk_manager.py @@ -7,149 +7,91 @@ class RiskManager: def __init__(self, config=None): self.config = config or {} - self.max_single_trade = self.config.get("max_single_trade", 500) # 最大单笔$500 - self.max_drawdown = self.config.get("max_drawdown", 0.10) # 最大回撤10% - self.min_liquidity = self.config.get("min_liquidity", 1000) # 最小流动性$1000 - self.max_slippage = self.config.get("max_slippage", 0.02) # 最大滑点2% - self.min_confidence = self.config.get("min_confidence", 0.65) # 最小置信度65% - + # 基础风控参数 + self.max_single_trade = self.config.get("max_single_trade", 50.0) # 最大单笔调整为 $50 + self.max_daily_exposure = 50.0 # 每日最高投入上限 + self.daily_used_exposure = 0.0 + self.last_reset_date = "" + + self.min_confidence = 0.5 self.peak_capital = 0 - self.current_drawdown = 0 self.is_trading_paused = False - logger.info("Initializing Risk Manager...") + logger.info("Initializing Pro Risk Manager...") - def check_trade_risk(self, - trade_size: float, - market_data: dict, - model_confidence: float) -> dict: + def _reset_daily_exposure(self): + """每日重置额度""" + from datetime import datetime + today = datetime.now().strftime("%Y-%m-%d") + if self.last_reset_date != today: + self.daily_used_exposure = 0.0 + self.last_reset_date = today + logger.info(f"Daily exposure reset for {today}") + + def calculate_position_size(self, + base_confidence_usd: float, + depth: float, + hours_to_settle: float, + is_high_relative_volume: bool) -> tuple[float, str]: """ - 检查单笔交易风险 - - Args: - trade_size: 交易金额 - market_data: 市场数据 (包含订单簿等) - model_confidence: 模型置信度 - - Returns: - dict: 风险检查结果 + 四层过滤仓位计算方法: + 仓位 = base_position(置信度) + × liquidity_factor(深度/仓位 >= 5x) + × time_decay(离结算衰减) + × budget_limit """ - risks = [] - passed = True + self._reset_daily_exposure() - # 1. 检查交易金额 - if trade_size > self.max_single_trade: - risks.append({ - "type": "TRADE_SIZE", - "message": f"Trade size ${trade_size:.2f} exceeds max ${self.max_single_trade}" - }) - passed = False - - # 2. 检查置信度 - if model_confidence < self.min_confidence: - risks.append({ - "type": "LOW_CONFIDENCE", - "message": f"Model confidence {model_confidence:.2f} below threshold {self.min_confidence}" - }) - passed = False - - # 3. 检查流动性 - orderbook = market_data.get("orderbook", {}) - total_liquidity = self._calculate_liquidity(orderbook) - if total_liquidity < self.min_liquidity: - risks.append({ - "type": "LOW_LIQUIDITY", - "message": f"Market liquidity ${total_liquidity:.2f} below threshold ${self.min_liquidity}" - }) - passed = False - - # 4. 检查滑点 - expected_slippage = self._estimate_slippage(trade_size, orderbook) - if expected_slippage > self.max_slippage: - risks.append({ - "type": "HIGH_SLIPPAGE", - "message": f"Expected slippage {expected_slippage:.2%} exceeds max {self.max_slippage:.2%}" - }) - passed = False - - # 5. 检查是否暂停交易 - if self.is_trading_paused: - risks.append({ - "type": "TRADING_PAUSED", - "message": "Trading is paused due to drawdown limit" - }) - passed = False - - return { - "passed": passed, - "risks": risks, - "liquidity": total_liquidity, - "expected_slippage": expected_slippage - } + final_pos = base_confidence_usd + reason = "Normal" - def _calculate_liquidity(self, orderbook: dict) -> float: - """计算订单簿总流动性""" - bids = orderbook.get("bids", []) - asks = orderbook.get("asks", []) + # 1. 流动性过滤: 深度 < $50 强制跳过; 深度 < 仓位的 5 倍则缩减 + if depth < 50: + return 0.0, "🚫深度不足 (min $50)" - bid_liquidity = sum(float(b.get("size", 0)) for b in bids) - ask_liquidity = sum(float(a.get("size", 0)) for a in asks) - - return bid_liquidity + ask_liquidity + if depth < final_pos * 5: + # 如果深度不足以承载期望仓位,按比例缩减至深度的 1/5 + final_pos = depth / 5.0 + reason = "⚠️深度限流" - def _estimate_slippage(self, trade_size: float, orderbook: dict) -> float: - """估算滑点""" - asks = orderbook.get("asks", []) - if not asks: - return 0.05 # 无数据时假设5%滑点 + # 2. 时间衰减因子 + # 离结算时间越近,预测越准但也存在剧烈博弈风险 + time_factor = 1.0 + if hours_to_settle <= 1.0: + time_factor = 0.0 # 最后 1 小时停止建仓 + reason = "🚫临近结算" + elif hours_to_settle <= 4.0: + time_factor = 0.4 # 1-4小时:缩小 60% + reason = "⏱️结算冲刺 (40%)" + elif hours_to_settle <= 12.0: + time_factor = 0.7 # 4-12小时:缩小 30% + reason = "⏳接近结算 (70%)" - best_ask = float(asks[0].get("price", 0)) if asks else 0 - if best_ask == 0: - return 0.05 - - # 简单估算:交易额 / 流动性 * 基础滑点 - ask_liquidity = sum(float(a.get("size", 0)) for a in asks) - if ask_liquidity == 0: - return 0.05 - - impact_ratio = trade_size / ask_liquidity - estimated_slippage = impact_ratio * 0.1 # 假设10%的市场冲击系数 - - return min(estimated_slippage, 0.1) # 最大10% + final_pos *= time_factor + if final_pos <= 0: return 0.0, reason - def update_drawdown(self, current_capital: float) -> dict: - """ - 更新回撤状态 + # 3. 预算上限过滤 + remaining_daily = self.max_daily_exposure - self.daily_used_exposure + if remaining_daily <= 0: + return 0.0, "🚫今日总额度已满 ($50)" - Args: - current_capital: 当前资金 - - Returns: - dict: 回撤状态 - """ - # 更新峰值 - if current_capital > self.peak_capital: - self.peak_capital = current_capital - - # 计算回撤 - if self.peak_capital > 0: - self.current_drawdown = (self.peak_capital - current_capital) / self.peak_capital - else: - self.current_drawdown = 0 - - # 检查是否需要暂停交易 - if self.current_drawdown >= self.max_drawdown: - self.is_trading_paused = True - logger.warning(f"Trading PAUSED! Drawdown {self.current_drawdown:.2%} exceeds limit {self.max_drawdown:.2%}") - - return { - "peak_capital": self.peak_capital, - "current_capital": current_capital, - "drawdown": self.current_drawdown, - "is_paused": self.is_trading_paused - } + if final_pos > remaining_daily: + final_pos = remaining_daily + reason = "🛑触及日风控上限" - def resume_trading(self): - """手动恢复交易""" - self.is_trading_paused = False - logger.info("Trading resumed manually") + # 4. 高相对成交量加权 (如果是高成交量市场,且逻辑支持,可保持原状或微增) + # 这里逻辑设定为:如果不是高成交量,再次缩减 20% 防御 + if not is_high_relative_volume: + final_pos *= 0.8 + if reason == "Normal": reason = "📉低活缩减" + + return round(final_pos, 2), reason + + def record_trade(self, amount: float): + """记录成交额以扣除额度""" + self.daily_used_exposure += amount + logger.debug(f"Applied exposure: ${amount}. Daily Total: ${self.daily_used_exposure}") + + def check_trade_risk(self, trade_size: float, market_data: dict, model_confidence: float) -> dict: + """保持基础接口兼容""" + return {"passed": True, "risks": []} diff --git a/src/utils/notifier.py b/src/utils/notifier.py index 0765e4d1..5fdc7136 100644 --- a/src/utils/notifier.py +++ b/src/utils/notifier.py @@ -45,7 +45,7 @@ class TelegramNotifier: for cid in chat_ids: if not cid: continue - + payload = { "chat_id": cid, "text": text, @@ -125,33 +125,49 @@ class TelegramNotifier: ) return self._send_message(text) - def send_combined_alert(self, city: str, alerts: list, local_time: str = None): - """发送合并后的城市预警""" + def send_combined_alert( + self, + city: str, + alerts: list, + local_time: str = None, + forecast_temp: str = None, + total_volume: float = 0, + brackets_count: int = 0, + strategy_tips: list = None, + ): + """发送简约版合并预警""" if not alerts: return from datetime import datetime, timedelta # UTC+8 北京时间 - timestamp_bj = (datetime.utcnow() + timedelta(hours=8)).strftime("%H:%M") + now_bj = datetime.utcnow() + timedelta(hours=8) + timestamp_bj = now_bj.strftime( + "%H:%M" + ) # 简化为仅显示时间,日期通常与当地一致或不重要 + # 1. 信号详情构建 items_text = "" for a in alerts: - type_icon = "⚡" if a["type"] == "price" else "🐋" - # 买入标签:显示金额 - if a.get("bought"): - amount = a.get("amount", 5.0) - confidence = a.get("confidence", "") - buy_tag = f" [🛒 ${amount} {confidence}]" - else: - buy_tag = "" - items_text += f"{type_icon} {a['market']}: {a['msg']}{buy_tag}\n" + items_text += f"{a['msg']}\n\n" + # 2. 策略建议(如果有) + tips_text = "" + if strategy_tips: + tips_text = ( + "💡 策略建议:\n" + + "\n".join([f"• {self._escape_html(tip)}" for tip in strategy_tips]) + + "\n\n" + ) + + # 3. 总体布局 (回归清爽风格) text = ( f"🔔 城市监控报告 #{self._escape_html(city)}\n\n" f"📍 城市: {self._escape_html(city)}\n" f"📊 实时异动:\n" - f"{items_text}\n" + f"{items_text}" + f"{tips_text}" f"═══════════════════\n" f"🕒 当地时间: {self._escape_html(local_time or 'N/A')}\n" f"⏰ 预警时间: {timestamp_bj} (北京时间)" diff --git a/test_template.py b/test_template.py new file mode 100644 index 00000000..186e1caf --- /dev/null +++ b/test_template.py @@ -0,0 +1,65 @@ + +import os +import json +from src.utils.config_loader import load_config +from src.utils.notifier import TelegramNotifier + +def send_test_template(): + config_data = load_config() + notifier = TelegramNotifier(config_data["telegram"]) + + city = "Nyc" + target_date = "2026-02-07" + + # 模拟异动信号数据 + alerts = [ + { + "market": target_date, + "msg": ( + "🟢 26°F+\n" + "执行动作: BUY YES ⬆️\n" + "Ask: 80¢ | Bid: -- | Mid: 79.5¢\n" + "Spread: 1.2¢ | 深度: $1,847\n" + "流动性: ✅ 充裕 | 可交易: ✅\n" + "📐 预测偏差: -3.0°F (预测 23.0°F)" + ), + "bought": True, + "amount": 7.0, + "confidence": "⭐中置信" + }, + { + "market": target_date, + "msg": ( + "🔴 18-19°F\n" + "执行动作: SELL YES ⬇️\n" + "Ask: -- | Bid: 5.0¢ | Mid: 4.5¢\n" + "Spread: 0.5¢ | 深度: $312\n" + "流动性: ✅ 正常 | 可交易: ✅\n" + "📐 预测偏差: -4.0°F (预测 23.0°F)" + ), + "bought": False, + "amount": 0.0, + "confidence": "" + } + ] + + strategy_tips = [ + "预测温度 23.0°F 落在 22-23°F 区间,市场与模型一致", + "26°F+ 区间出现主力大额买入,建议跟随", + "18-19°F 流动性正常但偏差过大,已执行调仓" + ] + + print("🚀 正在发送测试模板到 Telegram...") + notifier.send_combined_alert( + city=city, + alerts=alerts, + local_time="10:56 EST", + forecast_temp="23.0°F", + total_volume=40113, + brackets_count=7, + strategy_tips=strategy_tips + ) + print("✅ 发送成功!请检查手机。") + +if __name__ == "__main__": + send_test_template()