diff --git a/deploy/profile-scripts/batch_reassess.py b/deploy/profile-scripts/batch_reassess.py index e69de29b..83f2a329 100644 --- a/deploy/profile-scripts/batch_reassess.py +++ b/deploy/profile-scripts/batch_reassess.py @@ -0,0 +1,1154 @@ +#!/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()