#!/usr/bin/env python3 """batch_reassess.py — 批量补全12维矩阵LLM分析(逐只处理,间隔防限流) 用法: python3 batch_reassess.py # 所有缺分析/过期的 active 策略 python3 batch_reassess.py --type holding # 只处理持仓策略 python3 batch_reassess.py --type watchlist # 只处理自选策略 python3 batch_reassess.py --type holding --today # 持仓每日刷新(今早未评过的强制重评) python3 batch_reassess.py --code XXXXXX # 单只 流程:收集最新数据 → 调LLM(gateway)写12维分析+策略 → 保存到DB """ import sys, json, subprocess, sqlite3, re, time, os from datetime import datetime # ── 共享 LLM 客户端 + DB 工具(profile-scripts 硬链到同目录)── sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) sys.path.insert(0, "/home/hmo/MoFin") from llm_client import call_llm, REASSESS_MODEL, FALLBACK_MODEL, gateway_alive, ocg_alive from mofin_db import snapshot_strategy_history, sync_recommend_tag DB = "/home/hmo/MoFin/data/mofin.db" COOLDOWN_HOURS = 1 STALE_HOURS = 20 # 分析超过20小时视为过期,需要重评 def has_llm_analysis(code): """检查是否为LLM生成的12维分析(>500字)""" conn = sqlite3.connect(DB, timeout=30) r = conn.execute("SELECT LENGTH(full_analysis) FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() conn.close() return r and r[0] and r[0] > 500 FORCE_REASSESS = "--force" in sys.argv def in_cooldown(code): """冷却期检查(--force 时全量强制重评)""" if FORCE_REASSESS: return False conn = sqlite3.connect(DB, timeout=30) r = conn.execute("SELECT reassessed_at FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() conn.close() if not r or not r[0]: return False try: last = datetime.fromisoformat(r[0]) diff = (datetime.now() - last).total_seconds() / 3600 return diff < COOLDOWN_HOURS except: return False def analysis_stale(code, force_today=False): """分析是否过期(>STALE_HOURS 或 force_today 时今早4点前未重评)""" conn = sqlite3.connect(DB, timeout=30) r = conn.execute("SELECT reassessed_at FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() conn.close() if not r or not r[0]: return True try: last = datetime.fromisoformat(r[0]) if force_today: today4am = datetime.now().replace(hour=4, minute=0, second=0, microsecond=0) return last < today4am return (datetime.now() - last).total_seconds() / 3600 > STALE_HOURS except: return True def get_portfolio(): """从 portfolio_summary 读实时现金/总资产(不再硬编码)""" try: conn = sqlite3.connect(DB, timeout=30) r = conn.execute("SELECT cash, total_assets FROM portfolio_summary WHERE id=1").fetchone() conn.close() if r and r[1]: return int(r[0] or 0), int(r[1]) except Exception: pass return 0, 0 def collect_data(code): """收集最新数据(含完整策略原文)""" data = {"code": code} # 从DB读策略(含 full_analysis / changelog_json / position_advice) conn = sqlite3.connect(DB, timeout=30) r = conn.execute("SELECT name, entry_low, entry_high, stop_loss, take_profit, timing_signal, action, rr_ratio, tech_snapshot, sector_context, stock_category, full_analysis, changelog_json, reassessed_at, position_advice, strategy_name FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() if r: data["name"] = r[0] data["entry_low"] = r[1] or 0 data["entry_high"] = r[2] or 0 data["stop_loss"] = r[3] or 0 data["take_profit"] = r[4] or 0 data["timing_signal"] = r[5] or "" data["action"] = r[6] or "" data["rr_ratio"] = r[7] or 0 data["tech_snapshot"] = r[8] or "" data["sector_context"] = r[9] or "" data["stock_category"] = r[10] or "" data["full_analysis"] = r[11] or "" data["changelog_json"] = r[12] or "" data["reassessed_at"] = r[13] or "" data["strategy_name"] = r[15] or "" # 2026-08-18 来源策略 # 2026-08-18 策略状态感知:switched → 用 strategy_attributed(切到的策略)的定义评估 try: _st0 = conn.execute("SELECT strategy_state, strategy_attributed FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() _eval_strat = data["strategy_name"] if _st0 and _st0[0] == "switched" and _st0[1]: _eval_strat = _st0[1] except Exception: _eval_strat = data["strategy_name"] # 2026-08-18 策略语义注入:读 strategy_defs 定义卡 + 情势体检 try: _sd = conn.execute("SELECT display_name, summary, entry_logic, exit_logic, review_focus, holding_style, status, retired_reason, superseded_by, regime FROM strategy_defs WHERE strategy_name=?", (_eval_strat,)).fetchone() if _sd: data["strategy_def"] = { "display_name": _sd[0] or data["strategy_name"], "summary": _sd[1] or "", "entry_logic": _sd[2] or "", "exit_logic": _sd[3] or "", "review_focus": _sd[4] or "", "holding_style": _sd[5] or "", "status": _sd[6] or "active", "retired_reason": _sd[7] or "", "superseded_by": _sd[8] or "", "regime": _sd[9] or "", } else: data["strategy_def"] = None except Exception: data["strategy_def"] = None # ── 技术指标收集(策略输入数据)── # 从 stock_indicators 读取实时技术指标(由 price_monitor 联动更新) try: _si = _db2.execute( "SELECT ma5, ma10, ma20, ma60, rsi, bias60, dist_ma20, strong_support, weak_support, pivot, weak_resist, strong_resist, candle_pattern, vol_ratio " "FROM stock_indicators WHERE code=? ORDER BY date DESC LIMIT 1", (code,)).fetchone() if _si: data["ta_ma5"] = _si[0] if _si[0] else None data["ta_ma10"] = _si[1] if _si[1] else None data["ta_ma20"] = _si[2] if _si[2] else None data["ta_ma60"] = _si[3] if _si[3] else None data["ta_rsi"] = round(_si[4], 1) if _si[4] else None data["ta_dist_ma20"] = _si[6] if _si[6] else None data["ta_strong_support"] = _si[7] if _si[7] else None data["ta_weak_support"] = _si[8] if _si[8] else None data["ta_pivot"] = _si[9] if _si[9] else None data["ta_weak_resist"] = _si[10] if _si[10] else None data["ta_strong_resist"] = _si[11] if _si[11] else None data["ta_candle"] = _si[12] or "" data["ta_volume"] = "" except Exception: pass try: import multi_timeframe as _mtf _mtf_r = _mtf.full_multi_tf_analysis(code) if _mtf_r: _adj = _mtf_r.get("strategy_adjustment", {}) data["mtf_trend_alignment"] = _adj.get("trend_alignment", "未知") data["mtf_daily_trend"] = _mtf_r.get("daily", {}).get("trend", {}).get("description", "") data["mtf_weekly_trend"] = _mtf_r.get("weekly", {}).get("trend", {}).get("description", "") data["mtf_monthly_trend"] = _mtf_r.get("monthly", {}).get("trend", {}).get("description", "") _rsi = _mtf_r.get("daily", {}).get("rsi") if _rsi is not None: data["ta_rsi"] = round(_rsi, 1) except Exception: pass try: _db2 = sqlite3.connect(DB, timeout=30) _fi = _db2.execute("SELECT mcap_q, pe_q, bias60, rsi, r5f, dist_lo20, dist_ma20, vol_ratio FROM stock_indicators WHERE code=? ORDER BY date DESC LIMIT 1", (code,)).fetchone() if _fi: data["factor_mcap_q"] = _fi[0] if _fi[0] is not None else None data["factor_pe_q"] = _fi[1] if _fi[1] is not None else None data["factor_bias60"] = _fi[2] if _fi[2] is not None else None if data.get("ta_rsi") is None and _fi[3] is not None: data["ta_rsi"] = round(_fi[3], 1) data["factor_ret5d"] = _fi[4] if _fi[4] is not None else None data["factor_dist_lo20"] = _fi[5] if _fi[5] is not None else None data["factor_dist_ma20"] = _fi[6] if _fi[6] is not None else None data["factor_vol_ratio"] = _fi[7] if _fi[7] is not None else None _db2.close() except Exception: pass # 行业强度(sector_snapshots: market_watch.py 写入的板块涨跌排名) try: _db3 = sqlite3.connect(DB, timeout=30) _sector = data.get("sector_context", "") if _sector: _ss = _db3.execute( "SELECT change_pct, rank_in_market FROM sector_snapshots " "WHERE sector_name LIKE ? ORDER BY date DESC LIMIT 1", (f"%{_sector}%",)).fetchone() if _ss: data["sector_change_pct"] = _ss[0] if _ss[0] is not None else None data["sector_rank"] = _ss[1] if _ss[1] is not None else None _db3.close() except Exception: pass # 情势体检:温区 + 高风险消息 + 执行红线 _sit = {"regime_a": "unknown", "regime_7d_ago": "unknown", "high_risk": "", "breach_stop": "否", "reach_tp": "否", "out_zone": "否", "over_hold": "否"} try: _m = conn.execute("SELECT regime, date FROM market_regime WHERE market='a' ORDER BY date DESC LIMIT 1").fetchone() if _m: _sit["regime_a"] = _m[0] or "unknown" _m7 = conn.execute("SELECT regime FROM market_regime WHERE market='a' AND date <= date(?, '-7 days') ORDER BY date DESC LIMIT 1", (_m[1],)).fetchone() if _m7: _sit["regime_7d_ago"] = _m7[0] or "unknown" except Exception: pass try: _hr = conn.execute("SELECT summary FROM signal_news WHERE overall_sentiment LIKE '%HIGH%' OR summary LIKE '%高风险%' ORDER BY id DESC LIMIT 1").fetchone() if _hr: _sit["high_risk"] = (_hr[0] or "")[:60] except Exception: pass try: _p = data.get("price") or 0 _el2, _eh2 = data.get("entry_low") or 0, data.get("entry_high") or 0 _sl2, _tp2 = data.get("stop_loss") or 0, data.get("take_profit") or 0 if _p > 0 and _sl2 > 0 and _p < _sl2: _sit["breach_stop"] = "是" if _p > 0 and _tp2 > 0 and _p >= _tp2: _sit["reach_tp"] = "是" if _eh2 > 0 and _p > _eh2 * 1.05: _sit["out_zone"] = "是" except Exception: pass data["situation"] = _sit data["position_advice"] = r[14] or "" # 持仓状态(2026-07-22 老爸要求:LLM 必须知道是否持有/成本/股数) hr = conn.execute("SELECT shares, cost, price FROM holdings WHERE code=? AND is_active=1 AND shares>0", (code,)).fetchone() if hr and hr[0]: data["held"] = True data["held_shares"] = hr[0] data["held_cost"] = hr[1] or 0 else: data["held"] = False # 2026-08-18 持仓处置决策:策略状态感知(invalidated/none → position_management) try: _st_r = conn.execute("SELECT strategy_state, strategy_attributed, strategy_provenance FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone() data["strategy_state"] = _st_r[0] or "active" if _st_r else "active" data["strategy_attributed"] = _st_r[1] or "" if _st_r else "" # 无策略/失效 + 持仓 → 用 position_management 框架 _need_pos_mgmt = (data.get("strategy_state") in ("invalidated", "none") or not data.get("strategy_def")) and data.get("held") if _need_pos_mgmt: _pm = conn.execute("SELECT display_name, summary, entry_logic, exit_logic, review_focus, holding_style, status, retired_reason, superseded_by, regime FROM strategy_defs WHERE strategy_name='position_management'").fetchone() if _pm: data["strategy_def"] = { "display_name": _pm[0], "summary": _pm[1], "entry_logic": _pm[2], "exit_logic": _pm[3], "review_focus": _pm[4], "holding_style": _pm[5], "status": _pm[6], "retired_reason": _pm[7], "superseded_by": _pm[8], "regime": _pm[9], } data["pos_mgmt_triggered"] = True # current_regime_strategies 已移除(旧策略名与strategy_defs不匹配) except Exception: pass conn.close() # 从腾讯API拉最新价和基本面 # 代码前缀:5位=港股(hk),6/9开头=沪(sh),其他=深(sz) _c = str(code) if len(_c) == 5: prefix = "hk" elif _c.startswith(("6", "9")): prefix = "sh" else: prefix = "sz" try: _pconn = sqlite3.connect(DB, timeout=30) # 价格: stock_daily(收盘价,全市场覆盖) 优先, live_prices(盘中) 其次 _sd = _pconn.execute("SELECT close FROM stock_daily WHERE code=? ORDER BY date DESC LIMIT 1", (code,)).fetchone() if _sd and _sd[0]: data["price"] = float(_sd[0]) else: _lp = _pconn.execute("SELECT price FROM live_prices WHERE code=?", (code,)).fetchone() data["price"] = float(_lp[0]) if _lp and _lp[0] else 0 data["change_pct"] = "0" # PE/市值: 从stock_fundamentals读(fundamentals_full_refresh全市场覆盖) _ff = _pconn.execute("SELECT pe, pb, mcap_total FROM stock_fundamentals WHERE code=?", (code,)).fetchone() if _ff: data["pe"] = str(_ff[0]) if _ff[0] else "" data["mcap"] = str(_ff[1]) if _ff[1] else "" # PB暂时不用,但mcap_total是市值 else: data["pe"] = "" data["mcap"] = "" _pconn.close() except: data["price"] = 0 # 行业上下文修正:sector_context 被"大盘上涨比"污染或为空时,用 stock_sectors 的行业名兜底; # 未映射的股票明确标注"行业未映射"(不让大盘指标伪装成行业信息) _sector_ctx = data.get('sector_context', '') or '' if (not _sector_ctx) or _sector_ctx.startswith('大盘上涨比') or len(_sector_ctx) < 4: _resolved = "" try: _sdb = sqlite3.connect(DB, timeout=30) _sr = _sdb.execute("SELECT sector_name FROM stock_sectors WHERE code=? LIMIT 1", (code,)).fetchone() if _sr and _sr[0]: _resolved = f"行业{_sr[0]}" else: # 2026-08-17 修复:stock_sectors 仅898只覆盖不全,回退 stock_sectors_em(5061只) _sr2 = _sdb.execute("SELECT sector FROM stock_sectors_em WHERE code=? LIMIT 1", (code,)).fetchone() if _sr2 and _sr2[0]: _resolved = f"行业{_sr2[0]}" _sdb.close() except Exception: pass _sector_ctx = _resolved if _resolved else "行业未映射(仅大盘环境参考)" data['sector_context'] = _sector_ctx # 大盘 try: conn = sqlite3.connect(DB, timeout=30) mr = conn.execute("SELECT structure FROM macro_context_log ORDER BY id DESC LIMIT 1").fetchone() if mr and mr[0]: s = json.loads(mr[0]) data["macro"] = s.get("description", "大盘震荡") conn.close() except: data["macro"] = "大盘震荡" # ── 确定性技术位(2026-07-22 老爸:系统计算的技术位必须作为客观锚喂给LLM)── # 与技术路径(per_stock)同源: technical_analysis.full_analysis # 强撑/弱撑/枢轴/弱压/强压/有效区间 + 均线,全部确定性计算,非LLM估计。 data["ta"] = {} try: import technical_analysis as ta_mod _ta = ta_mod.full_analysis(code) if _ta and "error" not in _ta: _sr = _ta.get("support_resistance", {}) or {} _mtf = _ta.get("multi_tf", {}) or {} _mas = (_mtf.get("mas") or {}) data["ta"] = { "strong_support": _sr.get("strong_support"), "weak_support": _sr.get("weak_support"), "pivot": _sr.get("pivot"), "weak_resist": _sr.get("weak_resist"), "strong_resist": _sr.get("strong_resist"), "effective_range": _sr.get("effective_range"), "ma5": _mas.get("ma5"), "ma10": _mas.get("ma10"), "ma20": _mas.get("ma20"), "ma60": _mas.get("ma60"), } except Exception as _te: print(f" ⚠️ 技术位计算失败({code}): {_te}", flush=True) # 深套重算(价格已拉取,与 strategy_lifecycle is_deep_loss 同口径 -20%) _cost3 = data.get("held_cost") or 0 _px3 = data.get("price") or 0 if data.get("held") and _cost3 > 0 and _px3 > 0: _pnl3 = (_px3 - _cost3) / _cost3 * 100 data["deep_loss"] = _pnl3 < -20 data["pnl_pct"] = round(_pnl3, 1) else: data["deep_loss"] = False data["pnl_pct"] = None return data def build_prompt(data): """构建LLM prompt,先审阅原策略再结合实时数据输出修改判断+12维矩阵分析""" cash, total = get_portfolio() if not total: cash, total = 241330, 929727 # 兜底(DB读不到时) # 拉取资金流数据 _flow_note = "暂无资金流数据" try: import sqlite3 as _sq, json as _j _db = _sq.connect("/home/hmo/MoFin/data/mofin.db") _fr = _db.execute("SELECT cache_json FROM capital_flow_cache WHERE id=1 ORDER BY updated_at DESC LIMIT 1").fetchone() if _fr and _fr[0]: _fc = _j.loads(_fr[0]) _stocks = _fc.get("stocks", {}) _s = _stocks.get(data['code'], {}) if _s and _s.get("analysis"): _a = _s["analysis"] _net = _a.get("net_flow", 0) _main = _a.get("main_force", 0) _retail = _a.get("retail_flow", 0) _trend = _a.get("trend", "中性") _flow_note = f"净流入{_net:.0f}万 主力{_main:.0f}万 散户{_retail:.0f}万 趋势{_trend}" _db.close() except: pass # 拉取近期消息面(不再要求情绪标签——原始新闻直接喂给12维LLM,由LLM自行判断情绪。 # 2026-07-22:情绪分类器已退役,利好/利空标签停在07-09,强制过滤=自断新闻源) _news_note = "暂无近期消息" _news_items = [] try: import sqlite3 as _sq _db = _sq.connect("/home/hmo/MoFin/data/mofin.db") # ① 个股直接相关(searched_stocks 含本代码 或 行业名匹配),不限情绪标签 _sector_name = "" try: _sr = _db.execute( "SELECT sector_name FROM stock_sectors WHERE code=? LIMIT 1", (data['code'],)).fetchone() _sector_name = _sr[0] if _sr else "" except Exception: pass _nr = _db.execute( "SELECT summary, overall_sentiment, created_at FROM signal_news " "WHERE searched_stocks LIKE ? OR sector LIKE ? " "ORDER BY id DESC LIMIT 3", (f'%{data["code"]}%', f'%{_sector_name}%')).fetchall() _news_items.extend(_nr) # ② 大盘兜底(独立 try,不被①的失败拖累;取最新3条,不限标签) if not _news_items: _nr2 = _db.execute( "SELECT summary, overall_sentiment, created_at FROM signal_news " "ORDER BY id DESC LIMIT 3").fetchall() _news_items.extend(_nr2) _db.close() except: pass if _news_items: def _fmt(r): senti = r[1] if r[1] and r[1] != 'unknown' else '未标注' return f"{r[2][:10]} [{senti}] {r[0][:40]}" _news_note = " | ".join([_fmt(r) for r in _news_items]) # ── 构建【原策略全文】section ── _params_parts = [] if data.get('action'): _params_parts.append(f"当前策略: {data['action']}") if data.get('timing_signal'): _params_parts.append(f"信号: {data['timing_signal']}") if data.get('entry_low') or data.get('entry_high'): _params_parts.append(f"买入区间: {data.get('entry_low',0)}~{data.get('entry_high',0)}") if data.get('stop_loss'): _params_parts.append(f"止损: {data['stop_loss']}") if data.get('take_profit'): _params_parts.append(f"止盈: {data['take_profit']}") if data.get('position_advice'): _params_parts.append(f"仓位: {data['position_advice']}") _params_str = " | ".join(_params_parts) if _params_parts else "无策略参数" # 最近3条变更记录 _changelog_str = "无变更记录" try: _cl_raw = data.get('changelog_json', '') if _cl_raw: _cl = json.loads(_cl_raw) if isinstance(_cl_raw, str) else _cl_raw if isinstance(_cl, list) and _cl: _recent = _cl[-3:] if len(_cl) > 3 else _cl _cl_lines = [] for i, c in enumerate(_recent): _act = c.get('action', c.get('reason', '')) if isinstance(c, dict) else str(c) _ts = c.get('timestamp', '') if isinstance(c, dict) else '' _cl_lines.append(f" {i+1}. {_ts[:16]} {_act[:80]}") if _cl_lines: _changelog_str = "\n".join(_cl_lines) except: pass # 完整分析原文(不截断) _full_analysis = data.get('full_analysis', '') or '' _fa_display = _full_analysis if _full_analysis else '(首次分析,无历史)' # ── 持仓上下文(2026-07-22 老爸要求:LLM 必须知道持有状态,建议不得两头都写)── if data.get('held'): _sh = data.get('held_shares', 0) _cost = data.get('held_cost', 0) _px = data.get('price', 0) or 0 _pnl = ((_px - _cost) / _cost * 100) if _cost else 0 _position_context = (f"⚠️ 我当前【已持有】{data['code']}:{_sh}股,成本{_cost:.2f}元," f"现价{_px}元(盈亏{_pnl:+.1f}%)。你的建议必须基于「已持有」状态给出" f"(加减仓/止损止盈/持有观察),禁止给「未持有者」的建仓建议。") else: _position_context = (f"⚠️ 我当前【未持有】{data['code']}。你的建议必须基于「未持有」状态给出" f"(是否建仓/什么价位建仓/仓位多大),禁止假设我有浮盈、" f"禁止出现「已持仓者」视角的建议。") # ── 行业强度数据 ── _sector_extra = "" if data.get("sector_change_pct") is not None: _sector_extra = f" 行业涨跌={data['sector_change_pct']}%" if data.get("sector_rank") is not None: _sector_extra += f" 行业排名={data['sector_rank']}" _tech_parts = [] if data.get("ta_strong_support"): _tech_parts.append(f"强支撑={data['ta_strong_support']} 弱支撑={data['ta_weak_support']} 枢轴={data['ta_pivot']} 弱压={data['ta_weak_resist']} 强压={data['ta_strong_resist']}") if data.get("ta_ma20"): _ma_info = f"MA5={data.get('ta_ma5','?')} MA10={data.get('ta_ma10','?')} MA20={data.get('ta_ma20','?')} MA60={data.get('ta_ma60','?')}" if data.get("ta_dist_ma20") is not None: _ma_info += f" 距MA20={data['ta_dist_ma20']}%" _tech_parts.append(_ma_info) if data.get("ta_rsi"): _tech_parts.append(f"RSI={data['ta_rsi']}") if data.get("ta_candle"): _tech_parts.append(f"K线形态={data['ta_candle']}") if data.get("mtf_trend_alignment"): _tech_parts.append(f"多周期趋势={data['mtf_trend_alignment']}") if data.get("mtf_daily_trend"): _tech_parts.append(f"日线={data['mtf_daily_trend']}") if data.get("mtf_weekly_trend"): _tech_parts.append(f"周线={data['mtf_weekly_trend']}") if data.get("mtf_monthly_trend"): _tech_parts.append(f"月线={data['mtf_monthly_trend']}") _factor_parts = [] for _fk, _fl in [("factor_mcap_q", "市值分位"), ("factor_pe_q", "PE分位"), ("factor_bias60", "bias60"), ("factor_ret5d", "5日涨幅"), ("factor_dist_lo20", "距20日低点"), ("factor_dist_ma20", "距MA20"), ("factor_vol_ratio", "量比")]: if data.get(_fk) is not None: _factor_parts.append(f"{_fl}={data[_fk]}") if _factor_parts: _tech_parts.append("基本面分位: " + " ".join(_factor_parts)) _tech_str = " | ".join(_tech_parts) if _tech_parts else "技术指标数据待刷新" # ── 换仓上下文(2026-07-24 老爸:现金不足时给出具体换股建议)── _rotation_context = "" if not data.get('held'): try: _rc = sqlite3.connect(DB, timeout=30) _all_holdings = _rc.execute(""" SELECT hs.code, hs.name, hs.timing_signal, hs.rr_ratio, h.position_pct, h.cost, lp.price, lp.change_pct, hs.full_analysis FROM holding_strategies hs JOIN holdings h ON hs.code = h.code AND h.is_active = 1 LEFT JOIN live_prices lp ON hs.code = lp.code WHERE hs.status='active' AND h.shares > 0 AND hs.timing_signal IN ('弱势持有','观望','持有') ORDER BY h.position_pct DESC""").fetchall() _rc.close() if _all_holdings: _hl = ";".join(f"{w[1]}({w[0]}){w[2]} RR={w[3]:.1f} 仓位{w[4]:.1f}% 浮盈{w[7]:+.1f}%" for w in _all_holdings) _rotation_context = ( f"\n【换仓分析】当前现金{cash:.0f}元不足以建仓,需换仓。" f"以下持仓可供减持(代码/信号/RR/仓位%/浮盈):{_hl}。" f"\n请基于你对该股的12维分析结论,评估以上每只持仓的持有价值," f"推荐最应该卖出的一只(代码+理由+建议减持金额),在SIGNAL_JSON的rotation_candidate字段中明确输出。" f"判断标准:因子评分低、技术面最弱、信号最差、浮亏最大、仓位可腾出的优先卖出。") _position_context += _rotation_context # 2026-08-18 换仓决策注入(老莫:需要资金时对比预期收益,卖E_hold最低的) # 触发:本票信号为买入/加仓 + 现金不足(cash < 建议仓位金额估算) try: from swap_decision import decide_swap, format_swap_advice, get_strategy_expected _sig_now = data.get("timing_signal") or "" _need_fund = _sig_now in ("买入", "可买入", "可加仓") if _need_fund: # 估算本票建议仓位金额(按 RR 5%-15%,取中 10%) _est_need = total * 0.10 if cash < _est_need: # 拉当前持仓(深套判定:cost vs price) try: import sqlite3 as _sq3 _c3 = _sq3.connect(DB, timeout=30) _c3.row_factory = _sq3.Row _holds = _c3.execute("SELECT code, name, shares, cost, price FROM holdings WHERE is_active=1 AND shares>0").fetchall() _c3.close() _hold_list = [] for _h in _holds: _cost_v = _h["cost"] or 0; _price_v = _h["price"] or 0 _pnl_v = (_price_v - _cost_v) / _cost_v * 100 if _cost_v > 0 and _price_v > 0 else None # 查该持仓的 strategy_name _sn = "" try: _c4 = _sq3.connect(DB, timeout=30) _sn_r = _c4.execute("SELECT strategy_name FROM holding_strategies WHERE code=? AND status='active'", (_h["code"],)).fetchone() _c4.close() _sn = _sn_r[0] if _sn_r else "" except Exception: pass _hold_list.append({ "code": _h["code"], "name": _h["name"] or _h["code"], "cost": _cost_v, "price": _price_v, "shares": _h["shares"] or 0, "strategy_name": _sn, "pnl_pct": _pnl_v, }) _sd = data.get("strategy_name") or data.get("strategy_attributed") or "" _dec = decide_swap(need_cash=_est_need - cash, holdings=_hold_list, new_strategy=_sd, market="a") _swap_advice = format_swap_advice(_dec) _position_context += f"\n\n{_swap_advice}\n⚠️ 若你给出买入建议但现金不足,以上为系统算好的换仓方案(卖预期收益最低的持仓凑钱);你只需确认是否采纳(可结合消息面/基本面修正),不要重新设计换仓逻辑。" print(f" [SWAP] 换仓决策已注入: {_dec.get('reason', '')[:80]}", flush=True) except Exception as _se: print(f" [SWAP] 换仓决策注入失败: {_se}", flush=True) except Exception: pass # ── 策略语义注入(2026-08-18 老莫:重评按策略定义,不是裸标签)── _sd = data.get("strategy_def") _sit = data.get("situation") or {} if _sd: _retired_note = "" if _sd.get("status") == "retired": _retired_note = (f"\n⚠️ 该策略已被系统标记【已淘汰】:{_sd.get('retired_reason') or ''}" f" 不要沿用其入场/出场参数作为默认锚,按当前市场环境重新评估。") _strategy_sec = f"""【策略定义】(系统策略库 strategy_defs 提供,非LLM生成) 策略名: {_sd.get('display_name')}({data.get('strategy_name') or 'unknown'}) 策略逻辑: {_sd.get('summary') or ''} 入场逻辑: {_sd.get('entry_logic') or ''} 出场规则: {_sd.get('exit_logic') or ''} 重评侧重: {_sd.get('review_focus') or ''} 持仓风格: {_sd.get('holding_style') or ''}{_retired_note}""" else: _strategy_sec = """【策略定义】(无策略记录) 来源策略未知。请基于技术形态/估值/资金特征判断该股当前最接近的策略画像,并在【策略判断】中说明; 若判断需要归类到某个策略,从以下候选中选择(禁止自创): accumulation(主力建仓) / b_td1_v3(超跌原池优选) / v_mr(弱市深超跌) / hk_pe_mom(港股动量) / hk_pe_oversold(港股超卖) / p_oversold(预测超跌) / s2_panic(恐慌买强势) / leader(龙头回调) / v_next(趋势龙头回调) / v8.1(顺势波段)""" _sit_sec = f"""【情势体检】(系统确定性计算,非LLM估计) 当前温区(A股): {_sit.get('regime_a','unknown')}(7天前: {_sit.get('regime_7d_ago','unknown')}) 温区突变: {'是' if _sit.get('regime_a') != _sit.get('regime_7d_ago') and _sit.get('regime_a') not in ('unknown','') else '否'} 高风险消息: {_sit.get('high_risk') or '无'} 执行红线: 破止损未执行={_sit.get('breach_stop')} 达止盈未执行={_sit.get('reach_tp')} 现价超买入区上沿={_sit.get('out_zone')} 超持有期={_sit.get('over_hold')} ⚠️ 若温区突变/高风险消息/执行红线任一命中,或重评侧重检查发现原入场逻辑不成立,【策略判断】必须考虑「策略失效重定」或「更换策略归类」。""" _orig_strategy_section = f"""{_strategy_sec} {_sit_sec} 来源策略: {data.get("strategy_name") or "unknown"}(按此策略选股逻辑重评,可据最新情况调整参数) 当前策略参数: {_params_str} 变更记录(最近3条): {_changelog_str} 完整分析原文: {_fa_display}""" # ── 技术位锚 section(2026-07-22 老爸:确定性计算值,LLM 必须尊重)── _ta = data.get("ta") or {} if _ta.get("weak_support"): _ma_parts = [f"{k.upper()}={_ta[k]}" for k in ("ma5", "ma10", "ma20", "ma60") if _ta.get(k)] _ma_line = (" ".join(_ma_parts) + "\n") if _ma_parts else "" _ta_sec = f"""【技术位锚】(以下数值由系统基于K线/均线确定性计算,是客观事实,不是你的估计值) 强撑={_ta.get('strong_support')} 弱撑={_ta.get('weak_support')} 枢轴={_ta.get('pivot')} 弱压={_ta.get('weak_resist')} 强压={_ta.get('strong_resist')} 有效区间={_ta.get('effective_range')} {_ma_line}⚠️ 参数锚定纪律(必须遵守,输出前自检): 【趋势判断】(必须给出,影响信号方向) - 上升趋势:MA5>MA10>MA20>MA60,价格在MA20上方 → 回调到支撑位可买入 - 下跌趋势:MA5= 30) or (10 <= _h < 15)) _price_label = "(盘中实时)" if _is_market_open else "(上次收盘/非交易时段)" _macro_label = "(盘中实时)" if _is_market_open else "(最近更新)" return f"""你是一个资深A股分析师。请先审阅以下【原策略全文】,判断是否需要修改策略,然后做出完整的12维矩阵分析。 【原策略全文】 {_orig_strategy_section} ── 以上是已有的策略,以下是当前实时数据,请结合两者做出判断 ── ⚠️ 重要:以下12个维度不是独立分析的,你必须交叉对比后给出综合结论。 例如:如果消息面利好但资金流在流出,说明利好可能是出货;如果基本面强但技术面破位,说明估值可能还没到底。 {_ta_sec} 当前数据(以下数据均来自实时API,每条标注时间窗口,禁止使用模型内部训练数据): 大盘:{data.get('macro','震荡')}{_macro_label} 最新价:{data.get('price',0)} 涨跌:{data.get('change_pct','0')}%{_price_label} PE={data.get('pe','?')}(最新财报) 市值={data.get('mcap','?')}亿 行业:{data.get('sector_context','?')}(近一个交易日){(' 涨跌='+str(data.get('sector_change_pct',''))+'% 排名='+str(data.get('sector_rank',''))) if data.get('sector_change_pct') is not None else ''}{_sector_extra} 技术面:{data.get('tech_snapshot','')[:300]}(MA=5/10/20/60日 支撑阻力=近20日 量价=当日+近5日趋势) {_tech_str} {_tech_str} 资金流:{_flow_note}(近5日累计) 消息面:{_news_note}(最近3条,自动标注抓取时间) 当前信号:{data.get('timing_signal','?')} 分类:{data.get('stock_category','?')} 我的总资产={total}元,可用现金={cash}元。 {_position_context} 请严格按以下格式输出(注意节标题不可省略): 【策略判断】三选一:维持原策略 / 修改参数 / 策略失效需更换 - 默认倾向「维持原策略」;只有情势体检命中红线、或重评侧重检查发现原入场逻辑不成立时,才选「策略失效需更换」 - 「修改参数」= 策略框架仍适用,仅微调参数(止损/止盈/区间) - 「策略失效需更换」= 原策略在当前情势下已不适合(温区切换/黑天鹅/入场逻辑不成立),必须换一个策略来评估 【策略失效理由】(策略判断≠维持原策略时必填,逐条对照重评侧重说明哪条不成立;维持时写"无需修改") 【更换策略归类】策略失效需更换时【必须】写候选列表中的【策略代码名】(如 v_next、v8.1、leader),禁止写描述性文字(如"趋势回调""顺势追涨");候选列表: accumulation / b_td1_v3 / v_mr / hk_pe_mom / hk_pe_oversold / p_oversold / s2_panic / leader / v_next / v8.1 若确实无合适策略,写"无合适策略" 【维持或修改】明确二选一判断:维持原策略 / 需要修改策略 【修改点及理由】 如果维持原策略 → 写"无需修改" 如果需要修改 → 逐条列出(每条格式:"- 修改点名称:理由说明") 【最终新策略】 用自然语言输出完整的最终策略全文(200-400字),自包含核心交易逻辑、买入区间价格、止损价、止盈价、仓位比例、风险提示。 ⚠️ 本段不要使用【综合结论】【买入区间】等标签——用自然语言描述即可。 【交叉分析】用2-3句话说明哪些维度出现矛盾/共振,最关键的信号是什么 ① 大盘×基本面 [一句话,说明矛盾关系] ② 大盘×消息面 [一句话] ③ 大盘×技术面 [一句话] ④ 大盘×资金面 [一句话] ⑤ 行业×基本面 [一句话] ⑥ 行业×消息面 [一句话] ⑦ 行业×技术面 [一句话] ⑧ 行业×资金面 [一句话] ⑨ 个股×基本面 [一句话] ⑩ 个股×消息面 [一句话] ⑪ 个股×技术面 [一句话] ⑫ 个股×资金面 [一句话] 【综合结论】(买入/关注/观望/卖出) 【操作建议】具体操作建议 【买入区间】最低价~最高价(锚定技术位:下沿参考弱撑/强撑,上沿参考枢轴/弱压) 【建议止损】数字(必须严格低于买入区下沿,放弱撑下方1-3%或强撑附近) 【建议止盈】数字(参考弱压/强压) 【参数自检】一行,格式"止损X < 区下沿Y < 区上沿Z < 止盈W:通过/不通过+原因" 【建议仓位】⚠️不可省略。综合结论非"买入"时写"不新建仓";为"买入"时按以下公式: 基础仓位按RR确定:RR<1.5→不推荐,RR1.5~3→8%,RR3~5→12%,RR5+→15% 大盘偏弱×0.8,大盘偏强×1.15 蓝筹/白马×1.2,成长×0.85,题材/短线×0.6 最终仓位范围:5%~20% ⚠️ 仓位硬约束(必须遵守): - 必须输出具体数字%(如"8%(理由)"),禁止"中等仓位/减仓或观望/轻仓/适量"等模糊表述——模糊仓位视为输出作废 - 必须结合可用现金{cash}元计算可买股数,仓位%对应的金额不得超过可用现金 - 若理想仓位超出现金,明确写出:实际可执行仓位X%(受现金限制),并给出换仓建议(减持哪只弱持仓换入本股) 输出格式:"X%(理由:含现金可行性的一句话说明)" ⚠️ 输出纪律(必须遵守): 1. 直接以【维持或修改】开头,禁止任何寒暄、开场白、分隔线 2. 禁止输出 或任何 XML/代码块(唯一例外:末尾的 SIGNAL_JSON 行) 3. 所有【】节标题一个都不能少 4. 止损<区间下沿<区间上沿<止盈,违反任一条=输出作废重想 5. ⚠️ 最后一行必须是机器可读结论(单行,不加代码块标记): SIGNAL_JSON: {{"signal":"买入|可买入|可加仓|关注|观望|持有|弱势持有|卖出|止盈","actionable":true或false,"entry_low":数字或0,"entry_high":数字或0,"stop_loss":数字或0,"take_profit":数字或0,"position_pct":数字0到20,"reason":"一句话","rotation_candidate":{{"code":"代码或空","reason":"若现金不足需换仓,这里写最应卖出的持仓代码及理由,不需换仓则为空"}}}} - signal 必须与【综合结论】完全一致 - actionable=true 仅当"当前价格下立即可执行的买入/卖出操作";观望/等待/条件未满足/不符合建仓条件一律 false - 数字字段与【买入区间】【建议止损】【建议止盈】【建议仓位】一致,无则0""" def parse_response(text): """从LLM回复中提取策略参数。 ⚠️ 节标题精确匹配:只认行首【买入区间】【综合结论】等节行。 绝不用"包含关键词的第一行"——修改点段落会引用旧脏值(如"原买入区间95.0~99.0"), 曾导致脏数据被反复写回(17只股票背着95~99区间,LLM新区间形同虚设)。""" result = {"signal": "", "entry_low": 0, "entry_high": 0, "stop_loss": 0, "take_profit": 0, "position": "", "zone_cleared": False, "action_advice": ""} def _section_line(name): """匹配节标题行:行首(可含空白)【名称】,返回该行内容""" for l in text.split("\n"): if re.match(r'^\s*【' + name + r'】', l): return l return "" # 信号(只认【综合结论】节行,且锚定】后的首个词,防"观望(不建议买入)"误判为买入) sl = _section_line("综合结论") if sl: m = re.search(r'综合结论】\s*[((]?\s*(弱势持有|可加仓|可买入|买入|卖出|止盈|关注|观望|持有)', sl) if m: result["signal"] = m.group(1) # 买入区间(只认【买入区间】节行;"无"→显式清空,不保留旧值) zl = _section_line("买入区间") if zl: if re.search(r'】\s*(无|不设|不参与|空仓)', zl): result["zone_cleared"] = True elif len(re.findall(r'\d+\.?\d*', zl)) >= 2 and all(float(n) <= 0.01 for n in re.findall(r'\d+\.?\d*', zl)[:2]): # 2026-08-18 修复:LLM 输出"买入区0.0~0.0"也视为 zone_cleared(清空脏值) result["zone_cleared"] = True else: # 2026-08-18 修复:LLM 说"95~99(异常,应为14.87~15.01)"时,取"应为"后的正确区间 # 否则取前两个数字(避免旧脏值被当推荐值写回,老莫 600262 教训) _corrected = re.search(r'(?:应为|修正为|改为|更正为)\s*[((]?([\d.]+)\s*[~~-]\s*([\d.]+)', zl) if _corrected: result["entry_low"] = float(_corrected.group(1)) result["entry_high"] = float(_corrected.group(2)) else: nums = re.findall(r'\d+\.?\d*', zl) if len(nums) >= 2: a, b = float(nums[0]), float(nums[1]) result["entry_low"] = min(a, b) result["entry_high"] = max(a, b) # 止损(只认【建议止损】节行) for name in ("建议止损", "止损"): l = _section_line(name) if l: nums = re.findall(r'\d+\.?\d*', l) if nums: result["stop_loss"] = float(nums[0]) break # 止盈(只认【建议止盈】节行) for name in ("建议止盈", "止盈"): l = _section_line(name) if l: nums = re.findall(r'\d+\.?\d*', l) if nums: result["take_profit"] = float(nums[0]) break # 操作建议(只认【操作建议】节行)→ action 字段,前端"当前操作策略"列的唯一新鲜来源 al = _section_line("操作建议") if al: result["action_advice"] = re.sub(r'^\s*【操作建议】\s*', '', al).strip()[:200] # 仓位:只有买入信号才需要,提取百分比数字(只认【建议仓位】节行) result["position"] = "" if result["signal"] == "买入": l = _section_line("建议仓位") if l: nums = re.findall(r'\d+\.?\d*', l) for n in nums: f = float(n) if 1 <= f <= 30: # 合理的仓位范围 result["position"] = f"{f:.0f}%" break # 2026-08-18 策略判断(只认【策略判断】节行) # 老莫 2026-08-18:失效=需更换(同一决策),失效时必须给新归属;给不出才停在失效未归属 result["strategy_judge"] = "" _sj = _section_line("策略判断") if _sj: _m = re.search(r'策略判断】\s*(维持原策略|修改参数|策略失效需更换|策略失效重定)', _sj) if _m: result["strategy_judge"] = _m.group(1) # 策略失效理由 result["strategy_reason"] = "" _sr = _section_line("策略失效理由") if _sr: result["strategy_reason"] = re.sub(r'^\s*【策略失效理由】\s*', '', _sr).strip()[:200] # 更换策略归类(封闭集校验:必须在候选列表内;失效和更换都解析) result["strategy_switch_to"] = "" _CANDIDATES = ("accumulation", "b_td1_v3", "v_mr", "hk_pe_mom", "hk_pe_oversold", "p_oversold", "s2_panic", "leader", "v_next", "v8.1") _st = _section_line("更换策略归类") if _st: _body = re.sub(r'^\s*【更换策略归类】\s*', '', _st).strip().lower() if "无合适策略" in _body or "无合适" in _body: result["strategy_switch_to"] = "" # 明确无合适 → 失效未归属 else: for c in _CANDIDATES: if c in _body: result["strategy_switch_to"] = c break # ── 2026-08-25 机器可读结论优先(老莫:结构化标注,不靠关键词匹配散文)── _m_json = re.search(r'^SIGNAL_JSON:\s*(\{.*\})\s*$', text, re.M) if _m_json: try: import json as _js _sj = _js.loads(_m_json.group(1)) result["signal_json"] = _sj if _sj.get("signal"): result["signal"] = _sj["signal"] for _k, _rk in (("entry_low","entry_low"),("entry_high","entry_high"), ("stop_loss","stop_loss"),("take_profit","take_profit")): if _sj.get(_k): result[_rk] = float(_sj[_k]) if _sj.get("position_pct") is not None: result["position"] = f"{_sj['position_pct']}%(SIGNAL_JSON)" if result["signal"] in ("买入","可买入","可加仓") and _sj.get("actionable") is False: result["signal"] = "关注" result["actionable_downgraded"] = True except Exception: pass return result def save_result(code, full_text, parsed, ta_levels=None): """保存LLM结果到DB(先快照再UPDATE)。空分析拒绝写入。 ta_levels: collect_data 计算的确定性技术位,用于止损锚定校验。""" if not (full_text or "").strip(): print(f" \u274c 拒绝写入空分析(LLM输出为空,保护已有数据)") return conn = sqlite3.connect(DB, timeout=30) conn.execute("PRAGMA busy_timeout=30000") now = datetime.now().isoformat() # ── 修改前快照 ── snapshot_strategy_history(conn, code, 'batch_12d') updates = ["full_analysis=?", "reassessed_at=?"] params = [full_text, now] if parsed["signal"]: updates.append("timing_signal=?") params.append(parsed["signal"]) # 区间写入门禁:上下沿都必须为正且 下沿<上沿<下沿x3,否则视为解析错误整体跳过 # (防 214.68~2.52 类解析污染,与 GATE_ZONE_SANITY 同级防护) _el, _eh = parsed["entry_low"], parsed["entry_high"] if parsed.get("zone_cleared"): # LLM 显式输出【买入区间】无 → 清空区间(不再保留可能脏的旧值) updates.append("entry_low=0") updates.append("entry_high=0") elif _el > 0 and _eh > _el and _eh < _el * 3: # 区间-现价距离门禁:整体偏离现价过远(区上沿<现价0.5x 或 区下沿>现价1.5x) # → 判定脏数据/解析错误,拒写并清空(不再"保留原值"养脏,如95~99 vs 现价60) _px = 0.0 try: _pr = conn.execute("SELECT price FROM live_prices WHERE code=?", (code,)).fetchone() _px = float(_pr[0]) if _pr and _pr[0] else 0.0 except Exception: pass if _px > 0 and (_eh < _px * 0.5 or _el > _px * 1.5): print(f" ⚠️ 买入区{_el}~{_eh}偏离现价{_px}过远,拒写并清空(防脏数据残留)", flush=True) updates.append("entry_low=0") updates.append("entry_high=0") else: updates.append("entry_low=?") params.append(_el) updates.append("entry_high=?") params.append(_eh) elif _el > 0 or _eh > 0: print(f" ⚠️ 买入区解析异常({_el}~{_eh}),跳过区间写入(保留原值)", flush=True) # 止损/止盈一致性门禁:损>0 时必须在区间下沿之下(0.5x~1.0x),盈>0 时必须在区间上沿之上 _sl, _tp = parsed["stop_loss"], parsed["take_profit"] # ── 止损锚定门禁(2026-07-22 老爸:止损=区间下沿=把止损设在买入价上,RR恒0)── # LLM 输出 sl>=el 时:用确定性技术位自动修正(弱撑×0.985),无技术位则拒写止损。 if _sl > 0 and _el > 0 and _sl >= _el: _ws = (ta_levels or {}).get("weak_support") or 0 _ss = (ta_levels or {}).get("strong_support") or 0 if _ws > 0 and _ws < _el: _fixed = round(_ws * 0.985, 2) print(f" ⚠️ 止损{_sl}≥区下沿{_el},按技术锚修正为 弱撑{_ws}×0.985={_fixed}", flush=True) _sl = _fixed elif _ss > 0 and _ss < _el: _fixed = round(_ss * 0.99, 2) print(f" ⚠️ 止损{_sl}≥区下沿{_el},按技术锚修正为 强撑{_ss}×0.99={_fixed}", flush=True) _sl = _fixed else: print(f" ⚠️ 止损{_sl}≥区下沿{_el}且无可用技术位,拒写止损(保留原值)", flush=True) _sl = 0 if _sl > 0 and (not _el or _sl < _el) and (not _tp or _sl < _tp): updates.append("stop_loss=?") params.append(_sl) elif _sl > 0: print(f" ⚠️ 止损{_sl}与区间/止盈不一致,跳过写入(保留原值)", flush=True) if _tp > 0 and (not _eh or _tp > _eh) and (not _sl or _tp > _sl): updates.append("take_profit=?") params.append(_tp) elif _tp > 0: print(f" ⚠️ 止盈{_tp}与区间/止损不一致,跳过写入(保留原值)", flush=True) # 2026-08-18 修复:LLM 算的 entry 区间也写回 entry_low/entry_high(纠正 promote 错误区间) # 例:promote 入库 entry 95~99 vs 现价 7.59 严重偏离,重评算正确 7.55~7.70 应覆盖 _el_new = parsed.get("entry_low") or 0 _eh_new = parsed.get("entry_high") or 0 if _el_new > 0 and _eh_new > _el_new: # 校验 entry 与现 stop_loss/take_profit 一致(sl < el < eh < tp) _ok = True if "stop_loss=?" in updates: _cur_sl = params[updates.index("stop_loss=?")] if _cur_sl >= _el_new: _ok = False if "take_profit=?" in updates: _cur_tp = params[updates.index("take_profit=?")] if _cur_tp <= _eh_new: _ok = False if _ok: updates.append("entry_low=?") params.append(_el_new) updates.append("entry_high=?") params.append(_eh_new) if parsed["position"]: updates.append("position_advice=?") params.append(parsed["position"]) if parsed.get("action_advice"): # 12维操作建议 → action(前端"当前操作策略"列;防技术路径旧值与分析矛盾) updates.append("action=?") params.append(parsed["action_advice"]) # ── 2026-08-17 同步 trigger_json:price_monitor 的监控阈值来源 ── if updates: _cur = conn.execute( "SELECT stop_loss, take_profit, entry_low, entry_high FROM holding_strategies " "WHERE code=? AND status='active'", (code,)).fetchone() _sl = _tp = _el = _eh = 0 if _cur: _sl, _tp, _el, _eh = _cur if "stop_loss=?" in updates: _sl = params[updates.index("stop_loss=?")] if "take_profit=?" in updates: _tp = params[updates.index("take_profit=?")] import json as _json _trigger = {"stop_loss": _sl, "entry_zone": f"{_el}~{_eh}" if _el and _eh else "", "take_profit_zone": f"0~{_tp}" if _tp else ""} updates.append("trigger_json=?") params.append(_json.dumps(_trigger, ensure_ascii=False)) params.append(code) sql = f"UPDATE holding_strategies SET {', '.join(updates)} WHERE code=? AND status='active'" conn.execute(sql, params) conn.commit() # 2026-08-18 策略状态落库(strategy_judge/switch_to,老莫:失效=需更换,同一决策) _j = parsed.get("strategy_judge", "") _st = parsed.get("strategy_switch_to", "") if _j in ("策略失效需更换", "策略失效重定"): if _st: # 失效+给归属 → 切换到新策略 conn.execute("UPDATE holding_strategies SET strategy_attributed=?, strategy_state='switched', strategy_provenance='llm_attributed' WHERE code=? AND status='active'", (_st, code)) conn.commit() print(f" ✅ 策略切换: {_j} → 归属 {_st}(strategy_state=switched)") else: # 失效+无合适归属 → 标记失效 conn.execute("UPDATE holding_strategies SET strategy_state='invalidated' WHERE code=? AND status='active'", (code,)) conn.commit() print(f" ✅ 策略失效: {_j} 无合适归属(strategy_state=invalidated)") # ── 信号以分析为唯一事实源(防信号/分析脱节)── from mofin_db import reconcile_signal_from_analysis final_sig = reconcile_signal_from_analysis(conn, code) # ── 推荐操作 tag 同步(跟随对齐后的信号)── sync_recommend_tag(conn, code, final_sig) # 买入信号推送已统一收拢到 sync_recommend_tag 的转场推送(防双重告警)。 # 本路径只负责写库+tag,推送由 mofin_db.push_recommend_alert 在 tag 转场时触发。 conn.close() def process_stock(code, force_today=False): """处理单只股票""" print(f"\n{'='*50}") print(f"处理: {code}") print(f"{'='*50}") if in_cooldown(code): print(f" \u23ed 冷却期内,跳过") return False # 有分析且未过期 \u2192 跳过(除非 force_today 且今早未评) if not FORCE_REASSESS and has_llm_analysis(code) and not analysis_stale(code, force_today): print(f" \u23ed 已有12维分析且未过期,跳过") return False print(f" 收集数据...", flush=True) data = collect_data(code) if not data.get("price"): print(f" \u26a0\ufe0f 无价格数据,跳过") return False print(f" 调LLM生成12维分析...", flush=True) prompt = build_prompt(data) # ── 使用共享 LLM 客户端(替代 curl subprocess)── result = call_llm(prompt, model=REASSESS_MODEL, max_tokens=None, concurrent=True) # 2026-08-24 并发模式: router round-robin分key(4worker不再全压同一key) if not result["ok"] or not (result.get("content") or "").strip(): print(f" \u274c LLM调用失败或空输出: {result.get('error') or 'empty content'}") return False full_text = result["content"] print(f" \u2705 LLM返回({len(full_text)}字, {result['elapsed']:.1f}s, 尝试{result['attempts']}次)", flush=True) parsed = parse_response(full_text) # ── 截断保护:输出过短且无信号 = 低质输出,升级 pro 重试一次 ── if not parsed.get("signal") and len(full_text) < 1500: print(f" ⚠️ 输出截断({len(full_text)}字)且无信号,升级 {FALLBACK_MODEL} 重试...", flush=True) result2 = call_llm(prompt, model=FALLBACK_MODEL, max_tokens=None, concurrent=True) # 2026-08-24 并发模式 if result2["ok"] and len((result2.get("content") or "").strip()) > len(full_text): full_text = result2["content"] parsed = parse_response(full_text) print(f" \u2705 升级后({len(full_text)}字)", flush=True) print(f" 信号={parsed['signal']} 区间={parsed['entry_low']}~{parsed['entry_high']} 损={parsed['stop_loss']} 盈={parsed['take_profit']} 仓位={parsed['position']}") save_result(code, full_text, parsed, ta_levels=data.get("ta")) print(f" \u2705 已保存到DB") return True def main(): # ── 双通道预检:OCG直连 + hermes gateway 兜底,全挂才退出 ── _ocg_ok = ocg_alive() _gw_ok = gateway_alive() if not _ocg_ok and not _gw_ok: print("[FATAL] OCG上游与hermes gateway均不可用,退出") sys.exit(1) if not _ocg_ok: print("[WARN] OCG直连不可用,将使用gateway兜底(agent运行时,较慢)") if not _gw_ok: print("[WARN] hermes gateway不可用,仅使用OCG直连") codes = [] force_today = "--today" in sys.argv dtype = None if "--type" in sys.argv: idx = sys.argv.index("--type") dtype = sys.argv[idx + 1] # holding | watchlist | all if "--code" in sys.argv: idx = sys.argv.index("--code") codes = [sys.argv[idx+1]] else: # 按类型筛选 active 策略 type_map = {"holding": "持仓策略", "watchlist": "自选策略"} conn = sqlite3.connect(DB, timeout=30) if dtype in type_map: rows = conn.execute( "SELECT code FROM holding_strategies WHERE status='active' AND decision_type=? ORDER BY code", (type_map[dtype],)).fetchall() else: rows = conn.execute( "SELECT code FROM holding_strategies WHERE status='active' ORDER BY decision_type, code").fetchall() conn.close() codes = [r[0] for r in rows] # ── 分片并发(2026-07-24 老爸:这么多key不能并发?)── # --shard K/N:本 worker 只处理 index%N==K 的股票,N 个进程并发互不重叠。 _shard_k, _shard_n = 0, 1 if "--shard" in sys.argv: _sk = sys.argv[sys.argv.index("--shard") + 1] # 格式 K/N _shard_k, _shard_n = int(_sk.split("/")[0]), int(_sk.split("/")[1]) if _shard_n > 1: codes = [c for i, c in enumerate(codes) if i % _shard_n == _shard_k] print(f"待处理: {len(codes)}只 (type={dtype or 'all'}, force_today={force_today}" + (f", shard={_shard_k}/{_shard_n}" if _shard_n > 1 else "") + ")") ok = 0 fail = 0 skip = 0 failed_codes = [] for i, code in enumerate(codes): if not FORCE_REASSESS and has_llm_analysis(code) and not analysis_stale(code, force_today): print(f" [{i+1}/{len(codes)}] \u23ed {code} 已有12维分析且未过期") skip += 1 continue print(f" [{i+1}/{len(codes)}] ", end="", flush=True) if process_stock(code, force_today): ok += 1 else: fail += 1 failed_codes.append(code) # 间隔8秒(pro model较重但gateway可承受;retry逻辑吸收瞬断) if i < len(codes) - 1: print(f" 等待8秒...", flush=True) time.sleep(8) # ── 失败二轮:主跑结束后休息 60s 让上游恢复,失败股整体重试一次 ── # (凌晨上游空输出高发,二轮可救回大半;仍失败的留给下一轮调度) if failed_codes: print(f"\n{'='*50}") print(f"失败二轮: {len(failed_codes)}只,休息60s后重试...") time.sleep(60) retry_ok = 0 for code in failed_codes: print(f" [retry] {code} ", end="", flush=True) if process_stock(code, force_today): retry_ok += 1 ok += 1 fail -= 1 print(f" 等待8秒...", flush=True) time.sleep(8) print(f"失败二轮: {retry_ok}/{len(failed_codes)} 救回") print(f"\n{'='*50}") print(f"完成: {ok}成功, {fail}失败, {skip}跳过") print(f"{'='*50}") # ── 推荐摘要:本轮新增推荐聚成一条推送(防逐只轰炸)── # 并发分片模式(SKIP_FLUSH=1)下由 launcher 统一 flush,避免先到者发半成品摘要 if not os.environ.get("SKIP_FLUSH"): try: from mofin_db import flush_rec_digest flush_rec_digest() except Exception as _e: print(f" ⚠️ 推荐摘要发送失败: {_e}") if __name__ == "__main__": main()