1581 lines
75 KiB
Python
1581 lines
75 KiB
Python
# 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
|