muse-agent-example/tests/集成/test_正式变更事务.py
zizi a6f4f14a81 W09 正式变更、审阅、CAS 与原子提交:S01 正式变更、作者审阅、CAS 并发与跨模块原子提交。
按 R2 串行阶段整理提交;包内文件为该阶段交付(含后续小增量),状态以工作包清单为准。
2026-09-10 19:25:41 +08:00

333 lines
14 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""S01共用机制的真实PG验收;参与者使用合成表,不冒充正文/事实/章后主链。"""
from concurrent.futures import ThreadPoolExecutor
from dataclasses import dataclass, replace
from threading import Barrier
import pytest
from muse.共享.标识 import 任务标识
from muse.共享.调用身份 import 内容用途, 用途, 调用身份
from muse.正式变更.接口 import (
作者动作,
依赖引用,
修改影响,
参与操作,
参与者目录,
变更命令,
变更错误,
固定哈希,
审阅内容,
提交计划,
正式变更服务,
)
pytestmark = pytest.mark.数据库
@dataclass(frozen=True)
class 合成编辑:
内容: str
写后失败: bool = False
class 合成入口:
身份 = "synthetic-edit"
请求类型 = 合成编辑
def 查看(self, 连, 身份, 命令):
目标 = 连.execute(
"SELECT owner,revision,content FROM s01_test_document WHERE id=%s", (命令.目标,)
).fetchone()
if 目标 is None or 目标[0] != 身份.作者:
raise 变更错误("SCOPE_DENIED", "合成目标不属于当前作者")
候选 = 连.execute(
"SELECT revision,content,check_state FROM s01_test_candidate WHERE id=%s",
("candidate-1",),
).fetchone()
来源 = 连.execute(
"SELECT revision FROM s01_test_source WHERE id=%s", ("source-1",)
).fetchone()[0]
有候选 = 命令.动作 is not 作者动作.人工保存
依赖 = [
依赖引用("target", 命令.目标, str(目标[1])),
依赖引用("source", "source-1", str(来源)),
]
if 有候选:
依赖.append(依赖引用("candidate", "candidate-1", str(候选[0])))
return 审阅内容(
命令.目标,
目标[1],
"candidate-1" if 有候选 else None,
候选[0] if 有候选 else None,
固定哈希(候选[1]) if 有候选 else None,
固定哈希([目标[2], 命令.业务请求.内容]),
None,
None,
固定哈希(来源),
("replace-content",),
tuple(依赖),
)
def 准备(self, 连, 身份, 命令, 当前):
if 命令.动作 in {作者动作.拒绝, 作者动作.暂缓}:
return 提交计划(())
if 命令.动作 is 作者动作.采纳:
行 = 连.execute(
"SELECT content,check_state FROM s01_test_candidate WHERE id='candidate-1'"
).fetchone()
if 行[1] != "ready":
raise 变更错误("PRECONDITION_PENDING", "候选存在未处理检查问题")
if 行[0] != 命令.业务请求.内容:
raise 变更错误("REVIEW_MISMATCH", "只能采纳被审候选")
return 提交计划(
(参与操作("write", 命令.业务请求),),
(修改影响(命令.目标, "document_changed", str(当前.数据版本), str(当前.数据版本 + 1)),),
)
class 合成写入:
身份 = "write"
请求类型 = 合成编辑
def 提交(self, 连, 请求):
行 = 连.execute(
"UPDATE s01_test_document SET revision=revision+1,content=%s "
"WHERE id='document-1' RETURNING revision",
(请求.内容,),
).fetchone()
连.execute("INSERT INTO s01_test_audit (revision) VALUES (%s)", (行[0],))
if 请求.写后失败:
raise RuntimeError("合成参与者在业务写入后失败")
return {"revision": 行[0]}
class 合成依赖:
def __init__(self, 种类):
self.种类 = 种类
def 锁定当前(self, 连, 对象ID):
if self.种类 == "target":
行 = 连.execute(
"SELECT revision FROM s01_test_document WHERE id=%s FOR UPDATE", (对象ID,)
).fetchone()
elif self.种类 == "candidate":
行 = 连.execute(
"SELECT revision FROM s01_test_candidate WHERE id=%s FOR UPDATE", (对象ID,)
).fetchone()
else:
行 = 连.execute(
"SELECT revision FROM s01_test_source WHERE id=%s FOR SHARE", (对象ID,)
).fetchone()
if 行 is None:
raise 变更错误("SOURCE_STALE", "合成当前指针已删除")
return str(行[0])
@pytest.fixture
def 变更环境(应用测试库):
with 应用测试库[用途.维护].连接() as 连:
连.execute("""CREATE TABLE s01_test_document(
id text PRIMARY KEY, owner text,revision int,content text);
CREATE TABLE s01_test_candidate(
id text PRIMARY KEY,revision int,content text,check_state text);
CREATE TABLE s01_test_source(id text PRIMARY KEY,revision int);
CREATE TABLE s01_test_audit(revision int);
INSERT INTO s01_test_document VALUES ('document-1','author',1,'原内容');
INSERT INTO s01_test_candidate VALUES ('candidate-1',1,'候选内容','ready');
INSERT INTO s01_test_source VALUES ('source-1',1);""")
目录 = 参与者目录()
目录.登记入口(合成入口())
目录.登记参与者(合成写入())
for kind in ("target", "candidate", "source"):
目录.登记依赖(合成依赖(kind))
服务 = 正式变更服务(应用测试库[用途.生产], 目录)
身份 = 调用身份("author", 任务标识("author-operation"), 用途.生产, 内容用途.生成)
命令 = 变更命令(
"command-1", "synthetic-edit", "document-1", 作者动作.采纳, 1, 合成编辑("候选内容")
)
审阅 = 服务.打开审阅(身份, 命令)
命令 = replace(
命令, 审阅ID=审阅["review_id"], 审阅哈希=审阅["review_hash"], 批准变更=("replace-content",)
)
return 服务, 身份, 命令, 应用测试库
def test_审阅不是批准且提交回执幂等__a91001(变更环境):
服务, 身份, 命令, 库 = 变更环境
with 库[用途.生产].连接() as 连:
assert 连.execute("SELECT count(*) FROM muse_candidate_decision").fetchone()[0] == 0
assert 连.execute("SELECT revision FROM s01_test_document").fetchone()[0] == 1
回执 = 服务.提交变更(身份, 命令)
assert 回执["author_review_id"] == 命令.审阅ID
assert 回执["basis"]["候选版本"] == 1 and 回执["basis"]["数据版本"] == 1
assert 回执 == 服务.提交变更(身份, 命令) == 服务.读取回执(身份, 命令.命令ID)
with pytest.raises(变更错误) as 捕获:
服务.提交变更(身份, replace(命令, 批准变更=()))
assert 捕获.value.错误码 == "COMMAND_PAYLOAD_MISMATCH"
with 库[用途.生产].连接() as 连:
assert 连.execute("SELECT revision,content FROM s01_test_document").fetchone() == (
2,
"候选内容",
)
assert 连.execute("SELECT count(*) FROM s01_test_audit").fetchone()[0] == 1
assert 连.execute("SELECT count(*) FROM muse_change_impact").fetchone()[0] == 1
def test_参与者失败回滚业务决定影响与命令__a91002(变更环境):
服务, 身份, 命令, 库 = 变更环境
with pytest.raises(RuntimeError, match="合成参与者"):
服务.提交变更(身份, replace(命令, 业务请求=合成编辑("候选内容", True)))
with 库[用途.生产].连接() as 连:
assert 连.execute("SELECT revision FROM s01_test_document").fetchone()[0] == 1
assert 连.execute("SELECT count(*) FROM s01_test_audit").fetchone()[0] == 0
assert 连.execute("SELECT count(*) FROM muse_candidate_decision").fetchone()[0] == 0
assert 连.execute("SELECT count(*) FROM muse_change_command").fetchone()[0] == 0
assert 连.execute("SELECT count(*) FROM muse_change_impact").fetchone()[0] == 0
assert 服务.提交变更(身份, 命令)["results"] == [{"revision": 2}]
def test_两份命令并发采纳只有一份当前写入__a91003(变更环境):
服务, 身份, 命令, 库 = 变更环境
栅栏 = Barrier(2)
def 竞争(id_):
栅栏.wait()
try:
return 服务.提交变更(身份, replace(命令, 命令ID=id_))
except 变更错误 as exc:
return exc.错误码
with ThreadPoolExecutor(max_workers=2) as pool:
结果 = list(pool.map(竞争, ["command-A", "command-B"]))
assert sum(isinstance(v, dict) for v in 结果) == 1
assert next(v for v in 结果 if isinstance(v, str)) in {"SOURCE_STALE", "REVISION_CONFLICT"}
with 库[用途.生产].连接() as 连:
assert 连.execute("SELECT count(*) FROM s01_test_audit").fetchone()[0] == 1
assert 连.execute("SELECT revision FROM s01_test_document").fetchone()[0] == 2
def test_来源改变拒绝旧审阅且零业务写入__a91004(变更环境):
服务, 身份, 命令, 库 = 变更环境
with 库[用途.维护].连接() as 连:
连.execute("UPDATE s01_test_source SET revision=2")
with pytest.raises(变更错误) as 捕获:
服务.提交变更(身份, 命令)
assert 捕获.value.错误码 == "SOURCE_STALE"
with 库[用途.生产].连接() as 连:
assert 连.execute("SELECT revision FROM s01_test_document").fetchone()[0] == 1
assert 连.execute("SELECT count(*) FROM muse_change_command").fetchone()[0] == 0
def test_检查未过仍可拒绝或暂缓且不改正文__a91005(变更环境):
服务, 身份, 命令, 库 = 变更环境
with 库[用途.维护].连接() as 连:
连.execute("UPDATE s01_test_candidate SET check_state='issues'")
with pytest.raises(变更错误) as 捕获:
服务.提交变更(身份, 命令)
assert 捕获.value.错误码 == "PRECONDITION_PENDING"
for action in [作者动作.暂缓, 作者动作.拒绝]:
回执 = 服务.提交变更(身份, replace(命令, 命令ID=action.value, 动作=action, 批准变更=()))
assert 回执["action"] == action.value and 回执["results"] == []
with 库[用途.生产].连接() as 连:
assert 连.execute("SELECT revision FROM s01_test_document").fetchone()[0] == 1
assert 连.execute("SELECT decision FROM muse_candidate_decision").fetchone()[0] == "reject"
def test_非作者评测用途和旧候选版本拒绝__a91006(变更环境):
服务, 身份, 命令, 库 = 变更环境
for bad in [
replace(身份, 作者=None),
replace(身份, 作者="other"),
replace(身份, 用途=用途.评测),
]:
with pytest.raises(变更错误) as 捕获:
服务.提交变更(bad, 命令)
assert 捕获.value.错误码 == "SCOPE_DENIED"
with 库[用途.维护].连接() as 连:
连.execute("UPDATE s01_test_candidate SET revision=2,content='新版候选'")
with pytest.raises(变更错误) as 捕获:
服务.提交变更(身份, 命令)
assert 捕获.value.错误码 == "REVISION_CONFLICT"
def test_人工保存不要求候选检查或审阅__a91007(变更环境):
服务, 身份, 命令, 库 = 变更环境
with 库[用途.维护].连接() as 连:
连.execute("UPDATE s01_test_candidate SET check_state='issues'")
手写 = replace(
命令,
动作=作者动作.人工保存,
业务请求=合成编辑("作者亲自写的内容"),
审阅ID=None,
审阅哈希=None,
批准变更=(),
)
assert 服务.提交变更(身份, 手写)["results"] == [{"revision": 2}]
with 库[用途.生产].连接() as 连:
assert (
连.execute("SELECT content FROM s01_test_document").fetchone()[0] == "作者亲自写的内容"
)
assert 连.execute("SELECT count(*) FROM muse_candidate_decision").fetchone()[0] == 0
def test_预览返回具体影响而不记录批准或写入__a91008(变更环境):
服务, 身份, 命令, 库 = 变更环境
预览 = 服务.准备变更(身份, 命令)
assert 预览["participants"] == ["write"]
assert 预览["changes"] == ["replace-content"]
assert 预览["impacts"][0]["新版本"] == "2"
with 库[用途.生产].连接() as 连:
assert 连.execute("SELECT count(*) FROM muse_change_command").fetchone()[0] == 0
assert 连.execute("SELECT revision FROM s01_test_document").fetchone()[0] == 1
def test_来源当前指针在正式提交期间不可改变__a91009(变更环境):
from threading import Event
from time import monotonic
服务, 身份, 命令, 库 = 变更环境
已准备, 放行, 更新开始 = Event(), Event(), Event()
pid = []
class 等待参与者(合成写入):
def 提交(self, 连, 请求):
已准备.set()
assert 放行.wait(5)
return super().提交(连, 请求)
# 仅替换本例已登记的合成参与者;不向产品接入开放运行时回调。
服务.目录.参与者["write"] = 等待参与者()
def 更新来源():
with 库[用途.生产].连接() as 连:
pid.append(连.execute("SELECT pg_backend_pid()").fetchone()[0])
更新开始.set()
连.execute("UPDATE s01_test_source SET revision=2")
with ThreadPoolExecutor(max_workers=2) as pool:
提交 = pool.submit(服务.提交变更, 身份, 命令)
try:
assert 已准备.wait(3)
更新 = pool.submit(更新来源)
assert 更新开始.wait(2)
截止 = monotonic() + 2
阻塞 = False
with 库[用途.生产].连接() as 连:
while monotonic() < 截止:
阻塞 = 连.execute(
"SELECT EXISTS(SELECT 1 FROM pg_locks WHERE pid=%s AND NOT granted)",
(pid[0],),
).fetchone()[0]
if 阻塞:
break
Event().wait(0.01)
assert 阻塞, "来源更新没有被当前指针锁阻止"
assert 连.execute("SELECT revision FROM s01_test_source").fetchone()[0] == 1
finally:
放行.set()
assert 提交.result(timeout=3)["results"] == [{"revision": 2}]
更新.result(timeout=3)
with 库[用途.生产].连接() as 连:
assert 连.execute("SELECT revision FROM s01_test_source").fetchone()[0] == 2