fix(pipeline): 假设状态机条件 UPDATE+decompose 按卡互斥——update_hypothesis_state 带 allowed_from 守卫(单条条件 UPDATE+rowcount 二次分因,graveyard 终态竞态不可迁出)+decompose 全程 _DECOMPOSE_LOCK(仿 _GRADUATE_LOCK,跨存储两步写丢更新窗口闭合) (审计 P2-5/P2-6) [vps] [no-doc]

This commit is contained in:
2026-10-05 15:45:51 +08:00
parent 93d7acef50
commit a13ce79b93
4 changed files with 151 additions and 73 deletions
+72 -57
View File
@@ -7,6 +7,7 @@ P2 零新增跨机推送。promotion 执行(起实盘实例)仍是用户手动
"""
from __future__ import annotations
import asyncio
import hmac
import json
import logging
@@ -36,6 +37,7 @@ logger = logging.getLogger(__name__)
# graduate 全程进程内互斥:FastAPI sync 路由跑线程池,防并发双发双转移双事件
_GRADUATE_LOCK = threading.Lock()
_DECOMPOSE_LOCK = threading.Lock() # P2-6:decompose 全程持锁(仿 _GRADUATE_LOCK)
def _registry_path() -> str:
@@ -889,69 +891,82 @@ async def hypotheses_decompose(hyp_id: str) -> dict:
零特权:分解产物与手写因子同注册表同求值器,月度批评经动态注册链
(monthly_batch resolve_with_yaml)自动接管.失败者结构化返回不落库.
P2-6:全程持 _DECOMPOSE_LOCK 进程内串行化(同卡并发 decompose 的
registry 整文件覆盖丢更新窗口闭合;串行后二次请求重读卡片态与
existing_names,由 P2-5 状态条件+查重自然兜住).
"""
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"] not in _DECOMPOSE_ALLOWED_STATES:
raise HTTPException(409, f"卡片状态 {card['state']} 不可分解")
# graduate 是 sync 端点(FastAPI 线程池执行)可直接 with 持锁;本端点是
# async,裸 with 会在 await LLM 挂起期间把事件循环线程堵死在 acquire
# 上(第二个并发 decompose 到达即全服死锁),故经 executor 阻塞获取,
# 不占事件循环线程.
await asyncio.get_running_loop().run_in_executor(
None, _DECOMPOSE_LOCK.acquire)
try:
config = config_store.resolve_effective_config(_pipeline_db())
except LLMConfigError as e:
raise HTTPException(503, str(e))
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"] 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))
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
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)
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)
try:
result = await _run_decompose(
LLMClient(config), card, existing_names=existing_names,
library_exprs=library_exprs, unwired_domains=unwired)
except LLMError as e:
raise HTTPException(502, str(e))
except ValueError as e:
raise HTTPException(502, f"LLM 草稿不合契约: {e}")
try:
result = await _run_decompose(
LLMClient(config), card, existing_names=existing_names,
library_exprs=library_exprs, unwired_domains=unwired)
except LLMError as e:
raise HTTPException(502, str(e))
except ValueError as e:
raise HTTPException(502, f"LLM 草稿不合契约: {e}")
registered: list[dict[str, Any]] = []
if result["passed"]:
today = date.today().isoformat()
for cand in result["passed"]:
src = derive_source(cand["expression"]) # 已过门,必为 str
vr.upsert_factor(reg, cand["name"], hypothesis=hyp_id,
status="incubating")
vr.add_version(reg, cand["name"], commit="decompose",
params={"expression": cand["expression"],
"source": src, "origin": "decomposer",
"justification": cand["justification"]},
effective_from=today)
registered.append({"name": cand["name"], "source": src,
"expression": cand["expression"],
"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"])
items = pipeline_store.list_hypotheses(_pipeline_db())
item = next(i for i in items if i["id"] == hyp_id)
return {"registered": registered, "failed": result["failed"],
"rounds": result["rounds"], "item": item}
registered: list[dict[str, Any]] = []
if result["passed"]:
today = date.today().isoformat()
for cand in result["passed"]:
src = derive_source(cand["expression"]) # 已过门,必为 str
vr.upsert_factor(reg, cand["name"], hypothesis=hyp_id,
status="incubating")
vr.add_version(reg, cand["name"], commit="decompose",
params={"expression": cand["expression"],
"source": src, "origin": "decomposer",
"justification": cand["justification"]},
effective_from=today)
registered.append({"name": cand["name"], "source": src,
"expression": cand["expression"],
"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"])
items = pipeline_store.list_hypotheses(_pipeline_db())
item = next(i for i in items if i["id"] == hyp_id)
return {"registered": registered, "failed": result["failed"],
"rounds": result["rounds"], "item": item}
finally:
_DECOMPOSE_LOCK.release()
# —— P5.1 LLM 配置面(spec §13):三层解析的管理端点 ——
+34 -14
View File
@@ -263,27 +263,47 @@ def list_hypotheses(path: str) -> list[dict[str, Any]]:
def update_hypothesis_state(path: str, hyp_id: str, state: str,
death_reason: str | None = None,
factor_id: str | None = None) -> bool:
factor_id: str | None = None, *,
allowed_from: tuple[str, ...] | None = None
) -> bool | None:
"""状态转移+墓园回写(P4-2 分解器消费).
D3 转移矩阵校验:非法迁移 raise ValueError;未知 id→False;非法 state→
ValueError.COALESCE 语义保留(矩阵禁复活后无清空需求,注:未来开复活口
须同时引入显式清空 ""=置空).
返回三分:True=成功,False=状态不符被拒(仅 allowed_from 显式传入时),
None=卡片不存在(旧版 False 与「被拒」歧义,P2-5 改版).
D3 转移矩阵校验:非法迁移 raise ValueError;非法 state→ValueError.
COALESCE 语义保留(矩阵禁复活后无清空需求,注:未来开复活口须同时引入
显式清空 ""=置空).
P2-5 竞态守卫:单条条件 UPDATE(WHERE id AND state IN 允许源态)原子完成
判读+写,末写者胜窗口闭合(graveyard 终态竞态不可迁出);rowcount==0 时
二次 SELECT 分因.不传 allowed_from 则按 D3 矩阵推导允许源态,状态不符
维持 raise ValueError(既有调用点零行为变化).
"""
if state not in _HYP_STATES:
raise ValueError(f"非法卡片状态: {state!r}")
if allowed_from is None:
allowed = tuple(s for s, targets in _VALID_HYP_TRANSITIONS.items()
if state in targets)
else:
allowed = tuple(allowed_from)
placeholders = ",".join("?" * len(allowed))
with _connect(path) as conn:
cur = conn.execute(
"UPDATE hypotheses SET state=?,"
"death_reason=COALESCE(?,death_reason),"
"factor_id=COALESCE(?,factor_id),"
"updated_at=datetime('now','localtime') "
f"WHERE id=? AND state IN ({placeholders})",
(state, death_reason, factor_id, hyp_id, *allowed))
if cur.rowcount:
return True
row = conn.execute(
"SELECT state FROM hypotheses WHERE id=?", (hyp_id,)).fetchone()
if row is None:
return None
if allowed_from is not None:
return False
if state not in _VALID_HYP_TRANSITIONS[row["state"]]:
raise ValueError(
f"非法卡片状态迁移: {row['state']} → {state}"
f"(合法目标: {sorted(_VALID_HYP_TRANSITIONS[row['state']]) or '无,终态'})")
conn.execute(
"UPDATE hypotheses SET state=?, death_reason=COALESCE(?,death_reason),"
"factor_id=COALESCE(?,factor_id),"
"updated_at=datetime('now','localtime') WHERE id=?",
(state, death_reason, factor_id, hyp_id))
return True
raise ValueError(
f"非法卡片状态迁移: {row['state']} → {state}"
f"(合法目标: {sorted(_VALID_HYP_TRANSITIONS[row['state']]) or '无,终态'})")
@@ -237,3 +237,14 @@ class TestDecompose:
assert body["registered"][0]["name"] == "llm_vol5"
factors = _yaml.safe_load(reg.read_text())["factors"]
assert "llm_mom20" in factors and "llm_vol5" in factors # 两次注册都在
def test_decompose_lock_exists(self):
# P2-6:decompose 全程持 _DECOMPOSE_LOCK(仿 _GRADUATE_LOCK,进程内
# 串行化同卡并发,registry 整文件覆盖丢更新窗口闭合)。并发真测不做
# ——线程池+async 端点组合测不划算;锁内注册+转态的行为回归由本类
# 既有各例钉死,此处只钉锁存在防回退。
import threading
from sanguo_api import routes_pipeline as rp
assert isinstance(rp._DECOMPOSE_LOCK, type(threading.Lock()))
assert isinstance(rp._GRADUATE_LOCK, type(threading.Lock()))
@@ -50,9 +50,10 @@ def test_update_state_graveyard_writeback(tmp_path):
assert item["factorId"] == "fa_x"
def test_update_state_unknown_id_returns_false(tmp_path):
def test_update_state_unknown_id_returns_none(tmp_path):
# P2-5 接口语义三分:None=卡片不存在(旧 False 与「被拒」歧义,故改版)
db = str(tmp_path / "pipeline.db")
assert ps.update_hypothesis_state(db, "hyp-nope", "graveyard") is False
assert ps.update_hypothesis_state(db, "hyp-nope", "graveyard") is None
import pytest
@@ -90,3 +91,34 @@ def test_unknown_state_still_valueerror(tmp_path):
hyp = _card(db)
with pytest.raises(ValueError):
ps.update_hypothesis_state(db, hyp, "flying")
# —— P2-5 条件 UPDATE(P2 排期批):allowed_from 守卫+三分返回语义 ——
# 语义:True=成功,False=状态不符被拒,None=卡片不存在;不传 allowed_from
# 则按 D3 矩阵推导允许源态,状态不符维持 raise ValueError(现行为).
def test_update_state_allowed_from_success(tmp_path):
db = str(tmp_path / "p.db")
hyp = _card(db)
assert ps.update_hypothesis_state(
db, hyp, "building", allowed_from=("queued", "data_check")) is True
assert ps.list_hypotheses(db)[0]["state"] == "building"
def test_update_state_allowed_from_graveyard_rejected(tmp_path):
# P2-5 核心场景:graveyard 终态卡带守卫 update→被拒(False)且状态/死因
# 原样——竞态下末写者无法把墓园卡迁出(D3 禁复活在并发下成立)
db = str(tmp_path / "p.db")
hyp = _card(db)
ps.update_hypothesis_state(db, hyp, "graveyard", death_reason="证伪")
assert ps.update_hypothesis_state(
db, hyp, "building", allowed_from=("queued", "data_check")) is False
item = ps.list_hypotheses(db)[0]
assert item["state"] == "graveyard"
assert item["deathReason"] == "证伪"
def test_update_state_allowed_from_unknown_id_none(tmp_path):
db = str(tmp_path / "p.db")
assert ps.update_hypothesis_state(
db, "hyp-nope", "building", allowed_from=("queued",)) is None