重评核心重构:ds-v4-pro + 原策略全文 + strategy_history + 前端三改造
后端(重评管线):
- 新增 llm_client.py 共享客户端: REASSESS_MODEL=deepseek-v4-pro 单点,
gateway预检(fail-fast), 150s超时+1次重试, 永不抛异常
- batch_reassess/per_stock_reassess: curl/urllib -> call_llm,
prompt传入原策略全文+当前参数+最近3条变更, 输出 维持/修改判断+
修改点理由+最终新策略, max_tokens 4096
- mofin_db: 新增 strategy_history 表 + snapshot_strategy_history(),
write_holding_strategy 覆写前自动快照(保留20条/code)
- mofin_db: holding_strategies 补 tag 列迁移 + 写入保留
(tag缺席=保留旧值, 显式传''=允许清除), 修复推荐标签被静默丢弃
- mo_data.read_decisions: SELECT 补 tag
- stale_detector/promote_candidates: 子进程超时 240/60 -> 480s
前端:
- 移除 报告Tab -> mofin_health 全部流程/Cron 表加 最后十次 列
(modal列表->详情), /api/reports 支持 cron+script 多路匹配
(jobs.json name->id 解析 + 文件名/标题子串兜底)
- 移除 决策库Tab
- 盯盘Tab 重构: 全部持仓+自选, sort_group 分组(推荐/持仓/自选),
推荐行琥珀高亮+🔥badge+行内策略, 新增 操作策略 列查看
最近3次完整策略(/api/strategy_history/<code>, 表缺失时降级当前行)
- 提示词Tab: registry.py 数据路径改回 /home/hmo/MoFin/data/prompts
(红线: 数据只在规范数据根), 空态提示初始化命令
This commit is contained in:
@@ -1,19 +1,30 @@
|
||||
#!/usr/bin/env python3
|
||||
"""batch_reassess.py — 批量补全九维分析(逐只处理,间隔防限流)
|
||||
"""batch_reassess.py — 批量补全12维(九维矩阵)LLM分析(逐只处理,间隔防限流)
|
||||
|
||||
用法: python3 batch_reassess.py [--all] [--code XXXXXX]
|
||||
用法:
|
||||
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)写九维分析+策略 → 保存到DB
|
||||
流程:收集最新数据 → 调LLM(gateway)写12维分析+策略 → 保存到DB
|
||||
"""
|
||||
import sys, json, subprocess, sqlite3, re, time
|
||||
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, gateway_alive
|
||||
from mofin_db import snapshot_strategy_history
|
||||
|
||||
DB = "/home/hmo/MoFin/data/mofin.db"
|
||||
GATEWAY = "http://127.0.0.1:8643/v1/chat/completions"
|
||||
COOLDOWN_HOURS = 1
|
||||
STALE_HOURS = 20 # 分析超过20小时视为过期,需要重评
|
||||
|
||||
def has_llm_analysis(code):
|
||||
"""检查是否为LLM生成的九维分析(>500字)"""
|
||||
"""检查是否为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()
|
||||
@@ -33,13 +44,41 @@ def in_cooldown(code):
|
||||
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读策略
|
||||
# 从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 FROM holding_strategies WHERE code=? AND status='active'", (code,)).fetchone()
|
||||
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
|
||||
@@ -52,10 +91,21 @@ def collect_data(code):
|
||||
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 ""
|
||||
conn.close()
|
||||
|
||||
# 从腾讯API拉最新价和基本面
|
||||
prefix = "sh" if str(code).startswith(("6","9")) else "sz"
|
||||
# 代码前缀: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("~")
|
||||
@@ -80,9 +130,10 @@ def collect_data(code):
|
||||
return data
|
||||
|
||||
def build_prompt(data):
|
||||
"""构建LLM prompt,要求输出完整策略"""
|
||||
cash = 321271 # 可用现金(从DB读取)
|
||||
total = 952879 # 总资产
|
||||
"""构建LLM prompt,先审阅原策略再结合实时数据输出修改判断+九维矩阵分析"""
|
||||
cash, total = get_portfolio()
|
||||
if not total:
|
||||
cash, total = 241330, 929727 # 兜底(DB读不到时)
|
||||
|
||||
# 拉取资金流数据
|
||||
_flow_note = "暂无资金流数据"
|
||||
@@ -122,7 +173,53 @@ def build_prompt(data):
|
||||
except:
|
||||
pass
|
||||
|
||||
return f"""你是一个资深A股分析师。请对{data['code']} {data.get('name','')}做一个完整的九维矩阵分析,并输出策略参数。
|
||||
# ── 构建【原策略全文】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 '(首次分析,无历史)'
|
||||
|
||||
_orig_strategy_section = f"""当前策略参数: {_params_str}
|
||||
|
||||
变更记录(最近3条):
|
||||
{_changelog_str}
|
||||
|
||||
完整分析原文:
|
||||
{_fa_display}"""
|
||||
|
||||
return f"""你是一个资深A股分析师。请先审阅以下【原策略全文】,判断是否需要修改策略,然后做出完整的九维矩阵分析。
|
||||
|
||||
【原策略全文】
|
||||
{_orig_strategy_section}
|
||||
|
||||
── 以上是已有的策略,以下是当前实时数据,请结合两者做出判断 ──
|
||||
|
||||
⚠️ 重要:以下9个维度不是独立分析的,你必须交叉对比后给出综合结论。
|
||||
例如:如果消息面利好但资金流在流出,说明利好可能是出货;如果基本面强但技术面破位,说明估值可能还没到底。
|
||||
@@ -136,11 +233,18 @@ PE={data.get('pe','?')}(最新财报) 市值={data.get('mcap','?')}亿
|
||||
资金流:{_flow_note}(近5日累计)
|
||||
消息面:{_news_note}(最近3条,自动标注抓取时间)
|
||||
当前信号:{data.get('timing_signal','?')} 分类:{data.get('stock_category','?')}
|
||||
原策略:{(data.get('action','') or '')[:200]}
|
||||
|
||||
我的总资产={total}元,可用现金={cash}元。
|
||||
|
||||
请严格按以下格式输出:
|
||||
请严格按以下格式输出(注意节标题不可省略):
|
||||
|
||||
【维持或修改】明确二选一判断:维持原策略 / 需要修改策略
|
||||
【修改点及理由】
|
||||
如果维持原策略 → 写"无需修改"
|
||||
如果需要修改 → 逐条列出(每条格式:"- 修改点名称:理由说明")
|
||||
【最终新策略】
|
||||
用自然语言输出完整的最终策略全文(200-400字),自包含核心交易逻辑、买入区间价格、止损价、止盈价、仓位比例、风险提示。
|
||||
⚠️ 本段不要使用【综合结论】【买入区间】等标签——用自然语言描述即可。
|
||||
|
||||
【交叉分析】用2-3句话说明哪些维度出现矛盾/共振,最关键的信号是什么
|
||||
① 大盘×基本面 [一句话,说明矛盾关系]
|
||||
@@ -217,10 +321,13 @@ def parse_response(text):
|
||||
return result
|
||||
|
||||
def save_result(code, full_text, parsed):
|
||||
"""保存LLM结果到DB"""
|
||||
"""保存LLM结果到DB(先快照再UPDATE)"""
|
||||
conn = sqlite3.connect(DB)
|
||||
now = datetime.now().isoformat()
|
||||
|
||||
# ── 修改前快照 ──
|
||||
snapshot_strategy_history(conn, code, 'batch_12d')
|
||||
|
||||
updates = ["full_analysis=?", "reassessed_at=?"]
|
||||
params = [full_text, now]
|
||||
|
||||
@@ -259,107 +366,109 @@ def save_result(code, full_text, parsed):
|
||||
_sl = parsed.get("stop_loss", 0)
|
||||
_tp = parsed.get("take_profit", 0)
|
||||
_pos = parsed.get("position", "")
|
||||
_msg = f"📈 {_name}({code}) 价{_p}→12维分析生成买入信号!区间{_el}~{_eh} 损{_sl} 盈{_tp} 仓位{_pos}"
|
||||
_msg = f"\U0001f4c8 {_name}({code}) 价{_p}\u219212维分析生成买入信号!区间{_el}~{_eh} 损{_sl} 盈{_tp} 仓位{_pos}"
|
||||
import urllib.request, json as _jj
|
||||
_req = urllib.request.Request("http://127.0.0.1:5805/",
|
||||
data=_jj.dumps({"body": _msg, "to": "hmo@yoin.fun", "type": "chat"}).encode(),
|
||||
headers={"Content-Type": "application/json"})
|
||||
urllib.request.urlopen(_req, timeout=5)
|
||||
print(f" 📨 XMPP推送成功: {_msg[:60]}")
|
||||
print(f" \U0001f4e8 XMPP推送成功: {_msg[:60]}")
|
||||
except Exception as _e:
|
||||
print(f" ⚠️ XMPP推送失败: {_e}")
|
||||
print(f" \u26a0\ufe0f XMPP推送失败: {_e}")
|
||||
|
||||
conn.close()
|
||||
|
||||
def process_stock(code):
|
||||
def process_stock(code, force_today=False):
|
||||
"""处理单只股票"""
|
||||
print(f"\n{'='*50}")
|
||||
print(f"处理: {code}")
|
||||
print(f"{'='*50}")
|
||||
|
||||
if has_llm_analysis(code):
|
||||
print(f" ⏭ 已有LLM九维分析,跳过")
|
||||
if in_cooldown(code):
|
||||
print(f" \u23ed 冷却期内,跳过")
|
||||
return False
|
||||
|
||||
if in_cooldown(code):
|
||||
print(f" ⏭ 冷却期内,跳过")
|
||||
# 有分析且未过期 \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" ⚠️ 无价格数据,跳过")
|
||||
print(f" \u26a0\ufe0f 无价格数据,跳过")
|
||||
return False
|
||||
|
||||
print(f" 调LLM生成九维分析...", flush=True)
|
||||
prompt = build_prompt(data)
|
||||
|
||||
try:
|
||||
r = subprocess.run(["curl", "-s", "--max-time", "300",
|
||||
"-H", "Content-Type: application/json",
|
||||
"-H", "Authorization: Bearer hermes123",
|
||||
"-d", json.dumps({"model":"deepseek-v4-flash","messages":[{"role":"user","content":prompt}],"max_tokens":2048}),
|
||||
GATEWAY], capture_output=True, timeout=310)
|
||||
|
||||
if r.returncode != 0:
|
||||
print(f" ❌ curl失败: {r.stderr.decode()[:100]}")
|
||||
return False
|
||||
|
||||
resp = json.loads(r.stdout)
|
||||
if "choices" not in resp:
|
||||
print(f" ❌ API异常: {str(resp)[:200]}")
|
||||
return False
|
||||
|
||||
full_text = resp["choices"][0]["message"]["content"]
|
||||
print(f" ✅ LLM返回({len(full_text)}字)", flush=True)
|
||||
|
||||
parsed = parse_response(full_text)
|
||||
print(f" 信号={parsed['signal']} 区间={parsed['entry_low']}~{parsed['entry_high']} 损={parsed['stop_loss']} 盈={parsed['take_profit']} 仓位={parsed['position']}")
|
||||
|
||||
save_result(code, full_text, parsed)
|
||||
print(f" ✅ 已保存到DB")
|
||||
return True
|
||||
|
||||
except subprocess.TimeoutExpired:
|
||||
print(f" ❌ 超时")
|
||||
return False
|
||||
except Exception as e:
|
||||
print(f" ❌ 错误: {e}")
|
||||
# ── 使用共享 LLM 客户端(替代 curl subprocess)──
|
||||
result = call_llm(prompt, model=REASSESS_MODEL, max_tokens=4096)
|
||||
|
||||
if not result["ok"]:
|
||||
print(f" \u274c LLM调用失败: {result.get('error','未知错误')}")
|
||||
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)
|
||||
print(f" 信号={parsed['signal']} 区间={parsed['entry_low']}~{parsed['entry_high']} 损={parsed['stop_loss']} 盈={parsed['take_profit']} 仓位={parsed['position']}")
|
||||
|
||||
save_result(code, full_text, parsed)
|
||||
print(f" \u2705 已保存到DB")
|
||||
return True
|
||||
|
||||
def main():
|
||||
# ── Gateway 预检:不可用则立即退出(不阻塞 cron)──
|
||||
if not gateway_alive():
|
||||
print("[FATAL] Hermes Gateway 不可用,退出(检查 http://127.0.0.1:8643/v1/models)")
|
||||
sys.exit(1)
|
||||
|
||||
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)
|
||||
rows = conn.execute("SELECT code FROM holding_strategies WHERE status='active' AND decision_type='自选策略' ORDER BY code").fetchall()
|
||||
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]
|
||||
|
||||
print(f"待处理: {len(codes)}只")
|
||||
print(f"待处理: {len(codes)}只 (type={dtype or 'all'}, force_today={force_today})")
|
||||
|
||||
ok = 0
|
||||
fail = 0
|
||||
skip = 0
|
||||
for i, code in enumerate(codes):
|
||||
if has_llm_analysis(code):
|
||||
print(f" [{i+1}/{len(codes)}] ⏭ {code} 已有LLM分析")
|
||||
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):
|
||||
if process_stock(code, force_today):
|
||||
ok += 1
|
||||
else:
|
||||
fail += 1
|
||||
|
||||
# 间隔15秒(防gateway过载)
|
||||
# 间隔8秒(pro model较重但gateway可承受;retry逻辑吸收瞬断)
|
||||
if i < len(codes) - 1:
|
||||
print(f" 等待15秒...", flush=True)
|
||||
time.sleep(15)
|
||||
print(f" 等待8秒...", flush=True)
|
||||
time.sleep(8)
|
||||
|
||||
print(f"\n{'='*50}")
|
||||
print(f"完成: {ok}成功, {fail}失败, {skip}跳过")
|
||||
|
||||
Reference in New Issue
Block a user