将角色与 Skill 从 .claude 迁入 .agent,移除 Claude CLI 运行时并接入固定 Opus 角色 profile、完整 schema、预算 deadline、raw 与回执证据链。 同步拆分 Skill 职责、复利 lesson、Gate 回放、Dashboard 人审入口、数据库登记和机械门禁;候选设计正文不包含在本提交中。
141 lines
5.9 KiB
Python
141 lines
5.9 KiB
Python
#!/usr/bin/env python3
|
||
"""候选 CAS 状态链的 PostgreSQL 持久化存储(生产用)。
|
||
|
||
实现 run_writer_pipeline 的 CasStateStore 协议,替换仅供实验的 InMemoryCasStateStore:
|
||
create/transition/start_next/latest 全部走短事务条件 UPDATE,竞争时返回 None(由编排层
|
||
收敛为 CAS_CONFLICT),不再因进程退出丢失状态。DB 触发器(db/ddl/109)另行锁死迁移
|
||
方向、revision 单调 +1 与身份不可变,调用方代码缺陷也无法把链改成非法形状。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import pathlib
|
||
import sys
|
||
from typing import Any, Callable
|
||
|
||
SCRIPT_DIR = pathlib.Path(__file__).resolve().parent
|
||
if str(SCRIPT_DIR) not in sys.path:
|
||
sys.path.insert(0, str(SCRIPT_DIR))
|
||
|
||
from muse_db import connect as _default_connect # noqa: E402
|
||
from run_writer_pipeline import CasToken, PipelineError # noqa: E402
|
||
|
||
|
||
class PostgresCasStateStore:
|
||
"""生产运行的持久 CAS 状态存储。
|
||
|
||
参数只影响 create 时登记的可观察归属(work_id/target_chapter)与审计 updater;
|
||
迁移语义完全由 CasToken 期望值与 DB 约束决定。connect 可注入以便测试替换。
|
||
"""
|
||
|
||
# 与 InMemoryCasStateStore 同形的方向闭集:本地先行拒绝非法方向(返回 None,
|
||
# 与内存实现行为一致),DB 触发器仍是最终权威。
|
||
_ALLOWED = {
|
||
"DRAFT": frozenset({"CHECKING"}),
|
||
"CHECKING": frozenset({"PASSED", "REJECTED"}),
|
||
"PASSED": frozenset(),
|
||
"REJECTED": frozenset(),
|
||
}
|
||
|
||
def __init__(
|
||
self,
|
||
*,
|
||
work_id: int | None = None,
|
||
target_chapter: int | None = None,
|
||
creator: str = "continuation",
|
||
connect: Callable[..., Any] | None = None,
|
||
) -> None:
|
||
self._work_id = work_id
|
||
self._target_chapter = target_chapter
|
||
self._creator = creator
|
||
self._connect = connect or _default_connect
|
||
|
||
def create(self, run_id: str, *, attempt: int, candidate_version: int) -> CasToken:
|
||
"""为未存在的 run 建链(DRAFT/revision=1);已存在即 CAS 冲突,失败关闭。"""
|
||
|
||
with self._connect() as conn:
|
||
try:
|
||
cur = conn.execute(
|
||
"INSERT INTO example_candidate_cas(run_id, work_id, target_chapter, attempt, "
|
||
"candidate_version, state, revision, creator, updater) "
|
||
"VALUES (%s,%s,%s,%s,%s,'DRAFT',1,%s,%s) ON CONFLICT (run_id) DO NOTHING",
|
||
(run_id, self._work_id, self._target_chapter, attempt, candidate_version,
|
||
self._creator, self._creator),
|
||
)
|
||
if cur.rowcount != 1:
|
||
raise PipelineError("CAS_CONFLICT", f"run 已存在 CAS 链,不能重复创建初始状态: {run_id}")
|
||
conn.commit()
|
||
except Exception:
|
||
conn.rollback()
|
||
raise
|
||
return CasToken(run_id, attempt, candidate_version, "DRAFT", 1)
|
||
|
||
def transition(self, expected: CasToken, target_state: str) -> CasToken | None:
|
||
"""仅当链上最新状态与期望 token 完全一致时迁移;竞争、已迁移或方向非法返回 None。"""
|
||
|
||
if target_state not in self._ALLOWED.get(expected.state, frozenset()):
|
||
return None
|
||
with self._connect() as conn:
|
||
try:
|
||
cur = conn.execute(
|
||
"UPDATE example_candidate_cas SET state=%s, revision=revision+1, updater=%s "
|
||
"WHERE run_id=%s AND revision=%s AND state=%s AND attempt=%s "
|
||
"AND candidate_version=%s AND deleted=false",
|
||
(target_state, self._creator, expected.run_id, expected.revision,
|
||
expected.state, expected.attempt, expected.candidate_version),
|
||
)
|
||
if cur.rowcount != 1:
|
||
return None
|
||
conn.commit()
|
||
except Exception:
|
||
conn.rollback()
|
||
raise
|
||
return CasToken(
|
||
expected.run_id, expected.attempt, expected.candidate_version,
|
||
target_state, expected.revision + 1,
|
||
)
|
||
|
||
def start_next(
|
||
self, expected: CasToken, *, attempt: int, candidate_version: int
|
||
) -> CasToken | None:
|
||
"""只允许从最新 REJECTED 开严格递增的新一轮 DRAFT;否则返回 None。"""
|
||
|
||
if (
|
||
expected.state != "REJECTED"
|
||
or attempt <= expected.attempt
|
||
or candidate_version <= expected.candidate_version
|
||
):
|
||
return None
|
||
with self._connect() as conn:
|
||
try:
|
||
cur = conn.execute(
|
||
"UPDATE example_candidate_cas SET state='DRAFT', attempt=%s, candidate_version=%s, "
|
||
"revision=revision+1, updater=%s WHERE run_id=%s AND revision=%s AND state='REJECTED' "
|
||
"AND attempt=%s AND candidate_version=%s AND deleted=false",
|
||
(attempt, candidate_version, self._creator, expected.run_id,
|
||
expected.revision, expected.attempt, expected.candidate_version),
|
||
)
|
||
if cur.rowcount != 1:
|
||
return None
|
||
conn.commit()
|
||
except Exception:
|
||
conn.rollback()
|
||
raise
|
||
return CasToken(expected.run_id, attempt, candidate_version, "DRAFT", expected.revision + 1)
|
||
|
||
def latest(self, run_id: str) -> CasToken | None:
|
||
"""回读链的最新不可变状态;不存在或已软删返回 None。"""
|
||
|
||
with self._connect(readonly=True) as conn:
|
||
row = conn.execute(
|
||
"SELECT run_id, attempt, candidate_version, state, revision "
|
||
"FROM example_candidate_cas WHERE run_id=%s AND deleted=false",
|
||
(run_id,),
|
||
).fetchone()
|
||
if row is None:
|
||
return None
|
||
return CasToken(row[0], row[1], row[2], row[3], row[4])
|
||
|
||
|
||
__all__ = ["PostgresCasStateStore"]
|