# sanguo_api/routes_pipeline.py """/pipeline/* 接真端点——spec §4.5 决议 L 六件之 4(驾驶舱接真=决议 L 增补). 前端零改动=删 mock 路由即接真(api/pipeline.ts 适配层签名不变); B' 架构=双机部署同一份代码各读各的本地数据,路径全走 env/cwd 相对, P2 零新增跨机推送。promotion 执行(起实盘实例)仍是用户手动动作,档案只记录。 """ from __future__ import annotations import asyncio import hmac import json import logging import os import re import sys import threading from dataclasses import asdict from datetime import date, datetime from typing import Any, Literal from fastapi import APIRouter, Depends, HTTPException from sanguo_api import decompose_jobs from sanguo_api.hypothesis_card import ( build_wizard_messages as _build_wizard_messages, load_domains as _wizard_domains, normalize_draft as _normalize_draft, ) from sanguo_api.hypothesis_decompose import run_decompose as _run_decompose from sanguo_api.llm import LLMClient, LLMConfigError, LLMError, config_store from sanguo_api.llm.probe import probe_connection from sanguo_api.routes import verify_token from sanguo_portfolio import pipeline_store router = APIRouter(dependencies=[Depends(verify_token)]) logger = logging.getLogger(__name__) # graduate 全程进程内互斥:FastAPI sync 路由跑线程池,防并发双发双转移双事件 _GRADUATE_LOCK = threading.Lock() _DECOMPOSE_LOCK = threading.Lock() # P2-6:decompose 全程持锁(仿 _GRADUATE_LOCK) def _registry_path() -> str: return os.environ.get("SANGUO_STRATEGY_REGISTRY", os.path.join("data", "strategy_registry.yaml")) def _pipeline_db() -> str: from sanguo_portfolio import pipeline_store return pipeline_store.default_pipeline_db_path() def _load_strategy_registry() -> dict: from sanguo_portfolio.strategy_registry import ( ensure_runtime_registry, load_registry) return load_registry(ensure_runtime_registry(_registry_path())) def _fmt_vs(v: float | None) -> str: return f"{v * 100:+.1f}%/周" if v is not None else "无基准" _eval_db_injected: dict[str, str | None] = {"path": None} def set_pipeline_eval_db_path(path: str) -> None: """app.create_app 注入本机 eval 库路径(与月度链同库;routes_factor 同款).""" _eval_db_injected["path"] = path def _current_eval_db() -> str: from sanguo_factor import eval_store return _eval_db_injected["path"] or eval_store.default_eval_db_path() def _monthly_dir() -> str: return os.environ.get("SANGUO_FACTOR_MONTHLY_DIR", os.path.join("reports", "factor_monthly")) # 因子晋级评审门:近窗 t≥2 才值得人评五问(与前端 PromotionReview 的 # promotionFilter.REVIEW_T_GATE 同值镜像,改门须两处同改+spec 注记) REVIEW_T_GATE = 2.0 def _replay_view(doc: dict | None) -> dict: """verdict_replay.json → 前端 replay 形状;无历史触发点 → 功课累积提示.""" from sanguo_portfolio import verdict_records labels = {t["level"]: t["action"] for t in verdict_records.SEVEN_TIERS} def _rows(outcome: str) -> list[dict]: return [{"option": f"档{lv} {labels[lv]}", "outcome": outcome} for lv in range(7)] na = "需精确验证(现有回测工具)" if not doc or not doc.get("episodes"): since = (doc or {}).get("since") or "监控起步" return {"period": f"暂无历史触发点(功课自 {since} 逐月累积)", "after3m": _rows("功课累积中")} replayable = [e for e in doc["episodes"] if e.get("replayable")] if not replayable: return {"period": f"共 {len(doc['episodes'])} 次历史触发,均在基准曲线窗外", "after3m": _rows(na)} e = replayable[-1] rows = [] for lv in range(7): st = (e.get("tiers") or {}).get(str(lv)) outcome = (f"{st['ret3m']:+.1%}/回撤 {st['max_dd']:.1%}(直加·估)" if st is not None else na) rows.append({"option": f"档{lv} {labels[lv]}", "outcome": outcome}) return {"period": f"最近触发 {e['month']}(共 {len(doc['episodes'])} 次)之后 3 个月", "after3m": rows} @router.get("/pipeline/ladder") def ladder() -> dict: """registry+周报+gate 态拼成前端 LadderStrategy 形状(retired 不上梯).""" from sanguo_portfolio import pipeline_store from sanguo_portfolio.promotion_gate import evaluate_shadow_gate from sanguo_portfolio.strategy_registry import ( shadow_total_weeks, shadow_weeks) reg = _load_strategy_registry() today = date.today().isoformat() thresholds = pipeline_store.get_effective_thresholds(_pipeline_db()) items = [] for name in sorted(reg["strategies"]): e = reg["strategies"][name] stage = e.get("stage") if stage not in ("backtest", "paper", "shadow", "live"): continue item: dict[str, Any] = {"id": name, "name": name, "stage": stage, "note": "补录档案(gate 快照留空)" if e.get("backfilled") else ""} if stage == "shadow": item["shadowWeeks"] = shadow_weeks(e, today) item["shadowTotalWeeks"] = shadow_total_weeks(e) w = pipeline_store.latest_weekly(_pipeline_db(), name) if w and w.get("te_annual") is not None \ and w.get("fill_rate") is not None \ and w.get("vs_backtest") is not None: item["weeklyReport"] = { "trackingError": w["te_annual"], "fillRate": w["fill_rate"], "vsBacktest": _fmt_vs(w["vs_backtest"])} gate = evaluate_shadow_gate( item["shadowWeeks"], w["te_annual"], w["fill_rate"], abs(w["vs_backtest"]), thresholds) item["note"] = (item["note"] + " " if item["note"] else "") + \ ("gate ✓" if gate.passed else "gate ✗ 未毕业") else: item["note"] = (item["note"] + " " if item["note"] else "") + \ "周报待算" cap = e.get("capital") if cap: v = cap.get("value") item["capitalShare"] = (f"{v}%" if cap.get("mode") == "占账户比例" else f"¥{float(v):,.0f}") items.append(item) return {"items": items} @router.get("/pipeline/gate-config") def gate_config_get() -> dict: """配置页:代码默认+本机运行值+审计行(决议 L 阈值可配置页面).""" from sanguo_portfolio import pipeline_store from sanguo_portfolio.promotion_gate import DEFAULT_THRESHOLDS return {"defaults": asdict(DEFAULT_THRESHOLDS), "effective": asdict(pipeline_store.get_effective_thresholds( _pipeline_db())), "audit": pipeline_store.get_audit(_pipeline_db())} @router.put("/pipeline/gate-config") def gate_config_put(body: dict) -> dict: from sanguo_portfolio import pipeline_store key = body.get("key") value = body.get("value") changed_by = body.get("changed_by") or "console" if not isinstance(key, str) or not isinstance(value, (int, float)) \ or isinstance(value, bool): raise HTTPException(422, "body 需要 {key, value:number, changed_by}") try: pipeline_store.set_threshold(_pipeline_db(), key, float(value), changed_by) except ValueError as e: raise HTTPException(422, str(e)) from e return {"ok": True} @router.post("/pipeline/promotion/graduate") def graduate(body: dict) -> dict: """毕业=口令+影子 gate 硬闸+分配方式→转移+落档+开 issue(触点④,决议 B).""" from sanguo_portfolio import pipeline_store from sanguo_portfolio.promotion_gate import evaluate_shadow_gate from sanguo_portfolio.strategy_registry import ( append_event, render_graduate_issue, save_registry, shadow_weeks, shadow_total_weeks, transition) name = body.get("name") or "" passphrase = body.get("passphrase") or "" expect = os.environ.get("SANGUO_PROMOTION_PASSPHRASE", "") if not expect: raise HTTPException(503, "未配置 SANGUO_PROMOTION_PASSPHRASE(首班用户定)") if not hmac.compare_digest(passphrase.encode("utf-8"), expect.encode("utf-8")): raise HTTPException(403, "口令不符") with _GRADUATE_LOCK: reg = _load_strategy_registry() entry = reg["strategies"].get(name) if entry is None: raise HTTPException(404, f"策略未注册: {name}") if entry["stage"] != "shadow": raise HTTPException(422, f"仅 shadow 态可毕业,当前 {entry['stage']}") today = date.today().isoformat() weeks = shadow_weeks(entry, today) w = pipeline_store.latest_weekly(_pipeline_db(), name) if w is None or w.get("te_annual") is None or w.get("fill_rate") is None: raise HTTPException(422, "周报数据不足(TE/fillRate 未算出),等 eod 后重算") thresholds = pipeline_store.get_effective_thresholds(_pipeline_db()) vs = w.get("vs_backtest") gate = evaluate_shadow_gate(weeks, w["te_annual"], w["fill_rate"], abs(vs) if vs is not None else 1.0, thresholds) if not gate.passed: raise HTTPException(422, { "error": "影子毕业 gate 未过(延长影子期或退回工厂)", "checks": [c.__dict__ for c in gate.checks], "weeks": weeks, "total_weeks": shadow_total_weeks(entry)}) try: alloc_value = float(str(body.get("allocValue", "0")).replace(",", "")) except ValueError as exc: raise HTTPException(422, "allocValue 须为数字") from exc capital = {"mode": body.get("allocMode") or "固定金额", "value": alloc_value, "at": today} snap = gate.to_dict() snap["weeks"] = weeks snap["te_annual"] = w["te_annual"] snap["fill_rate"] = w["fill_rate"] snap["vs_backtest"] = w.get("vs_backtest") extra: dict[str, Any] = {"gate_snapshot": snap, "capital": capital, "at": today} if body.get("reviewIssue"): try: extra["review_issue"] = int(body["reviewIssue"]) except (TypeError, ValueError): raise HTTPException(422, "reviewIssue 须为整数") new_entry = transition(reg, name, "live", **extra) save_registry(_registry_path(), reg) issue_no = None warning = None try: events_path = os.environ.get( "SANGUO_STRATEGY_EVENTS", os.path.join("data", "strategy_events.jsonl")) append_event(events_path, new_entry, "shadow", "live", capital=capital) title, body_text = render_graduate_issue(name, new_entry) issue_dir = os.environ.get("SANGUO_REGISTRY_ISSUE_DIR", os.path.join("reports", "registry_events")) os.makedirs(issue_dir, exist_ok=True) with open(os.path.join(issue_dir, f"{name}_graduate_issue.md"), "w", encoding="utf-8") as f: f.write(f"# {title}\n\n{body_text}\n") token = os.environ.get("SANGUO_GITEA_TOKEN") if token: from sanguo_portfolio.strategy_registry import open_gitea_issue issue_no = open_gitea_issue(title, body_text, token) except Exception as exc: # noqa: BLE001 转移已落档,尾巴失败不回滚不 500 logger.warning("graduate %s 转移已落档,事件/模板落盘失败: %s", name, exc) warning = "转移已落档,事件/模板落盘失败需人工补" resp: dict[str, Any] = {"ok": True, "issue": issue_no} if warning: resp["warning"] = warning return resp @router.get("/pipeline/todos") def todos() -> dict: """驾驶舱聚合总入口(决议 L 增补):读各域本机落盘工件拼接,域报告保持域内纯净.""" from sanguo_portfolio import pipeline_store from sanguo_portfolio.promotion_gate import evaluate_shadow_gate from sanguo_portfolio.strategy_registry import ( shadow_weeks, shadow_total_weeks) items: list[dict[str, Any]] = [] now = datetime.now().isoformat(timespec="seconds") thresholds = pipeline_store.get_effective_thresholds(_pipeline_db()) # ① 毕业候选(shadow 满 total 周且 gate 过) reg = _load_strategy_registry() today = date.today().isoformat() for name, e in reg["strategies"].items(): if e.get("stage") != "shadow": continue weeks = shadow_weeks(e, today) if weeks < shadow_total_weeks(e): continue w = pipeline_store.latest_weekly(_pipeline_db(), name) gate_ok = False if w and w.get("te_annual") is not None: vs = w.get("vs_backtest") gate_ok = evaluate_shadow_gate( weeks, w["te_annual"], w.get("fill_rate") or 0.0, abs(vs) if vs is not None else 1.0, thresholds).passed if gate_ok: items.append({"id": f"grad-{name}", "touchpoint": "promotion", "title": f"{name} 影子期满且 gate 过,等口令+初始分配", "detail": f"周数 {weeks}/{shadow_total_weeks(e)}", "path": "/pipeline/incubation", "severity": "warn", "updatedAt": now}) # ①′ 因子晋级评审候选(触点②另一半,10-08 补源):可判且近窗 t≥2 # = 评审页下拉同门——待办板与评审页从此同口径,「N 件等你决策」说真话 from sanguo_factor import version_registry as vr from sanguo_portfolio.strategy_registry import ensure_runtime_registry try: freg = vr.load_registry(ensure_runtime_registry( _factor_registry_path(), os.path.join("config", "factor_registry.yaml"))) points, _ = _eval_latest_points() for fname, fe in freg["factors"].items(): if fe.get("status") != "assessable": continue t = (points.get(fname) or {}).get("t") if t is None or t < REVIEW_T_GATE: continue items.append({"id": f"review-{fname}", "touchpoint": "promotion", "title": f"{fname} 待晋级评审(五问)", "detail": f"可判且近窗 t={t:.2f}≥{REVIEW_T_GATE:.0f}," "评审页等你判定", "path": "/pipeline/review", "severity": "warn", "updatedAt": now}) except Exception: # noqa: BLE001 注册表/eval 库异常不拖垮待办聚合 pass # ② 因子衰减告警+③ 集体水位(本机最新月度批评 JSON;最新有效 JSON,损坏跳过) from sanguo_factor.verdict_replay import load_latest_monthly rep = load_latest_monthly(_monthly_dir()) if rep: for fname, v in (rep.get("verdicts") or {}).items(): if v.get("state") == "alert": items.append({"id": f"decay-{fname}", "touchpoint": "alert", "title": f"{fname} 衰减告警", "detail": "月度批评:连续 3 月双条件满足", "path": "/pipeline/factors", "severity": "warn", "updatedAt": now}) if (rep.get("collective") or {}).get("is_collective_decay"): items.append({"id": "collective", "touchpoint": "alert", "title": "集体水位跌破阈——研判卡待裁决", "detail": "集体衰减判别门(决议 C):权重冻结,上交研判", "path": "/pipeline/attribution", "severity": "critical", "updatedAt": now}) # ④ 数据缺口(data_gaps.json);P2-11 标记件={"manifest_missing": true}—— # 旧循环 join(True) TypeError 被 except 吞=最需告警场景零输出,#91③ # 标记键单列转可见告警项,正常形状循环不变 gaps_path = os.path.join(_monthly_dir(), "data_gaps.json") if os.path.exists(gaps_path): try: with open(gaps_path, encoding="utf-8") as f: gaps = json.load(f) if gaps.get("manifest_missing"): items.append({"id": "gap-manifest-missing", "touchpoint": "hypothesis", "title": "数据底座 manifest 缺失", "detail": "月度批 manifest 缺失/不可读,缺口按" "缺失处理批照跑(P2-11 标记件)——" "修复 config/data_manifest.yaml 后" "下月批自愈", "path": "/pipeline/hypotheses", "severity": "warn", "updatedAt": now}) else: for fname in sorted(gaps): items.append({"id": f"gap-{fname}", "touchpoint": "hypothesis", "title": f"{fname} 数据缺口待补", "detail": f"缺源: {','.join(gaps[fname])}", "path": "/pipeline/hypotheses", "severity": "info", "updatedAt": now}) except Exception: pass # ⑤ 假设卡琥珀待办(D7 件③,人工卡点②显性化):已确认未分解的卡 # (queued/data_check=尚未成功过分解;分解成功即转 building 消行) for it in pipeline_store.list_hypotheses(_pipeline_db()): if it["state"] not in ("queued", "data_check"): continue items.append({"id": f"decompose-{it['id']}", "touchpoint": "hypothesis", "title": f"{it['title'][:24]}——已确认未分解", "detail": "人工卡点②:你点「AI 分解」,机器落册孵化因子", "path": "/pipeline/hypotheses", "severity": "warn", "updatedAt": now}) order = {"critical": 0, "warn": 1, "info": 2} items.sort(key=lambda i: order.get(i["severity"], 9)) return {"items": items} def _factor_registry_path() -> str: return os.environ.get("SANGUO_FACTOR_REGISTRY", os.path.join("data", "factor_registry.yaml")) @router.get("/pipeline/stations") def pipeline_stations() -> dict: """流水线站点条(D1 地铁图融总览,v4 高保真稿 SUBWAY 落地):七站计数+人工卡点 pending. 每站 {name,count,pending,who}:count=站内在册量;pending=人工卡点等你 (琥珀);who=等的是谁(hover 文案).数据全现成端点拼装,无新算. """ from sanguo_portfolio import pipeline_store stations: list[dict[str, Any]] = [] # 假设池:活跃卡(非终态);pending=已确认未分解(queued/data_check,同触点⑤门) cards = pipeline_store.list_hypotheses(_pipeline_db()) active = [c for c in cards if c["state"] not in ("graveyard", "promoted")] pend_hypo = [c for c in active if c["state"] in ("queued", "data_check")] stations.append({"name": "假设池", "count": len(active), "pending": len(pend_hypo), "who": "、".join(c["title"][:12] for c in pend_hypo[:3]) or None}) # 因子工厂:非 graveyard 因子数;pending 无(评审卡点在下一站) try: from sanguo_factor import version_registry as vr freg = vr.load_registry(_ensure_factor_runtime_registry()) n_fac = sum(1 for e in freg["factors"].values() if e.get("status") != "graveyard") except Exception: # noqa: BLE001 坏档如实 0,不拖垮站点条 n_fac = 0 stations.append({"name": "因子工厂", "count": n_fac, "pending": 0, "who": None}) # 晋级评审:可判且 t≥2(与触点②/评审页同门);who=候选因子名 try: points, _ = _eval_latest_points() cands = [n for n, e in freg["factors"].items() if e.get("status") == "assessable" and (points.get(n) or {}).get("t") is not None and (points.get(n) or {}).get("t", 0) >= REVIEW_T_GATE] except Exception: # noqa: BLE001 cands = [] stations.append({"name": "晋级评审", "count": len(cands), "pending": len(cands), "who": "、".join(f"{n} 待五问" for n in cands[:3]) or None}) # 合成层:在位+挑战者(quant12 族两版) stations.append({"name": "合成层", "count": 2, "pending": 0, "who": None}) # 策略出生:合成层晋级后的出生卡点(暂无数据源,诚实 0) stations.append({"name": "策略出生", "count": 0, "pending": 0, "who": None}) # 孵化梯+实盘:strategy_registry stage 计数 try: sreg = _load_strategy_registry() stages = [e.get("stage") for e in sreg["strategies"].values()] n_ladder = sum(1 for s in stages if s in ("backtest", "paper", "shadow")) n_live = sum(1 for s in stages if s == "live") except Exception: # noqa: BLE001 n_ladder = n_live = 0 stations.append({"name": "孵化梯", "count": n_ladder, "pending": 0, "who": None}) stations.append({"name": "实盘", "count": n_live, "pending": 0, "who": None}) # 健康水位五格(v4 HEALTH):衰减告警/集体水位/数据缺口/数据链路/as_of from sanguo_factor.verdict_replay import load_latest_monthly n_alert = n_gap = 0 collective_txt, asof_txt, link_txt, link_state = "正常", "—", "—", "idle" try: rep = load_latest_monthly(_monthly_dir()) if rep: n_alert = sum(1 for v in (rep.get("verdicts") or {}).values() if v.get("state") == "alert") if (rep.get("collective") or {}).get("is_collective_decay"): collective_txt = "跌破" asof_txt = str(rep.get("as_of") or "—") except Exception: # noqa: BLE001 pass gaps_path2 = os.path.join(_monthly_dir(), "data_gaps.json") if os.path.exists(gaps_path2): try: with open(gaps_path2, encoding="utf-8") as f: n_gap = len(json.load(f)) except Exception: # noqa: BLE001 pass health = [ {"name": "衰减告警", "value": str(n_alert), "state": "warn" if n_alert else "ok"}, {"name": "集体水位", "value": collective_txt, "state": "warn" if collective_txt == "跌破" else "ok"}, {"name": "数据缺口", "value": str(n_gap), "state": "warn" if n_gap else "ok"}, {"name": "数据链路", "value": link_txt, "state": link_state}, {"name": "as_of", "value": asof_txt, "state": "idle"}, ] return {"stations": stations, "health": health, "edges": [ {"from": "假设池", "to": "因子工厂", "cond": "落成+注册", "trig": "人工"}, {"from": "因子工厂", "to": "晋级评审", "cond": "t≥2", "trig": "自动上板"}, {"from": "晋级评审", "to": "合成层", "cond": "五问晋级→入族登记", "trig": "人工"}, {"from": "合成层", "to": "策略出生", "cond": "族晋级→卡点生成", "trig": "自动"}, {"from": "策略出生", "to": "孵化梯", "cond": "登记+考卷", "trig": "人工+自动判"}, {"from": "孵化梯", "to": "实盘", "cond": "毕业口令", "trig": "人工"}]} def _eval_latest_points() -> tuple[dict[str, dict], str]: """eval_db 最近批 per-factor 点(monthly_review.pick/extract 同款)+批 end. 缺库/无批/异常→({}, "")——看板不 500,IC 值如实 null(前端 fmt 显示 '—')。 """ from sanguo_factor import monthly_review try: db = _current_eval_db() if not os.path.exists(db): return {}, "" run = monthly_review.pick_eval_run(db, date.today().isoformat()) if run is None: return {}, "" return monthly_review.extract_points(db, run), str(run.get("end") or "") except Exception as exc: # noqa: BLE001 eval db 异常不拖垮看板 logger.warning("factors eval_db 读取失败,如实返空: %s", exc) return {}, "" def _factor_category(source: str) -> str: if "pledge" in source: return "pledge" if "research" in source: return "research" if "corpus" in source or "sentiment" in source: return "sentiment" if source.startswith("fund") or "fundamental" in source: return "fundamental" return "technical" @router.get("/pipeline/factors") def factors() -> dict: """因子看板真实数据:factor 注册表(运行副本)+eval_db 最近批 t 值(缺库如实 null). 候选下拉/看板从 mock 接真:graveyard 不上此板(墓园归假设池页); eval_db 不存在(月度批评首班未跑)→ IC 两窗 null,前端 fmt 显示 '—'。 """ from sanguo_factor import version_registry as vr from sanguo_portfolio.strategy_registry import ensure_runtime_registry reg = vr.load_registry(ensure_runtime_registry( _factor_registry_path(), os.path.join("config", "factor_registry.yaml"))) points, last_eval = _eval_latest_points() # 全史 t(=全期合成 t)从最新月报 factor_stats 一次性接真(spec §4.2 承诺 # factors 端点升级;单次读取,勿循环内重读——codex review HIGH); # ∪daily/ 取新(2026-10-11 补丁) reports = _monthly_reports(include_daily=True) full_t: dict[str, Any] = ( {n: (s or {}).get("tAll") for n, s in (reports[0][2].get("factor_stats") or {}).items()} if reports else {}) items: list[dict[str, Any]] = [] for name, e in reg["factors"].items(): if e.get("status") == "graveyard": continue # 墓园不上板 versions = e.get("versions") or [] source = str(((versions[-1].get("params") or {}).get("source")) or "") \ if versions else "" t = (points.get(name) or {}).get("t") items.append({"id": name, "category": _factor_category(source), "version": str(versions[-1]["v"]) if versions else "1", "status": e["status"], # D7 件②:origin 身份戳+来源卡回链血统(manual=缺省) "origin": e.get("origin") or "manual", "hypothesis": e.get("hypothesis"), "description": e.get("description") or "", "icRecentT": t, # 最近批(12M 滚动窗)即近窗 t "icFullT": full_t.get(name), # 全期合成 t(月报 factor_stats.tAll) "promotedAtT": e.get("promotion_t") or None, "similarity": e.get("similarity"), # 疑似换皮提示(注册时落) "decayMonths": 0, "lastEvalDate": last_eval}) order = {"decaying": 0, "promoted": 1, "assessable": 2, "incubating": 3, "retired": 4} items.sort(key=lambda i: (order.get(i["status"], 9), i["id"])) return {"items": items} @router.get("/pipeline/factors/graveyard") def factors_graveyard() -> dict: """因子墓园(10-09 可见性缺口补):registry graveyard 条目带死因出列表. 判死因子不上工厂板、详情页 404——本端点是唯一视图(假设池页墓园 tab 判死因子区);老编号因子(hypothesis=H-*)无卡可回写,也在此兜底可见. """ from sanguo_factor import version_registry as vr from sanguo_portfolio.strategy_registry import ensure_runtime_registry reg = vr.load_registry(ensure_runtime_registry( _factor_registry_path(), os.path.join("config", "factor_registry.yaml"))) items = [{"name": name, "cause": str(e.get("cause") or ""), "hypothesis": e.get("hypothesis"), "origin": e.get("origin") or "manual"} for name, e in reg["factors"].items() if e.get("status") == "graveyard"] items.sort(key=lambda i: i["name"]) return {"items": items} _CROSS_DIFF_MAX = 0.05 # 双机对拍分歧阈(起步默认,首年校准;spec 对拍规则未定数) def _monthly_reports(include_daily: bool = False) -> list[tuple[str, str, dict]]: """monthly dir 全部报告 [(host, as_of, doc)];按 as_of 降序,坏文件跳过. 文件名={host}_{as_of}.json(monthly_review.main 落盘原文)。 include_daily=True 时并入 daily/ 子目录件(2026-10-11 双层分工补丁: 速览读最新——根层∪daily 取 as_of 新者);判定层(月末快照语义)不传参, 仍只读根层,一行不动。 """ d = _monthly_dir() if not os.path.isdir(d): return [] out: list[tuple[str, str, dict]] = [] def _scan(folder: str) -> None: for fn in os.listdir(folder): if not fn.endswith(".json"): continue host, _, as_of = fn[:-5].rpartition("_") if not host or len(as_of) != 10 or as_of[4] != "-": continue try: with open(os.path.join(folder, fn), encoding="utf-8") as f: out.append((host, as_of, json.load(f))) except (OSError, json.JSONDecodeError): continue _scan(d) if include_daily: daily = os.path.join(d, "daily") if os.path.isdir(daily): _scan(daily) out.sort(key=lambda t: t[1], reverse=True) return out def _machine_label(host: str) -> str: return "NAS" if "nas" in host.lower() else "VPS" @router.get("/pipeline/monthly-review") def monthly_review_view() -> dict | None: """月度批评卡(10-01 用户拍板接真):读本机 monthly dir 最新报告 JSON 原文. 首班(10-04)未跑→null(前端 PanelCard 隐藏,不造数); crossCheck 只在同 as_of 双 host 报告都落在本机 dir 时才出双行 (对拍件汇合后自动生效),单机=空列表(不伪造"一致"); maxAbsDiff=两机同因子 t 差绝对值最大值,> _CROSS_DIFF_MAX 判"分歧"。 """ reports = _monthly_reports() if not reports: return None host, as_of, doc = reports[0] factors = doc.get("factors") or {} col = doc.get("collective") or {} peers = {h: d for h, a, d in reports if a == as_of} cross: list[dict[str, Any]] = [] if len(peers) >= 2: hosts = sorted(peers) fa = peers[hosts[0]].get("factors") or {} fb = peers[hosts[1]].get("factors") or {} common = [f for f in fa.keys() & fb.keys() if (fa.get(f) or {}).get("t") is not None and (fb.get(f) or {}).get("t") is not None] max_diff = max((abs(fa[f]["t"] - fb[f]["t"]) for f in common), default=0.0) verdict = "一致" if max_diff <= _CROSS_DIFF_MAX else "分歧" cross = [{"machine": _machine_label(h), "factors": len(peers[h].get("factors") or {}), "maxAbsDiff": round(max_diff, 6), "verdict": verdict} for h in hosts[:2]] return {"period": as_of[:7], "runDate": as_of, "factorsEvaluated": len(factors), "consensusMedianT": col.get("median"), "consensusThreshold": col.get("floor"), "crossCheck": cross} @router.get("/pipeline/factors/{name}/ic-trend") def factor_ic_trend(name: str) -> dict: """IC 趋势图(10-01 用户拍板接真):最新报告 monthly_points 的月度 t 序列. baseline 红线=已晋级因子 promotion_t×0.5(相对降幅过半位置,触点② 衰减条件①);未晋级/无基准=2.0 绝对地板(衰减条件②)——衰减告警需 两条同时满足+连续 3 月,t 为 None 的月如实跳过。 """ reports = _monthly_reports() if not reports: raise HTTPException(404, "月度批评首班未跑,无 monthly_points") _, _, doc = reports[0] pts = (doc.get("monthly_points") or {}).get(name) if not pts: raise HTTPException(404, f"因子无月度点: {name}") pairs = [(p["month"], p["t"]) for p in pts if p.get("t") is not None] from sanguo_factor import version_registry as vr from sanguo_portfolio.strategy_registry import ensure_runtime_registry reg = vr.load_registry(ensure_runtime_registry( _factor_registry_path(), os.path.join("config", "factor_registry.yaml"))) promotion_t = (reg["factors"].get(name) or {}).get("promotion_t") baseline = round(promotion_t * 0.5, 4) if promotion_t else 2.0 return {"dates": [m for m, _ in pairs], "rollingT": [t for _, t in pairs], "baseline": baseline} @router.get("/pipeline/factors/{name}/detail") def factor_detail(name: str) -> dict: """因子详情/出生档案(D7 验收反馈①):谁生的/哪批生的/点得回去. decomposer 因子附 birth=台账中 registered 含该名的最新批次 (jobId/startedAt/rounds);manual 因子 birth=null.时间线由前端 用这些原料拼(出生→入批→末评),不在后端造展示结构. """ from sanguo_factor import version_registry as vr from sanguo_portfolio.strategy_registry import ensure_runtime_registry reg = vr.load_registry(ensure_runtime_registry( _factor_registry_path(), os.path.join("config", "factor_registry.yaml"))) e = reg["factors"].get(name) if e is None or e.get("status") == "graveyard": raise HTTPException(404, f"因子不存在或已入墓园: {name}") versions = e.get("versions") or [] first = versions[0] if versions else {} params = first.get("params") or {} points, last_eval = _eval_latest_points() birth = None hyp_id = e.get("hypothesis") if e.get("origin") == "decomposer" and hyp_id: for job in pipeline_store.list_decompose_jobs(_pipeline_db(), hyp_id): if any(c.get("name") == name for c in (job.get("registered") or [])): birth = {"jobId": job["jobId"], "startedAt": job["startedAt"], "rounds": job["rounds"], "hypId": hyp_id} break stats = None reports = _monthly_reports(include_daily=True) if reports: stats = (reports[0][2].get("factor_stats") or {}).get(name) return {"factor": { "name": name, "status": e["status"], "origin": e.get("origin") or "manual", "hypothesis": hyp_id, "description": e.get("description") or "", "expression": params.get("expression"), "source": params.get("source"), "justification": params.get("justification"), "createdAt": first.get("effective_from"), "version": str(first.get("v")) if first else "1", "icRecentT": (points.get(name) or {}).get("t"), "icStats": stats, "similarity": e.get("similarity"), "lastEvalDate": last_eval}, "birth": birth} @router.post("/pipeline/factors/{name}/verdict") def factor_verdict(name: str, body: dict) -> dict: """晋级评审三按钮(触点②):promote/graveyard 走 version_registry 状态机; revise=留观(不改状态只留评审事件——升版本是 factor 域档流程,控制台不越权).""" from sanguo_factor import version_registry as vr from sanguo_portfolio.strategy_registry import ensure_runtime_registry verdict = body.get("verdict") answers = body.get("answers") or [] note = str(body.get("note") or "").strip() # codex review:纯空白死因不可过(必填语义) if verdict not in ("promote", "revise", "graveyard"): raise HTTPException(422, "verdict 须为 promote/revise/graveyard") reg_path = ensure_runtime_registry( os.environ.get("SANGUO_FACTOR_REGISTRY", os.path.join("data", "factor_registry.yaml")), os.path.join("config", "factor_registry.yaml")) events_path = os.environ.get("SANGUO_FACTOR_EVENTS", os.path.join("data", "registry_events.jsonl")) # codex 终审 C2:load→save 全程持 _DECOMPOSE_LOCK——判的因子未必属在飞卡 # (不查 is_running),锁即互斥;防 worker 旧内存对象整文件覆盖丢本次转移. with _DECOMPOSE_LOCK: reg = vr.load_registry(reg_path) if name not in reg["factors"]: raise HTTPException(404, f"因子未注册: {name}") old = reg["factors"][name]["status"] extra: dict[str, Any] = {"answers": answers, "note": note, "reviewed_via": "console"} try: if verdict == "promote": extra["promoted_at"] = date.today().isoformat() entry = vr.transition(reg, name, "promoted", **extra) action = "promoted" elif verdict == "graveyard": if not note: raise HTTPException(422, "入墓园必须带 note(死因,可检索防重复造轮)") entry = vr.transition(reg, name, "graveyard", cause=note, **extra) action = "graveyard" else: # revise=留观:五问留痕,状态不动 entry = reg["factors"][name] action = "revise(留观)" except ValueError as exc: # codex review:状态机拒绝(墓园终态重复判死/非法迁移)→422 中文原因,不再裸 500 raise HTTPException(422, str(exc)) from exc vr.save_registry(reg_path, reg) vr.append_event(events_path, entry, old, action, **extra) if verdict == "graveyard": # #91②: 控制台判死零回写(D5 回写只活在 CLI)——同病灶形态收编; # 复用 writeback_card_death(失败降级 stderr+death_writeback_missed # 附记,不阻断转移;pipeline_db 缺省 env SANGUO_PIPELINE_DB→ # data/pipeline.db,与 CLI 同源) vr.writeback_card_death(None, reg, entry, note, events_path=events_path) return {"ok": True, "action": action} def _append_card_verdict_half_state(events_path: str, hyp_id: str, stage: str, error: str) -> None: """判卡半态附记(codex 终审 I1):registry 已迁移但跨存储附属(events/ 卡片库)失败时留痕供恢复——非转移事件,带 event 键与转移事件区分. 附记本身尽力而为(再失败只 stderr,不让半态升级成 500).""" ev = {"ts": datetime.now().isoformat(timespec="seconds"), "event": "card_verdict_half_state", "hyp_id": hyp_id, "stage": stage, "error": error} try: os.makedirs(os.path.dirname(os.path.abspath(events_path)), exist_ok=True) with open(events_path, "a", encoding="utf-8") as f: f.write(json.dumps(ev, ensure_ascii=False) + "\n") except OSError as exc: print(f"[hypothesis_verdict] ⚠️ 半态附记也失败 hyp={hyp_id}: {exc}", file=sys.stderr) @router.post("/pipeline/hypotheses/{hyp_id}/verdict") def hypothesis_verdict(hyp_id: str, body: dict) -> dict: """决议O②:人工判卡片死(强断言「想法死了」,危险区操作)——卡下活因子陪葬. 原子性(codex 终审 I1 诚实化):registry 迁移内原子——全部 transition 在 内存 reg 完成后才 save_registry,中途 ValueError 零落盘.跨存储边界 (events 事件流/卡片库状态)不回滚:失败降级 stderr 告警+events 附记 card_verdict_half_state 供恢复(转移是正主,附属失败不阻断;对齐 writeback_card_death 同款降级哲学).update_hypothesis_state 的 ValueError(卡已终态竞态)单独 422 提示刷新. 陪葬 cause=卡片死因+「(随卡判死)」后缀;incubating→graveyard 状态机合法.0 因子卡=纯想法毙掉无陪葬.promoted 卡不可判死(毕业语义挂账, 决议O 只动墓园). 分解并发守卫(codex 终审 C2):在飞分解 409 快拒;load→save 全程持 _DECOMPOSE_LOCK(sync 端点线程池执行可直接 with,graduate 先例)—— 防 worker 旧内存对象整文件覆盖复活墓园因子(终态不变量). """ from sanguo_factor import version_registry as vr if decompose_jobs.is_running(hyp_id): raise HTTPException(409, "该卡分解进行中,判死请等分解结束再操作") reason = str(body.get("reason") or "").strip() if not reason: raise HTTPException(422, "判死必须带 reason(死因,可检索防重复造轮)") if len(reason) > _VERDICT_REASON_MAX: raise HTTPException(422, f"reason 过长(>{_VERDICT_REASON_MAX} 字)") items = pipeline_store.list_hypotheses(_pipeline_db()) card = next((i for i in items if i["id"] == hyp_id), None) if card is None: raise HTTPException(404, "卡片不存在") if card["state"] == "graveyard": raise HTTPException(422, "卡片已在墓园(终态不可再判)") if card["state"] == "promoted": raise HTTPException(422, "promoted 卡不可判死(决议O 挂账:毕业语义待定)") reg_path = _ensure_factor_runtime_registry() events_path = os.environ.get("SANGUO_FACTOR_EVENTS", os.path.join("data", "registry_events.jsonl")) buried: list[tuple[str, str, dict[str, Any]]] = [] # (name, old, entry) with _DECOMPOSE_LOCK: reg = vr.load_registry(reg_path) try: for name, e in sorted(reg["factors"].items()): if e.get("hypothesis") != hyp_id or e.get("status") == "graveyard": continue old = str(e.get("status")) entry = vr.transition(reg, name, "graveyard", cause=f"{reason}(随卡判死)", reviewed_via="console") buried.append((name, old, entry)) vr.save_registry(reg_path, reg) except ValueError as exc: raise HTTPException(422, str(exc)) from exc # 零副作用(未落盘) for name, old, entry in buried: try: vr.append_event(events_path, entry, old, "graveyard", cause=f"{reason}(随卡判死)", reviewed_via="console") except Exception as exc: print(f"[hypothesis_verdict] ⚠️ 事件流写入失败 hyp={hyp_id} " f"factor={name}: {exc}", file=sys.stderr) _append_card_verdict_half_state(events_path, hyp_id, "append_event", str(exc)) try: pipeline_store.update_hypothesis_state( _pipeline_db(), hyp_id, "graveyard", death_reason=reason) except ValueError as exc: # 卡片库竞态(他端点已把卡推到终态):因子已埋是正主,提示刷新不回滚 raise HTTPException(422, "卡片状态已变(竞态),因子已埋请刷新查看") from exc except Exception as exc: print(f"[hypothesis_verdict] ⚠️ 卡片埋葬回写失败 hyp={hyp_id}: {exc}", file=sys.stderr) _append_card_verdict_half_state(events_path, hyp_id, "bury_card", str(exc)) return {"ok": True, "buried": [n for n, _, _ in buried], "card": "graveyard"} @router.get("/pipeline/research-card") def research_card() -> dict: """研判卡(spec §4.7 P3 第一刀):最新月报 collective 触发才出卡;未触发 card=null.""" from sanguo_factor.verdict_replay import load_latest_monthly from sanguo_portfolio import verdict_records rep = load_latest_monthly(_monthly_dir()) if not rep: return {"card": None} col = rep.get("collective") or {} if not col.get("is_collective_decay"): return {"card": None} as_of = str(rep.get("as_of") or "") month = as_of[:7] ranked = sorted((rep.get("verdicts") or {}).items(), key=lambda kv: (kv[1].get("state") != "alert", kv[1].get("ratio") if kv[1].get("ratio") is not None else 9.0)) replay_doc: dict | None = None rp = os.path.join(_monthly_dir(), "verdict_replay.json") if os.path.exists(rp): try: with open(rp, encoding="utf-8") as f: doc = json.load(f) if not isinstance(doc, dict): doc = None replay_doc = doc except (json.JSONDecodeError, OSError): replay_doc = None ruling = next((r for r in verdict_records.load_rulings() if r.get("month") == month), None) card = {"id": f"collective-{month}", "triggeredAt": as_of, "medianT": col.get("median"), "threshold": col.get("floor"), "topDecayFactors": [k for k, _ in ranked[:3]], "options": verdict_records.SEVEN_TIERS, "replay": _replay_view(replay_doc), "ruling": ruling} return {"card": card} @router.post("/pipeline/research-card/verdict") def research_card_verdict(body: dict) -> dict: """裁决留痕(只落档不执行,决议 C);同月重复 409,非法入参 422.""" from sanguo_portfolio import verdict_records try: rec = verdict_records.append_ruling( str(body.get("month") or ""), body.get("level"), str(body.get("reason") or "")) except ValueError as exc: raise HTTPException(status_code=422, detail=str(exc)) except KeyError as exc: raise HTTPException(status_code=409, detail=str(exc)) return {"ok": True, "record": rec} # IS 五行(P3 接线棒):Perold 四分解+价格移动参考行,读偏差日报 is_daily 最新行。 # diff 口径=is_daily 无 traded_value 列,契约①的 bp 归一(金额/当日成交额)无分母 # 可除 → 按金额直显(f"{v:+,.2f}元",正=逆风);腿数注在 live 行内补量级语境。 _IS_METRICS = ( ("延迟成本(落地差距)", "delay_cost"), ("冲击成本(落地差距)", "impact_cost"), ("机会成本(落地差距)", "opportunity_cost"), ("显性费用(落地差距)", "fee_total"), ("价格移动(市场漂移·参考)", "price_movement"), ) def _attribution_dir() -> str: return os.environ.get("SANGUO_ATTRIBUTION_DIR", os.path.join("data", "attribution")) def _aggregate_monthly_contribution(rows: list, window_months: int = 12) -> list: """月度长表 → 前端短表契约 {factorId, contributionPct} 聚合. 长表(attribution.py 产物)=每因子×月一行 {factor, month, contribution_ret}, 含 market/size/industry/specific 三桶行; 前端(AttributionResearch.vue P0 mock 契约)={factorId, contributionPct}(toFixed(1)). 09-27 契位修复: 直通长表首行无 contributionPct → 前端 toFixed 整页崩. 口径=近 window_months 月 Σ contribution_ret×100(滚动年窗; 全期口径 market +65% 量级会溢出前端横条); 不足 window_months 月(挑战者件)退化 全窗 Σ; 三桶与源同板(研判核心结构); 脏行(缺字段/类型错)跳过不炸. """ if not rows: return [] months = sorted({r["month"] for r in rows if isinstance(r, dict) and isinstance(r.get("month"), str)}) recent = set(months[-window_months:]) agg: dict[str, float] = {} for r in rows: factor, ret = r.get("factor"), r.get("contribution_ret") if (not isinstance(factor, str) or not isinstance(ret, (int, float)) or r.get("month") not in recent): continue agg[factor] = agg.get(factor, 0.0) + float(ret) return [{"factorId": f, "contributionPct": round(v * 100, 2)} for f, v in sorted(agg.items(), key=lambda kv: -kv[1])] def _read_latest_attribution_doc(variant: str = "incumbent") -> dict | None: """attribution_*.json 文件名降序最新一份原文(dict); challenger=子目录 (件名带因子名也命中 glob). 缺/坏 → None(warning 留痕,不炸).""" adir = _attribution_dir() if variant == "challenger": adir = os.path.join(adir, "challenger") try: names = sorted((n for n in os.listdir(adir) if n.startswith("attribution_") and n.endswith(".json")), reverse=True) if not names: return None with open(os.path.join(adir, names[0]), encoding="utf-8") as f: doc = json.load(f) return doc if isinstance(doc, dict) else None except (OSError, ValueError) as exc: # 缺目录/坏 JSON/编码错一并按缺数据处理 logger.warning("attribution JSON 读取失败,按缺数据处理返空: %s", exc) return None def _latest_factor_contribution(variant: str = "incumbent") -> list: """最新归因件 → 月度长表聚合为前端短表(09-27 契位修复, 见聚合函数注).""" doc = _read_latest_attribution_doc(variant) if not doc: return [] rows = doc.get("factorContribution") return _aggregate_monthly_contribution(rows if isinstance(rows, list) else []) def _latest_attribution_meta(variant: str = "incumbent") -> dict: """归因件元信息(前端口径条): 数据截至/窗末/覆盖月数; 缺件 → 空字段.""" doc = _read_latest_attribution_doc(variant) or {} rows = doc.get("factorContribution") months = sorted({r.get("month") for r in rows if isinstance(r, dict) and isinstance(r.get("month"), str)}) \ if isinstance(rows, list) else [] win = doc.get("window") if isinstance(doc.get("window"), dict) else {} return {"asOf": doc.get("as_of"), "windowEnd": win.get("end"), "months": len(months)} @router.get("/pipeline/attribution") def attribution(variant: Literal["incumbent", "challenger"] = "incumbent") -> dict: """归因日报(P3 接线棒):周报三指标真值+IS 五行读 latest_is(偏差日报) +factorContribution 读 attribution JSON 聚合; ?variant= incumbent=在位者 (quant12_v2a_szneul)/challenger=挑战者(#69 等权 v1,challenger/ 子目录); attributionMeta=归因件元信息(前端口径条: 数据截至/覆盖月数). 价格移动行带「参考」不进评价,金额字段缺失如实 "—"(宁缺毋假).""" from sanguo_portfolio import pipeline_store reg = _load_strategy_registry() th = pipeline_store.get_effective_thresholds(_pipeline_db()) rows: list[dict] = [] latest_week = "" for name in sorted(reg["strategies"]): e = reg["strategies"][name] if e.get("stage") not in ("live", "shadow", "paper"): continue w = pipeline_store.latest_weekly(_pipeline_db(), name) if not w or not w.get("week_start"): continue latest_week = max(latest_week, w["week_start"]) te, fr, vs = w.get("te_annual"), w.get("fill_rate"), w.get("vs_backtest") rows.append({"metric": f"{name}·跟踪误差(年化)", "live": f"{te:.2%}" if te is not None else "—", "backtest": f"阈 {th.shadow_te_annual_max:.2%}", "diff": f"{te - th.shadow_te_annual_max:+.2%}" if te is not None else "—"}) rows.append({"metric": f"{name}·成交率", "live": f"{fr:.1%}" if fr is not None else "—", "backtest": f"阈 {th.shadow_fill_rate_min:.0%}", "diff": f"{fr - th.shadow_fill_rate_min:+.1%}" if fr is not None else "—"}) rows.append({"metric": f"{name}·周收益差(vs回测)", "live": _fmt_vs(vs), "backtest": f"阈 ±{th.shadow_weekly_dev_max:.0%}/周", "diff": ("超" if vs is not None and abs(vs) > th.shadow_weekly_dev_max else "过") if vs is not None else "—"}) latest_is = pipeline_store.latest_is(_pipeline_db()) legs = (latest_is or {}).get("legs_count") for label, col in _IS_METRICS: v = (latest_is or {}).get(col) if v is None: rows.append({"metric": label, "live": "—", "backtest": "—", "diff": "—"}) else: note = f"({legs}腿)" if legs is not None else "" rows.append({"metric": label, "live": f"{v:,.2f}元{note}", "backtest": "—", "diff": f"{v:+,.2f}元"}) return {"period": latest_week or "周报待算", "liveVsBacktest": rows, "factorContribution": _latest_factor_contribution(variant), "attributionMeta": _latest_attribution_meta(variant)} # 影子 A/B 挑战者策略 id(#69: 等权 v1 vs 在位 v2a, 四周对照) _SHADOW_CHALLENGER_ID = "szneul_eq_reb63" _SHADOW_TOTAL_WEEKS = 4 # 活树清单(以策略为锚: 新策略钉住新合成层上位时在此登记一棵树—— # 版本树=「血统进化史+谁在位」, 无策略消费的单代族不进树, 进了也是孤节点) _COMPOSITE_FAMILIES = [ {"key": "composite_quant12", "label": "quant12 · 12 量价源", "note": "在位=szneul_reb63(quant12_v2a_szneul) · 挑战者=#69 影子(等权 v1)"}, ] @router.get("/pipeline/composite") def composite(family: str | None = None) -> dict: """合成层版本树真值(P1 接真, 09-27): 多族就绪, 当前一棵活树. ?family= 选族(缺省=清单首棵); families=可用族清单(前端切换器). 在位者 quant12_v2a=L1 权重档案 resolved(唯一事实源,代码零硬编码); 挑战者 quant12_v1=等权基线(#69 影子 A/B,_szneul 变体 09-28 起跑). 影子对照读挑战者周报(weeks_counted/te_annual)——首跑前无周报,如实 0/4 周+TE null+进行中(不造数); 满 4 周按影子 TE 阈值判达标/超阈. 未知族/因子侧坏档 → fail-soft items=[](看板不 500, warning 留痕). """ from sanguo_portfolio import pipeline_store fam = family or _COMPOSITE_FAMILIES[0]["key"] known = fam in {f["key"] for f in _COMPOSITE_FAMILIES} if not known: logger.warning("composite 未知族(可用=%s): %s", [f["key"] for f in _COMPOSITE_FAMILIES], fam) return {"family": fam, "families": _COMPOSITE_FAMILIES, "items": []} try: from sanguo_factor import composite_weighting as cw from sanguo_factor.composite_library import QUANT_SOURCES prof = cw.load_profile(str(cw._DEFAULT_PROFILE_PATH)) w2a = cw.resolve_weights(prof) incumbent = { "version": "quant12_v2a", "releasedAt": str(prof.metadata.get("created", "")), "note": f"{prof.name}(method={prof.method}; 档案 {prof.profile_id}; " "szneul_reb63 live 钉住)", "members": [{"factorId": src, "weight": round(w2a[src], 4)} for src, _d in QUANT_SOURCES], } challenger = { "version": "quant12_v1", "releasedAt": "2026-08-31", "note": "等权基线 v1.1(08-31 族分析固化 12 源); _szneul 变体=" "#69 影子 A/B 挑战者, 与 v2a 只差权重", "members": [{"factorId": src, "weight": round(1 / len(QUANT_SOURCES), 4)} for src, _d in QUANT_SOURCES], } except Exception as exc: # 因子域坏档/注册链断: 宁空勿崩(同归因页纪律) logger.warning("composite 版本树因子侧读取失败,按缺数据处理返空: %s", exc) return {"family": fam, "families": _COMPOSITE_FAMILIES, "items": []} weekly = pipeline_store.latest_weekly(_pipeline_db(), _SHADOW_CHALLENGER_ID) week = int(weekly.get("weeks_counted") or 0) if weekly else 0 te = weekly.get("te_annual") if weekly else None if week < _SHADOW_TOTAL_WEEKS or te is None: verdict = "进行中" else: threshold = pipeline_store.get_effective_thresholds( _pipeline_db()).shadow_te_annual_max verdict = "达标" if te <= threshold else "超阈" challenger["shadowCompare"] = {"versus": "quant12_v2a", "week": week, "totalWeeks": _SHADOW_TOTAL_WEEKS, "trackingError": te, "verdict": verdict} return {"family": fam, "families": _COMPOSITE_FAMILIES, "items": [challenger, incumbent]} @router.get("/pipeline/composite/challenge") def composite_challenge() -> dict: """三路对拍最新件(等权/ICIR/LGBM 影子实验,spec §4.4 决议 O 前置可視性). 读报告根 challenger_lgbm/ 子目录(challenger_lgbm.py 产物,月度链 stage6 每月覆盖);多件取 as_of 最大,坏件 skip 留痕.无件=空态(11-01 首考), 前端显示"暂无对拍数据"——不造数. """ import glob as _glob sub = os.path.join(_monthly_dir(), "challenger_lgbm") docs: list[tuple[str, str, dict]] = [] # (as_of, file, doc) for p in sorted(_glob.glob(os.path.join(sub, "*.json"))): try: with open(p, encoding="utf-8") as f: d = json.load(f) as_of = str(d.get("as_of") or "") if as_of: docs.append((as_of, os.path.basename(p), d)) except Exception as exc: # 坏件 skip: 对拍件非权威数据面,宁缺勿崩 logger.warning("challenge 对拍件坏件 skip %s: %s", p, exc) if not docs: return {"challenge": None, "available": [], "note": "暂无对拍数据——月度链 stage6 每月产出" "(11-01 首考);challenger 永不进生产(影子实验)"} docs.sort(key=lambda t: t[0]) # as_of 字典序=时间序 _, fname, doc = docs[-1] doc["file"] = fname return {"challenge": doc, "available": [f for _, f, _ in docs], "note": ""} # —— P4-1 假设卡片向导(决议 M:LLM 只做文字→结构,draft 不落库,确认才落) —— _DRAFT_SENTENCE_MAX = 500 _VERDICT_REASON_MAX = 500 # codex 终审 NOTE2:判死 reason 上限(与 draft sentence 同款) # P3-8: 502 detail 固定文案——LLMError 文本含上游响应体片段(端点/账号/配额 # 上下文),透传前端=信息泄露;真因只进服务端日志 _LLM_UPSTREAM_502 = "LLM 上游返回异常" # P3-7: save 端点 source 白名单(值域=向导默认/手工录入/分解器;前端现不传, # 默认 wizard)。有限集天然隐含长度上界,不再另设截断 _SAVED_SOURCES = ("wizard", "manual", "decomposer") def _card_to_camel(card: dict) -> dict: return { "title": card["title"], "logic": card["logic"], "expectedSign": card["expected_sign"], "falsifiable": card["falsifiable"], "dataNeeds": card["data_needs"], } @router.get("/pipeline/hypotheses") def hypotheses_list() -> dict: items = pipeline_store.list_hypotheses(_pipeline_db()) # 每卡附最新分解 job 视图(job 化后前端唯一轮询源;内存→台账→重启判 failed) for it in items: it["decomposeJob"] = decompose_jobs.view(it["id"], _pipeline_db()) # D7 件①1a:已产因子清单(registry hypothesis 血统反查)——卡上单因子名 # (首条代表作指针,历史包袱)由此退役;规模小直接内联,膨胀再改按需拉取 by_hyp: dict[str, list[dict[str, Any]]] = {} try: from sanguo_factor import version_registry as vr reg = vr.load_registry(_ensure_factor_runtime_registry()) for name, e in reg["factors"].items(): hyp = e.get("hypothesis") if not hyp: continue versions = e.get("versions") or [] expr = ((versions[0].get("params") or {}).get("expression") if versions else None) by_hyp.setdefault(hyp, []).append({ "name": name, "expression": expr, "origin": e.get("origin") or "manual", "createdAt": (versions[0].get("effective_from") if versions else None), "status": e.get("status"), "cause": e.get("cause")}) # 决议O③死因透传(死因子行tooltip) except Exception: # noqa: BLE001 registry 读失败不拖垮卡片列表 logger.warning("hypotheses factors 反查失败,如实空清单", exc_info=True) by_hyp = {} # 批次归属(验收追加:清单加批次分割线)——台账新→旧先见者胜, # 每条因子附 birth batch;台账查不到(超 limit/远古)=None 归「早期」 batch_of: dict[tuple[str, str], dict[str, Any]] = {} for hyp_id in by_hyp: for job in pipeline_store.list_decompose_jobs(_pipeline_db(), hyp_id): for c in job.get("registered") or []: batch_of.setdefault((hyp_id, c["name"]), {"jobId": job["jobId"], "startedAt": job["startedAt"]}) for it in items: fs = by_hyp.get(it["id"]) or [] for f in fs: f["batch"] = batch_of.get((it["id"], f["name"])) fs.sort(key=lambda f: (f["batch"] or {}).get("startedAt") or "", # noqa: E501 批次新→旧,无批次垫底 reverse=True) it["factors"] = fs return {"items": items} @router.post("/pipeline/hypotheses/draft") async def hypotheses_draft(body: dict) -> dict: sentence = str(body.get("sentence") or "").strip() if not sentence: raise HTTPException(400, "sentence 不能为空") if len(sentence) > _DRAFT_SENTENCE_MAX: raise HTTPException(400, f"sentence 超长(>{_DRAFT_SENTENCE_MAX} 字)") try: config = config_store.resolve_effective_config(_pipeline_db()) except LLMConfigError as e: raise HTTPException(503, str(e)) domains = _wizard_domains() client = LLMClient(config) try: raw = await client.chat_json(_build_wizard_messages(sentence, domains)) card = _normalize_draft(raw, domains) except LLMError as e: logger.warning("llm draft upstream error: %s", e) raise HTTPException(502, _LLM_UPSTREAM_502) from e except ValueError as e: raise HTTPException(502, f"LLM 草稿不合契约: {e}") return {"draft": _card_to_camel(card)} @router.post("/pipeline/hypotheses", status_code=201) def hypotheses_save(body: dict) -> dict: sentence = str(body.get("sentence") or "") if len(sentence) > _DRAFT_SENTENCE_MAX: raise HTTPException(400, f"sentence 超长(>{_DRAFT_SENTENCE_MAX} 字)") snake = { "title": body.get("title"), "logic": body.get("logic"), "expected_sign": body.get("expectedSign"), "falsifiable": body.get("falsifiable"), "data_needs": body.get("dataNeeds") or [], "sentence": body.get("sentence"), "source": body.get("source") or "wizard", } if snake["source"] not in _SAVED_SOURCES: # P3-7: 白名单,非法 400 raise HTTPException( 400, f"source 非法: 只认 {'/'.join(_SAVED_SOURCES)}") domains = _wizard_domains() try: card = _normalize_draft(snake, domains) # 保存再验(tickflow 双点复用) except ValueError as e: raise HTTPException(400, str(e)) card["sentence"] = snake["sentence"] card["source"] = snake["source"] row = pipeline_store.insert_hypothesis(_pipeline_db(), card) items = pipeline_store.list_hypotheses(_pipeline_db()) item = next(i for i in items if i["id"] == row["id"]) return {"item": item} _DECOMPOSE_ALLOWED_STATES = ("queued", "data_check", "building") def _ensure_factor_runtime_registry() -> str: """因子注册表运行副本自举(data/ 首触从 config/ 种子拷贝,strategy 同款).""" from sanguo_portfolio.strategy_registry import ensure_runtime_registry return ensure_runtime_registry( _factor_registry_path(), os.path.join("config", "factor_registry.yaml")) @router.post("/pipeline/hypotheses/{hyp_id}/decompose", status_code=202) async def hypotheses_decompose(hyp_id: str) -> dict: """卡片→候选因子 job 化(决议 M②,10-09 QA 范式):POST 立即返 job, 前端轮询卡片 decomposeJob 到终态——浏览器不再扛分钟级同步 HTTP (响应孤儿化=10-09 全部「静默卡死」同源,调超时永远治不了). 同卡在飞 409;卡片不存在/状态不可分解/LLM 未配三类前置校验同步 fail-fast(不空起 job).worker 体内锁内重读卡片态(P2-5 竞态守卫). """ if decompose_jobs.is_running(hyp_id): raise HTTPException(409, "该卡分解进行中(后台任务),刷新可看进度") items = pipeline_store.list_hypotheses(_pipeline_db()) card = next((i for i in items if i["id"] == hyp_id), None) if card is None: raise HTTPException(404, "卡片不存在") if card["state"] == "graveyard": raise HTTPException(422, "终态卡(墓园)不可分解——重推想法=开新卡(决议O)") if card["state"] not in _DECOMPOSE_ALLOWED_STATES: raise HTTPException(409, f"卡片状态 {card['state']} 不可分解") try: config = config_store.resolve_effective_config(_pipeline_db()) except LLMConfigError as e: raise HTTPException(503, str(e)) async def worker(job_id: str) -> dict: return await _decompose_worker(hyp_id, config, job_id) job = decompose_jobs.start(hyp_id, worker, _pipeline_db()) return {"job": job} @router.get("/pipeline/hypotheses/{hyp_id}/decompose-jobs") def hypotheses_decompose_jobs(hyp_id: str) -> dict: """该卡分解批次历史(D7 件①1b):台账新→旧,卡片折叠区按需拉取. 端点形状照抄 QA GET /api/v1/mining/tasks/list(app.py:553-557); 内存只有最新一条,历史=台账单一真相(进度过程态不入此列). """ items = pipeline_store.list_hypotheses(_pipeline_db()) if not any(i["id"] == hyp_id for i in items): raise HTTPException(404, "卡片不存在") return {"items": pipeline_store.list_decompose_jobs( _pipeline_db(), hyp_id, limit=20)} async def _decompose_worker(hyp_id: str, config, job_id: str | None = None) -> dict: """分解后台体:锁内重读卡片→LLM+硬门循环→通过者注册挂血统→卡片转 building. 零特权:分解产物与手写因子同注册表同求值器,月度批评经动态注册链 (monthly_batch resolve_with_yaml)自动接管.失败者结构化留 job 不落库. P2-6:全程持 _DECOMPOSE_LOCK 进程内串行化(同卡并发 decompose 的 registry 整文件覆盖丢更新窗口闭合;串行后二次请求重读卡片态与 existing_names,由 P2-5 状态条件+查重自然兜住). """ # graduate 是 sync 端点(FastAPI 线程池执行)可直接 with 持锁;本函数是 # async,裸 with 会在 await LLM 挂起期间把事件循环线程堵死在 acquire # 上(第二个并发 decompose 到达即全服死锁),故经 executor 阻塞获取, # 不占事件循环线程. await asyncio.get_running_loop().run_in_executor( None, _DECOMPOSE_LOCK.acquire) try: items = pipeline_store.list_hypotheses(_pipeline_db()) card = next((i for i in items if i["id"] == hyp_id), None) if card is None: raise decompose_jobs.JobFailure("卡片不存在(排队期间被删?)") if card["state"] not in _DECOMPOSE_ALLOWED_STATES: raise decompose_jobs.JobFailure( f"卡片状态已变({card['state']}),不可分解") import sanguo_factor # noqa: F401 导入即注册全部手写因子库(零特权比对面) from sanguo_factor import version_registry as vr from sanguo_factor.factor_guard import WIRED_DOMAINS, derive_source from sanguo_factor.registry import list_factors rt = _ensure_factor_runtime_registry() reg = vr.load_registry(rt) library_exprs = {f["name"]: f["expression"] for f in list_factors()} for name, entry in reg["factors"].items(): versions = entry.get("versions") or [] expr = ((versions[-1].get("params") or {}).get("expression") if versions else None) if expr: library_exprs.setdefault(name, expr) existing_names = set(reg["factors"]) | set(library_exprs) unwired = tuple(d for d in (card.get("dataNeeds") or []) if d not in WIRED_DOMAINS) result = await _run_decompose( LLMClient(config), card, existing_names=existing_names, library_exprs=library_exprs, unwired_domains=unwired, on_progress=(lambda **kw: decompose_jobs.update_progress( job_id, **kw)) if job_id else None) registered: list[dict[str, Any]] = [] if result["passed"]: today = date.today().isoformat() for cand in result["passed"]: src = derive_source(cand["expression"]) # 已过门,必为 str # 仅建条纪律(codex review):已存在条目(同名同假设幂等重注册) # 不重写 similarity——与 description/origin 落盘纪律同款 is_new = cand["name"] not in reg["factors"] vr.upsert_factor(reg, cand["name"], hypothesis=hyp_id, status="incubating", origin="decomposer", description=cand.get("description")) vr.add_version(reg, cand["name"], commit="decompose", params={"expression": cand["expression"], "source": src, "origin": "decomposer", "justification": cand["justification"]}, effective_from=today) # AST 防换皮提示(spec §4.2 2026-10-10):对既有池比对,≥0.6 存 # entry.similarity——硬门(dup_subtree)外的部分同构仍可注册, # 提示非硬拒,判定人做.池=注册表条目当前表达式(versions[-1]). if is_new: from sanguo_factor import expression_match as em _cands = {n: str(((v.get("versions") or [{}])[-1] .get("params") or {}).get("expression") or "") for n, v in reg["factors"].items() if n != cand["name"]} _sim = em.top_similar(str(cand["expression"]), _cands) if _sim: reg["factors"][cand["name"]]["similarity"] = _sim registered.append({"name": cand["name"], "source": src, "expression": cand["expression"], "description": cand.get("description", ""), "justification": cand["justification"]}) vr.save_registry(rt, reg) if card["state"] != "building": # building 态二次分解=纯增量注册(D3 矩阵无自迁移边) pipeline_store.update_hypothesis_state( _pipeline_db(), hyp_id, "building", factor_id=registered[0]["name"]) return {"registered": registered, "failed": result["failed"], "rounds": result["rounds"]} finally: _DECOMPOSE_LOCK.release() # —— P5.1 LLM 配置面(spec §13):三层解析的管理端点 —— @router.get("/pipeline/llm/config") def llm_config_get() -> dict: """设置页 LLM 卡:打码视图+逐字段来源(db/env/default)+审计行.""" return {**config_store.get_masked_config(_pipeline_db()), "audit": config_store.get_llm_audit(_pipeline_db())} @router.put("/pipeline/llm/config") def llm_config_put(body: dict) -> dict: """写 DB 覆盖层;含 '*' 的 api_key(掩码回显)静默跳过;空串=清键回落.""" values = body.get("values") changed_by = body.get("changed_by") or "console" if not isinstance(values, dict): raise HTTPException(422, "body 需要 {values: {...}, changed_by}") try: changed = config_store.set_llm_config( _pipeline_db(), values, changed_by) except ValueError as e: raise HTTPException(422, str(e)) from e return {"ok": True, "changed": changed} @router.post("/pipeline/llm/test") async def llm_test() -> dict: """测试连接:按当前生效配置(三层解析后)单发最小探测,返回 ok/latency/error.""" try: config = config_store.resolve_effective_config(_pipeline_db()) except LLMConfigError as e: raise HTTPException(503, str(e)) return await probe_connection(config) # —— P5.2/P5.6 设置页只读探针(spec §13):口令类 env-only,不给 UI 编辑 —— @router.get("/pipeline/env-status") def env_status() -> dict: """环境依赖状态卡:configured 布尔+指引;缺什么前置可见(治毕业 503 暗坑).""" def _env_set(name: str) -> bool: return bool((os.environ.get(name) or "").strip()) try: config_store.resolve_effective_config(_pipeline_db()) llm_ok = True except LLMConfigError: llm_ok = False items = [ {"key": "SANGUO_PROMOTION_PASSPHRASE", "label": "毕业口令", "configured": _env_set("SANGUO_PROMOTION_PASSPHRASE"), "hint": "env-only(安全边界,不给 UI 编辑);未配置时毕业提交 503"}, {"key": "llm", "label": "LLM 连接(draft/分解)", "configured": llm_ok, "hint": "DB>env>默认三层;设置页 LLM 卡可配,env 兜底首配"}, {"key": "SANGUO_GITEA_TOKEN", "label": "Gitea token(毕业开 issue)", "configured": _env_set("SANGUO_GITEA_TOKEN"), "hint": "env-only;缺=毕业转移仍落档,issue 需人工补开"}, {"key": "factor_registry", "label": "因子注册表(运行副本)", "configured": os.path.exists(_factor_registry_path()), "hint": f"{_factor_registry_path()};首触自举,缺=月度批未跑过"}, {"key": "strategy_registry", "label": "策略注册表", "configured": os.path.exists(_registry_path()), "hint": f"{_registry_path()};梯子/毕业读此档"}, ] return {"items": items} @router.get("/pipeline/monthly-status") def monthly_status() -> dict: """月度批状态卡(P5.6 只读):cron 日志最后一段 start→done 各段 rc. 日志候选两位:monthly_dir 内(VPS wrapper 落 OUT_DIR)与上级(NAS runner 落 BASE 根);都缺=本机从未跑过,logFound=False 如实展示。 只展示不触发(手动触发按钮=P5.5 明确不做,40min 重批误触代价高)。 """ md = _monthly_dir() log = next((p for p in ( os.path.join(md, "factor_monthly_cron.log"), os.path.join(os.path.dirname(os.path.abspath(md)), "factor_monthly_cron.log")) if os.path.exists(p)), None) out: dict[str, Any] = {"logFound": False, "startedAt": None, "asOf": None, "stages": [], "finalRc": None} if not log: return out out["logFound"] = True try: with open(log, encoding="utf-8", errors="replace") as f: f.seek(0, os.SEEK_END) f.seek(max(0, f.tell() - 65536)) # 只看尾部 64KB,月度日志极小 lines = f.read().splitlines() except OSError: return out start_i = max((i for i, l in enumerate(lines) if "factor-monthly start" in l), default=-1) if start_i < 0: return out m = re.search(r"=== (\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2})", lines[start_i]) out["startedAt"] = m.group(1) if m else None m = re.search(r"as_of=([\d-]+)", lines[start_i]) out["asOf"] = m.group(1) if m else None for line in lines[start_i + 1:]: if "factor-monthly start" in line: break # 防御:尾部截断致旧段混入时只认最后一段 m = re.search(r"stage(\d+)\s+(\S+)\s+exit=(\d+)", line) if m: out["stages"].append({"n": int(m.group(1)), "name": m.group(2), "rc": int(m.group(3))}) m = re.search(r"factor-monthly done final_rc=(\d+)", line) if m: out["finalRc"] = int(m.group(1)) return out