Files
MoFin/deploy/profile-scripts/batch_reassess.py
T

1119 lines
59 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/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
# ── 技术指标收集(策略输入数据)──
try:
import technical_analysis as _ta
_tech = _ta.full_analysis(code)
if _tech:
_sr = _tech.get("support_resistance", {})
data["ta_strong_support"] = _sr.get("strong_support", 0)
data["ta_weak_support"] = _sr.get("weak_support", 0)
data["ta_pivot"] = _sr.get("pivot", 0)
data["ta_weak_resist"] = _sr.get("weak_resist", 0)
data["ta_strong_resist"] = _sr.get("strong_resist", 0)
_cs = _tech.get("candlestick", {})
data["ta_candle"] = _cs.get("pattern", "") + "/" + _cs.get("sentiment", "")
_vol = _tech.get("volume", {})
data["ta_volume"] = _vol.get("description", "")
import re as _re
_snap = data.get("tech_snapshot", "")
_ma = _re.search(r'MA5=([\d.]+).*?MA10=([\d.]+).*?MA20=([\d.]+).*?MA60=([\d.]+)', _snap)
if _ma:
data["ta_ma5"] = float(_ma.group(1))
data["ta_ma10"] = float(_ma.group(2))
data["ta_ma20"] = float(_ma.group(3))
data["ta_ma60"] = float(_ma.group(4))
if data["ta_ma20"] > 0:
data["ta_dist_ma20"] = round((data["price"] - data["ta_ma20"]) / data["ta_ma20"] * 100, 2)
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:
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, 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_em5061只)
_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}"""
# ── 技术位锚 section2026-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<MA10<MA20<MA60,价格在MA20下方 → 反弹到阻力位应观望,禁止买入
- 下跌趋势中的反弹:价格从低点回升但未突破MA20 → 属于技术性反弹,不是反转,禁止追高买入
- 震荡:MA交织,方向不明 → 观望为主
- 波段出场形态(k线形态判断,2026-08-12 老莫:k线形态由LLM判断,不在代码硬编码):价格连续2日收破MA10(站稳MA10下方)→ 趋势转弱,可波段先出/减仓;价格收回MA10上方且突破前一日高点(回稳+新高)→ 趋势回稳,可波段再进/接回
1. 买入区应落在技术位之间:下沿参考弱撑/强撑附近,上沿参考枢轴/弱压附近
2. 止损必须严格低于买入区下沿——放在弱撑下方1-3%或强撑附近;严禁止损≥区间下沿(等于把止损设在买入价上,下沿买入立即止损,RR恒为0)
3. 止盈应参考弱压/强压,不得明显高于强压
4. 止盈目标优先使用最近阻力位(20日新高/弱压/强压),不得超过最近阻力;若最近阻力使止盈空间压缩50%以上,直接降级为观望
5. 若按技术位计算的 RR(中值) < 2.0,信号必须降级为「关注」或「观望」,禁止给 RR<2.0 的买入建议
6. 现价已高于买入区上沿 5% 以上时,禁止追高推荐,信号降级为「观望」
7. 【关键】空仓/观望时买入区仍须填合理技术区间(供风报比计算),禁止填 0.0~0.0;
空仓仅改信号/仓位建议,不影响区间参数——区间是参考值,不是"买不买"的判断
4. 若你判断技术位不适用(如突发重大消息/基本面剧变),必须在【修改点及理由】中明确写出偏离理由,禁止静默偏离"""
else:
_ta_sec = "【技术位锚】本次计算不可用,请基于价格行为谨慎给出参数,并仍须满足:止损<区间下沿<区间上沿<止盈。"
# ── 市场状态标签(2026-07-27 老爸:盘前批次数据是昨收,不要骗LLM是"当日实时")──
_now = datetime.now()
_h, _m, _w = _now.hour, _now.minute, _now.weekday()
_is_market_open = _w < 5 and ((_h == 9 and _m >= 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. 禁止输出 <structured_data> 或任何 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
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
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_jsonprice_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) # 文档: 推理模型不指定max_tokens
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) # 文档: 推理模型不指定max_tokens
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()