From 050c397e17db517090524254fbc60c2f3801ae9d Mon Sep 17 00:00:00 2001 From: xxm Date: Wed, 26 Aug 2026 13:55:10 +0800 Subject: [PATCH] feat-rotation-llm-driven --- deploy/profile-scripts/batch_reassess.py | 1151 ---------------------- deploy/profile-scripts/mofin_db.py | 52 +- 2 files changed, 18 insertions(+), 1185 deletions(-) diff --git a/deploy/profile-scripts/batch_reassess.py b/deploy/profile-scripts/batch_reassess.py index d6596f1b..e69de29b 100644 --- a/deploy/profile-scripts/batch_reassess.py +++ b/deploy/profile-scripts/batch_reassess.py @@ -1,1151 +0,0 @@ -#!/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) - _weak = _rc.execute(""" - SELECT hs.code, hs.name, hs.timing_signal, h.position_pct, h.cost - FROM holding_strategies hs - JOIN holdings h ON hs.code = h.code AND h.is_active = 1 - WHERE hs.status='active' AND h.shares > 0 - AND hs.timing_signal IN ('弱势持有','观望','持有') - ORDER BY CASE hs.timing_signal WHEN '弱势持有' THEN 0 WHEN '观望' THEN 1 ELSE 2 END, - h.position_pct DESC LIMIT 3""").fetchall() - _rc.close() - if _weak: - _wl = ";".join(f"{w[1]}({w[0]}){w[2]}仓位{w[3]:.1f}%" for w in _weak) - _rotation_context = (f"\n我的最弱持仓(可减换仓候选):{_wl}。" - f"若你认为{data['code']}比它们更值得持有,在【操作建议】末尾明确写" - f"「换仓建议:减持XX换入本股」。") - except Exception: - pass - _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":"一句话"}} -- 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() diff --git a/deploy/profile-scripts/mofin_db.py b/deploy/profile-scripts/mofin_db.py index 606f610b..09e6ddf0 100644 --- a/deploy/profile-scripts/mofin_db.py +++ b/deploy/profile-scripts/mofin_db.py @@ -1876,40 +1876,24 @@ def flush_rec_digest(max_items=5): + (f"(合计≈{cum:.0f}%)" if buys else "")) # ── 换仓策略:有排队推荐时,找可减的弱持仓来腾挪 ── if queued: - weak = conn.execute(""" - SELECT hs.code, hs.name, hs.timing_signal, h.position_pct, h.cost, lp.price, lp.change_pct - 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 ('弱势持有','观望','持有') - """).fetchall() - # ── v7.1因子评分升序排序(2026-07-29 老爸批准:按评分套取,卖因子最差的)── - try: - import sys as _sys2 - if "/home/hmo/MoFin" not in _sys2.path: - _sys2.path.insert(0, "/home/hmo/MoFin") - from backtest_framework import prepare_bars as _pb, compute_single_score as _cs - from datetime import datetime as _dt2, timedelta as _td2 - _end2 = _dt2.now().strftime('%Y-%m-%d') - _start2 = (_dt2.now() - _td2(days=150)).strftime('%Y-%m-%d') - _scored = [] - for w in weak: - _sc = 0 - try: - _bars = _pb(w['code'], _start2, _end2) - if _bars and len(_bars) >= 25: - _r = _cs(_bars) - _sc = _r[0] if _r else 0 - except Exception: - pass - _scored.append((_sc, w)) - _scored.sort(key=lambda x: x[0]) # 评分最低 = 优先套取 - weak = [w for _, w in _scored] - print(" [换仓] 因子评分排序: " + ", ".join(f"{w['name']}({s})" for s, w in _scored[:5]), flush=True) - except Exception as _se: - print(f" [换仓] 评分排序失败(回退信号排序): {_se}", flush=True) - weak = sorted(weak, key=lambda w: ({'弱势持有': 0, '观望': 1}.get(w['timing_signal'], 2), + weak = [] + # 2026-08-26: LLM rotation + for item in items: + sj = item.get('signal_json') or {} + rc = sj.get('rotation_candidate') or {} + if rc.get('code') and rc.get('reason'): + _rot = conn.execute( + "SELECT code, name, timing_signal, position_pct, cost FROM holding_strategies hs " + "JOIN holdings h ON hs.code = h.code AND h.is_active = 1 " + "WHERE hs.code = ? AND h.shares > 0", (rc['code'],) + ).fetchone() + if _rot: + _d = dict(_rot) + _d['_rotation_reason'] = rc.get('reason', '') + weak.append(_d) + if weak: + print(f" [换仓] LLM推荐 {len(weak)} 只", flush=True) +weak = sorted(weak, key=lambda w: ({'弱势持有': 0, '观望': 1}.get(w['timing_signal'], 2), -(w['position_pct'] or 0))) if weak: need_pct = queued[0][1]