1118 lines
58 KiB
Python
1118 lines
58 KiB
Python
#!/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, bias20, rsi, r5f, dist_lo20 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
|
||
data["factor_bias20"] = _fi[3] if _fi[3] is not None else None
|
||
if data.get("ta_rsi") is None and _fi[4] is not None:
|
||
data["ta_rsi"] = round(_fi[4], 1)
|
||
data["factor_ret5d"] = _fi[5] if _fi[5] is not None else None
|
||
data["factor_dist_lo20"] = _fi[6] if _fi[6] 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_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_bias20", "bias20"), ("factor_ret5d", "5日涨幅"), ("factor_dist_lo20", "距20日低点")]:
|
||
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<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_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) # 文档: 推理模型不指定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()
|