#!/usr/bin/env python3 # -*- coding: utf-8 -*- """strategy_period_rollup.py — 策略周期统计每日滚动派生(2026-08-15 问题2:策略数据每日自动更新) 背景: dashboard 研究 Tab 的 1y/2y/5y 列依赖 strategy_research 对应 period_tag 行, 但 s2_panic/v_lurk_v3/v_mr4 等激活策略只有 10y 记录 → 列空、health_monitor 无 5y 基线。 本脚本每日从「最长周期主记录」的 trades 按窗口切出 1y/2y/5y,派生行 upsert 回表。 原则: 1. 窗口锚定主记录数据末端(max entry_date),不锚今天——回测未重跑时窗口不缩水 2. 不覆盖真实研究记录:只 upsert 带 derived_from 标记的派生行;真实行存在则跳过 3. 派生行含 calc_summary + portfolio_sim(10槽) + portfolio_full(100槽), 口径与 strategy_lab.list_strategies 的 read-time 切窗一致(1m/6m/1y 已有同类逻辑) 调度:每日盘后(策略路由后、regime_perf 前),hermes cron。 """ import sys import json import sqlite3 from pathlib import Path from datetime import datetime, timedelta from collections import Counter # ── 消息通道统一路由(broadcast/xmpp by delivery) ── try: from messenger import install_stdio_hook as _msh _msh() except Exception: pass _SCRIPT_DIR = Path(__file__).resolve().parent sys.path.insert(0, str(_SCRIPT_DIR)) sys.path.insert(0, "/home/hmo/MoFin") DB = "/home/hmo/MoFin/data/mofin.db" # 目标派生周期(10y 作为主记录源,不派生) TARGET_PERIODS = {"1y": 365, "2y": 730, "5y": 1826} # 周期强度排序(选主记录用):越长越优先 _PERIOD_RANK = {"10y": 6, "7y": 5, "5y": 4, "2y": 3, "1y": 2, "6m": 1, "1m": 0} def _period_rank(tag): return _PERIOD_RANK.get(tag or "2y", 3) def load_master_rows(conn): """每个 (version, market) 取周期最长、含 trades 的记录作为主记录""" rows = conn.execute( "SELECT id, version, name, summary, hypothesis, parent, config_json, results_json, " "period, COALESCE(market,'all'), COALESCE(period_tag,'2y') FROM strategy_research" ).fetchall() best = {} for r in rows: key = (r[1], r[9]) rank = _period_rank(r[10]) if key not in best or (rank, r[0]) > (best[key][0], best[key][1]): best[key] = (rank, r[0], r) return [v[2] for v in best.values()] def existing_real_periods(conn, version, market): """该 (version, market) 已存在的【真实】period_tag(results_json 无 derived_from 标记)""" rows = conn.execute( "SELECT COALESCE(period_tag,'2y'), results_json FROM strategy_research " "WHERE version=? AND COALESCE(market,'all')=?", (version, market)).fetchall() real, derived = set(), set() for tag, rj in rows: try: d = json.loads(rj or "{}") except Exception: d = {} if d.get("derived_from"): derived.add(tag) else: real.add(tag) return real, derived def derive_period(master, target_tag, days, calc_summary, portfolio_sim): """从主记录切窗派生 target_tag 记录。返回 (results_dict, period_str) 或 None""" res = json.loads(master["results_json"]) trades = res.get("trades") or [] if not trades: return None max_date = max(t.get("entry_date", "") for t in trades) if not max_date: return None cutoff = (datetime.strptime(max_date, "%Y-%m-%d") - timedelta(days=days)).strftime("%Y-%m-%d") sliced = [t for t in trades if t.get("entry_date", "") >= cutoff] if not sliced: return None for t in sliced: t.setdefault("boost", 1.0) capital = res.get("capital") or 1000000 summary = calc_summary(sliced, capital) try: summary["portfolio"] = portfolio_sim(sliced, 1000000, max_positions=10) summary["portfolio_full"] = portfolio_sim(sliced, 1000000, max_positions=100) except Exception as e: print(f" [warn] portfolio_sim 失败({e})", flush=True) year_dist = dict(Counter(t["entry_date"][:4] for t in sliced if t.get("entry_date"))) out = { "strategy": res.get("strategy") or master["version"], "strategy_name": res.get("strategy_name") or master["name"], "market": res.get("market") or master["market"], "period": f"{cutoff} ~ {max_date}", "period_tag": target_tag, "capital": capital, "total_stocks_screened": None, # 窗口派生无法还原,置空而非伪造 "scored_events": None, "trades": sliced, "summary": summary, "year_dist": year_dist, "derived_from": master["period_tag"], "derived_at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), } return out def main(): from strategy_lab import calc_summary, portfolio_sim now = datetime.now().strftime("%Y-%m-%d %H:%M:%S") conn = sqlite3.connect(DB, timeout=60) conn.execute("PRAGMA busy_timeout=60000") masters = load_master_rows(conn) print(f"[period_rollup] {now} 主记录 {len(masters)} 个 (version, market)", flush=True) inserted = skipped_real = skipped_empty = 0 for m in masters: master = { "id": m[0], "version": m[1], "name": m[2], "summary": m[3], "hypothesis": m[4], "parent": m[5], "config_json": m[6], "results_json": m[7], "period": m[8], "market": m[9], "period_tag": m[10], } real, derived = existing_real_periods(conn, master["version"], master["market"]) for tag, days in TARGET_PERIODS.items(): if tag == master["period_tag"]: continue # 主记录本身就是该周期 if _period_rank(tag) >= _period_rank(master["period_tag"]): continue # 只向更短周期派生 if tag in real: skipped_real += 1 continue # 真实研究记录存在,不覆盖 try: res = derive_period(master, tag, days, calc_summary, portfolio_sim) except Exception as e: print(f" [err] {master['version']}/{master['market']}/{tag}: {e}", flush=True) continue if not res: skipped_empty += 1 continue # 删旧派生行,插新派生行(保持每日最新) conn.execute( "DELETE FROM strategy_research WHERE version=? AND COALESCE(market,'all')=? " "AND COALESCE(period_tag,'2y')=? AND results_json LIKE '%\"derived_from\"%'", (master["version"], master["market"], tag)) conn.execute( "INSERT INTO strategy_research " "(version, name, summary, hypothesis, parent, config_json, results_json, " " analysis_json, period, created_at, market, period_tag) " "VALUES (?,?,?,?,?,?,?,?,?,?,?,?)", (master["version"], master["name"], master["summary"], master["hypothesis"], master["parent"], master["config_json"], json.dumps(res, ensure_ascii=False), json.dumps({"derived": True, "from": master["period_tag"]}, ensure_ascii=False), res["period"], now, master["market"], tag)) inserted += 1 s = res["summary"] print(f" + {master['version']:<14} {master['market']:<3} {tag}: " f"{s.get('total_trades')}笔 胜率{s.get('win_rate')}% (派生自{master['period_tag']})", flush=True) conn.commit() conn.close() print(f"[period_rollup] 完成: 派生写入 {inserted} 行 | 真实记录跳过 {skipped_real} | 空窗跳过 {skipped_empty}", flush=True) if __name__ == "__main__": main()