#!/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) 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 def in_cooldown(code): """冷却期检查""" conn = sqlite3.connect(DB) 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) 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) 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) 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 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["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 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: r = subprocess.run(["curl", "-s", f"http://qt.gtimg.cn/q={prefix}{code}"], capture_output=True, timeout=10) parts = r.stdout.decode("gbk", errors="ignore").split("~") data["price"] = float(parts[3]) if len(parts) > 3 and parts[3] else 0 data["pe"] = parts[39] if len(parts) > 39 and parts[39] else "" data["mcap"] = parts[44] if len(parts) > 44 and parts[44] else "" data["change_pct"] = parts[32] if len(parts) > 32 and parts[32] else "0" 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) _sr = _sdb.execute("SELECT sector_name FROM stock_sectors WHERE code=? LIMIT 1", (code,)).fetchone() _sdb.close() if _sr and _sr[0]: _resolved = f"行业{_sr[0]}" except Exception: pass _sector_ctx = _resolved if _resolved else "行业未映射(仅大盘环境参考)" data['sector_context'] = _sector_ctx # 大盘 try: conn = sqlite3.connect(DB) 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) 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 ORDER BY id 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"禁止出现「已持仓者」视角的建议。") # ── 换仓上下文(2026-07-24 老爸:现金不足时给出具体换股建议)── _rotation_context = "" if not data.get('held'): try: _rc = sqlite3.connect(DB) _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 _orig_strategy_section = f"""当前策略参数: {_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}⚠️ 参数锚定纪律(必须遵守,输出前自检): 1. 买入区应落在技术位之间:下沿参考弱撑/强撑附近,上沿参考枢轴/弱压附近 2. 止损必须严格低于买入区下沿——放在弱撑下方1-3%或强撑附近;严禁止损≥区间下沿(等于把止损设在买入价上,下沿买入立即止损,RR恒为0) 3. 止盈应参考弱压/强压,不得明显高于强压 4. 若你判断技术位不适用(如突发重大消息/基本面剧变),必须在【修改点及理由】中明确写出偏离理由,禁止静默偏离""" else: _ta_sec = "【技术位锚】本次计算不可用,请基于价格行为谨慎给出参数,并仍须满足:止损<区间下沿<区间上沿<止盈。" return f"""你是一个资深A股分析师。请先审阅以下【原策略全文】,判断是否需要修改策略,然后做出完整的12维矩阵分析。 【原策略全文】 {_orig_strategy_section} ── 以上是已有的策略,以下是当前实时数据,请结合两者做出判断 ── ⚠️ 重要:以下12个维度不是独立分析的,你必须交叉对比后给出综合结论。 例如:如果消息面利好但资金流在流出,说明利好可能是出货;如果基本面强但技术面破位,说明估值可能还没到底。 {_ta_sec} 当前数据(以下数据均来自实时API,每条标注时间窗口,禁止使用模型内部训练数据): 大盘:{data.get('macro','震荡')}(当日实时) 最新价:{data.get('price',0)} 涨跌:{data.get('change_pct','0')}%(当日实时) PE={data.get('pe','?')}(最新财报) 市值={data.get('mcap','?')}亿 行业:{data.get('sector_context','?')}(当日实时) 技术面:{data.get('tech_snapshot','')[:300]}(MA=5/10/20/60日 支撑阻力=近20日 量价=当日+近5日趋势) 资金流:{_flow_note}(近5日累计) 消息面:{_news_note}(最近3条,自动标注抓取时间) 当前信号:{data.get('timing_signal','?')} 分类:{data.get('stock_category','?')} 我的总资产={total}元,可用现金={cash}元。 {_position_context} 请严格按以下格式输出(注意节标题不可省略): 【维持或修改】明确二选一判断:维持原策略 / 需要修改策略 【修改点及理由】 如果维持原策略 → 写"无需修改" 如果需要修改 → 逐条列出(每条格式:"- 修改点名称:理由说明") 【最终新策略】 用自然语言输出完整的最终策略全文(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/JSON/代码块 3. 所有【】节标题一个都不能少 4. 止损<区间下沿<区间上沿<止盈,违反任一条=输出作废重想""" 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 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 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) 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) 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"]) params.append(code) sql = f"UPDATE holding_strategies SET {', '.join(updates)} WHERE code=? AND status='active'" conn.execute(sql, params) conn.commit() # ── 信号以分析为唯一事实源(防信号/分析脱节)── 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 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=4096) 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=4096) 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) 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 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()