W09 正式变更、审阅、CAS 与原子提交:S01 正式变更、作者审阅、CAS 并发与跨模块原子提交。

按 R2 串行阶段整理提交;包内文件为该阶段交付(含后续小增量),状态以工作包清单为准。
This commit is contained in:
zizi 2026-09-10 19:25:41 +08:00
parent d82581daee
commit a6f4f14a81
13 changed files with 882 additions and 0 deletions

View File

@ -0,0 +1 @@
"""正式变更模块;公开合同由接口.py提供。"""

View File

@ -0,0 +1,27 @@
"""被展示版本不是批准;正式决定需要命令、审阅和当前版本同时对应。"""
from muse.正式变更.存储 import 变更存储
from muse.正式变更.模型 import 作者动作, 变更错误
from muse.正式变更.版本检查 import 检查审阅
from muse.正式变更.状态迁移 import 检查决定
def 核对作者审阅(存储: 变更存储, 身份, 命令, 当前) -> None:
if 命令.动作 is 作者动作.人工保存:
return
if not 命令.审阅ID or 当前.候选ID is None or 当前.候选版本 is None:
raise 变更错误("REVIEW_MISMATCH", "候选决定需要实际候选版本和审阅记录")
行 = 存储.审阅行(命令.审阅ID, 身份.作者, 命令.入口, 当前.目标)
原快照 = 行["snapshot"]
for 字段, 代码 in (
("候选版本", "REVISION_CONFLICT"),
("候选哈希", "REVISION_CONFLICT"),
("有效结构哈希", "SCHEMA_STALE"),
("投影版本", "PROJECTION_STALE"),
("来源摘要哈希", "SOURCE_STALE"),
):
if 原快照[字段] != getattr(当前, 字段):
raise 变更错误(代码, "被审候选、结构、投影或来源已经改变")
检查审阅(命令, 当前, 行["content_hash"])
原决定 = 存储.决定行(当前.候选ID, 当前.候选版本)
检查决定(原决定["decision"] if 原决定 else "undecided", 命令.动作)

View File

@ -0,0 +1,18 @@
"""统一依赖锁序;通过所属owner保护当前指针,不锁历史版本假装保护现状。"""
from muse.正式变更.模型 import 依赖引用, 依赖维护者, 变更错误
def 锁定依赖(连, 引用: tuple[依赖引用, ...], 目录: dict[str, 依赖维护者]) -> None:
键 = [(r.种类, r.对象ID) for r in 引用]
if len(键) != len(set(键)):
raise 变更错误("INVALID_REQUEST", "同一当前指针不能登记多个依赖版本")
for 项 in sorted(引用):
维护者 = 目录.get(项.种类)
if 维护者 is None:
raise 变更错误("PRECONDITION_PENDING", "该类依赖的当前指针保护尚未登记")
if 维护者.锁定当前(连, 项.对象ID) != 项.版本:
代码 = {"schema": "SCHEMA_STALE", "projection": "PROJECTION_STALE"}.get(
项.种类, "SOURCE_STALE"
)
raise 变更错误(代码, "当前来源、结构或投影已改变,需要重新核对")

View File

@ -0,0 +1,11 @@
"""提交时登记持久复核影响;不暗改受影响正文或事实。"""
from muse.正式变更.存储 import 变更存储
from muse.正式变更.模型 import 修改影响, 变更错误
def 登记影响(存储: 变更存储, 回执ID: str, 影响: tuple[修改影响, ...]) -> list[str]:
键 = [(v.对象, v.原因) for v in 影响]
if len(set(键)) != len(键):
raise 变更错误("INVALID_OUTPUT", "变更计划重复登记同一对象的影响")
return [存储.保存影响(回执ID, v) for v in 影响]

View File

@ -0,0 +1,111 @@
"""S01唯一SQL归属;所有写入沿用传入的显式事务。"""
from dataclasses import asdict
from typing import LiteralString
from uuid import uuid4
from psycopg.rows import dict_row
from psycopg.types.json import Jsonb
from muse.共享.调用身份 import 调用身份
from muse.正式变更.模型 import 修改影响, 变更命令, 变更错误, 固定哈希, 审阅内容
class 变更存储:
def __init__(self, 连):
self.连 = 连
def 查询(self, SQL: LiteralString, 参数=()):
return self.连.cursor(row_factory=dict_row).execute(SQL, 参数)
def 占用命令(self, 身份: 调用身份, 命令: 变更命令) -> dict | None:
哈希 = 固定哈希(命令)
键 = (身份.作者, 身份.用途.value, 命令.命令ID)
self.查询(
"INSERT INTO muse_change_command (author_id,run_purpose,command_id,request_hash) "
"VALUES (%s,%s,%s,%s) ON CONFLICT DO NOTHING",
(*键, 哈希),
)
行 = self.查询(
"SELECT request_hash,receipt FROM muse_change_command "
"WHERE author_id=%s AND run_purpose=%s AND command_id=%s FOR UPDATE",
键,
).fetchone()
if 行 is None or 行["request_hash"] != 哈希:
raise 变更错误("COMMAND_PAYLOAD_MISMATCH", "命令身份已用于不同请求")
return 行["receipt"]
def 保存回执(self, 身份: 调用身份, 命令ID: str, 回执: dict) -> dict:
行 = self.查询(
"UPDATE muse_change_command SET receipt=%s WHERE author_id=%s "
"AND run_purpose=%s AND command_id=%s RETURNING receipt",
(Jsonb(回执), 身份.作者, 身份.用途.value, 命令ID),
).fetchone()
if 行 is None:
raise 变更错误("RECEIPT_NOT_FOUND", "命令未登记,不能保存提交回执")
return 行["receipt"]
def 读取回执(self, 身份: 调用身份, 命令ID: str) -> dict:
行 = self.查询(
"SELECT receipt FROM muse_change_command WHERE author_id=%s "
"AND run_purpose=%s AND command_id=%s",
(身份.作者, 身份.用途.value, 命令ID),
).fetchone()
if 行 is None or 行["receipt"] is None:
raise 变更错误("RECEIPT_NOT_FOUND", "没有对应作者命令的已提交回执")
return 行["receipt"]
def 新审阅(self, 作者: str, 入口: str, 当前: 审阅内容) -> dict:
身份 = str(uuid4())
哈希 = 固定哈希(当前)
self.查询(
"INSERT INTO muse_author_review (review_id,author_id,entry_id,target_ref,"
"content_hash,snapshot) VALUES (%s,%s,%s,%s,%s,%s)",
(身份, 作者, 入口, 当前.目标, 哈希, Jsonb(asdict(当前))),
)
return {"review_id": 身份, "review_hash": 哈希, "snapshot": asdict(当前)}
def 审阅行(self, 身份: str, 作者: str, 入口: str, 目标: str) -> dict:
行 = self.查询(
"SELECT * FROM muse_author_review WHERE review_id=%s AND author_id=%s "
"AND entry_id=%s AND target_ref=%s FOR SHARE",
(身份, 作者, 入口, 目标),
).fetchone()
if 行 is None:
raise 变更错误("REVIEW_MISMATCH", "审阅记录不属于本作者、操作或目标")
return 行
def 决定行(self, 候选ID: str, 版本: int) -> dict | None:
return self.查询(
"SELECT decision FROM muse_candidate_decision WHERE candidate_id=%s "
"AND candidate_revision=%s FOR UPDATE",
(候选ID, 版本),
).fetchone()
def 读取决定(self, 候选ID: str, 版本: int) -> str:
行 = self.查询(
"SELECT decision FROM muse_candidate_decision WHERE candidate_id=%s "
"AND candidate_revision=%s",
(候选ID, 版本),
).fetchone()
return 行["decision"] if 行 else "undecided"
def 保存决定(self, 身份: 调用身份, 命令: 变更命令, 当前: 审阅内容) -> None:
self.查询(
"INSERT INTO muse_candidate_decision "
"(candidate_id,candidate_revision,decision,author_id,review_id,command_id) "
"VALUES (%s,%s,%s,%s,%s,%s) ON CONFLICT (candidate_id,candidate_revision) "
"DO UPDATE SET decision=excluded.decision,author_id=excluded.author_id,"
"review_id=excluded.review_id,command_id=excluded.command_id",
(当前.候选ID, 当前.候选版本, 命令.动作.value, 身份.作者, 命令.审阅ID, 命令.命令ID),
)
def 保存影响(self, 回执ID: str, 影响: 修改影响) -> str:
身份 = str(uuid4())
self.查询(
"INSERT INTO muse_change_impact "
"(impact_id,receipt_id,target_ref,reason,before_version,after_version) "
"VALUES (%s,%s,%s,%s,%s,%s)",
(身份, 回执ID, 影响.对象, 影响.原因, 影响.原版本, 影响.新版本),
)
return 身份

View File

@ -0,0 +1,9 @@
"""调用身份先于幂等回读校验;模型任务不能冒充作者确认。"""
from muse.共享.调用身份 import 用途, 调用身份
from muse.正式变更.模型 import 变更错误
def 核对正式身份(身份: 调用身份, 数据库用途: 用途) -> None:
if not 身份.允许写正式内容 or 数据库用途 is not 用途.生产:
raise 变更错误("SCOPE_DENIED", "正式变更只接受生产用途的已认证作者动作")

View File

@ -0,0 +1,131 @@
"""正式变更的唯一公开入口;生产作者、当前依赖和同事务回执一起验证。"""
from dataclasses import asdict
from uuid import uuid4
from muse.共享.调用身份 import 调用身份
from muse.基础设施.数据库.连接 import 数据库工厂
from muse.正式变更.作者审阅 import 核对作者审阅
from muse.正式变更.依赖锁 import 锁定依赖
from muse.正式变更.修改影响 import 登记影响
from muse.正式变更.存储 import 变更存储
from muse.正式变更.幂等命令 import 核对正式身份
from muse.正式变更.提交事务 import 参与者目录
from muse.正式变更.模型 import (
事务参与者,
作者动作,
依赖引用,
依赖维护者,
修改影响,
参与操作,
变更入口,
变更命令,
变更错误,
固定哈希,
审阅内容,
提交计划,
)
from muse.正式变更.版本检查 import 检查目标
from muse.正式变更.状态迁移 import 检查决定
def 读取作者决定(连, 候选ID: str, 版本: int) -> str:
"""业务owner完成归属验证后,在原事务中读取S01拥有的决定。"""
return 变更存储(连).读取决定(候选ID, 版本)
class 正式变更服务:
def __init__(self, 数据库: 数据库工厂, 目录: 参与者目录):
self.数据库, self.目录 = 数据库, 目录
def 打开审阅(self, 身份: 调用身份, 命令: 变更命令) -> dict:
核对正式身份(身份, self.数据库.用途)
assert 身份.作者 is not None
入口 = self.目录.查入口(命令.入口, 命令.业务请求)
with self.数据库.连接() as 连, 连.transaction():
当前 = 入口.查看(连, 身份, 命令)
检查目标(命令, 当前)
return 变更存储(连).新审阅(身份.作者, 命令.入口, 当前)
def 提交变更(self, 身份: 调用身份, 命令: 变更命令) -> dict:
核对正式身份(身份, self.数据库.用途)
入口 = self.目录.查入口(命令.入口, 命令.业务请求)
with self.数据库.连接() as 连, 连.transaction():
存储 = 变更存储(连)
既有 = 存储.占用命令(身份, 命令)
if 既有 is not None:
return 既有
预读 = 入口.查看(连, 身份, 命令)
检查目标(命令, 预读)
必锁 = {预读.目标}
if 预读.候选ID is not None:
必锁.add(预读.候选ID)
if not 必锁 <= {r.对象ID for r in 预读.依赖}:
raise 变更错误("PRECONDITION_PENDING", "业务入口缺少目标或候选的当前指针保护")
锁定依赖(连, 预读.依赖, self.目录.依赖)
当前 = 入口.查看(连, 身份, 命令)
检查目标(命令, 当前)
if 预读.依赖 != 当前.依赖:
raise 变更错误("SOURCE_STALE", "依赖范围改变,需要重新打开审阅")
核对作者审阅(存储, 身份, 命令, 当前)
计划 = 入口.准备(连, 身份, 命令, 当前)
if 命令.动作 in {作者动作.采纳, 作者动作.人工保存} and not 计划.操作:
raise 变更错误("INVALID_OUTPUT", "正式写入计划没有业务参与者")
结果 = self.目录.提交(连, 计划)
回执ID = str(uuid4())
影响 = 登记影响(存储, 回执ID, 计划.影响)
if 命令.动作 is not 作者动作.人工保存:
存储.保存决定(身份, 命令, 当前)
回执 = {
"receipt_id": 回执ID,
"command_id": 命令.命令ID,
"target_ref": 命令.目标,
"action": 命令.动作.value,
"author_review_id": 命令.审阅ID,
"basis": asdict(当前),
"results": 结果,
"impact_ids": 影响,
}
return 存储.保存回执(身份, 命令.命令ID, 回执)
def 准备变更(self, 身份: 调用身份, 命令: 变更命令) -> dict:
"""只读预览后果;提交时重新准备,不接收客户端传回的计划作为授权。"""
核对正式身份(身份, self.数据库.用途)
入口 = self.目录.查入口(命令.入口, 命令.业务请求)
with self.数据库.连接(只读=True) as 连:
当前 = 入口.查看(连, 身份, 命令)
检查目标(命令, 当前)
计划 = 入口.准备(连, 身份, 命令, 当前)
self.目录.检查计划(计划)
return {
"target_ref": 当前.目标,
"review_hash": 固定哈希(当前),
"changes": list(当前.变更项),
"participants": [o.参与者 for o in 计划.操作],
"impacts": [asdict(v) for v in 计划.影响],
}
def 读取回执(self, 身份: 调用身份, 命令ID: str) -> dict:
核对正式身份(身份, self.数据库.用途)
with self.数据库.连接(只读=True) as 连:
return 变更存储(连).读取回执(身份, 命令ID)
__all__ = [
"正式变更服务",
"参与者目录",
"作者动作",
"依赖引用",
"依赖维护者",
"修改影响",
"参与操作",
"变更入口",
"变更命令",
"变更错误",
"固定哈希",
"审阅内容",
"提交计划",
"事务参与者",
"检查决定",
"读取作者决定",
]

View File

@ -0,0 +1,42 @@
"""具体类型参与者共享一个事务;不接受任意数据库写入描述。"""
from muse.正式变更.模型 import 事务参与者, 依赖维护者, 变更入口, 变更错误, 提交计划
class 参与者目录:
def __init__(self):
self.入口: dict[str, 变更入口] = {}
self.参与者: dict[str, 事务参与者] = {}
self.依赖: dict[str, 依赖维护者] = {}
def 登记入口(self, 入口: 变更入口) -> None:
if not 入口.身份 or 入口.身份 in self.入口:
raise 变更错误("INVALID_REQUEST", "业务入口身份为空或重复")
self.入口[入口.身份] = 入口
def 登记参与者(self, 参与者: 事务参与者) -> None:
if not 参与者.身份 or 参与者.身份 in self.参与者:
raise 变更错误("INVALID_REQUEST", "事务参与者身份为空或重复")
self.参与者[参与者.身份] = 参与者
def 登记依赖(self, 维护者: 依赖维护者) -> None:
if not 维护者.种类 or 维护者.种类 in self.依赖:
raise 变更错误("INVALID_REQUEST", "当前指针维护者为空或重复")
self.依赖[维护者.种类] = 维护者
def 查入口(self, 身份: str, 请求: object) -> 变更入口:
入口 = self.入口.get(身份)
if 入口 is None or type(请求) is not 入口.请求类型:
raise 变更错误("INVALID_REQUEST", "业务入口或请求类型未登记")
return 入口
def 检查计划(self, 计划: 提交计划) -> None:
# 先核对整份计划的参与者与类型,不执行到一半才发现陌生写入器。
for 操作 in 计划.操作:
参与者 = self.参与者.get(操作.参与者)
if 参与者 is None or type(操作.请求) is not 参与者.请求类型:
raise 变更错误("INVALID_OUTPUT", "变更计划包含未登记的参与者或错误请求类型")
def 提交(self, 连, 计划: 提交计划) -> list[dict]:
self.检查计划(计划)
return [self.参与者[o.参与者].提交(连, o.请求) for o in 计划.操作]

View File

@ -0,0 +1,134 @@
"""共用审阅与提交外壳;业务值和业务校验归具体参与者。"""
from __future__ import annotations
import hashlib
import json
from dataclasses import asdict, dataclass, is_dataclass
from enum import StrEnum
from typing import Any, Protocol
import psycopg
from muse.共享.调用身份 import 调用身份
from muse.共享.错误 import Muse错误
class 变更错误(Muse错误):
def __init__(self, 代码: str, 说明: str):
super().__init__(说明)
self.错误码 = 代码
def 固定哈希(内容: Any) -> str:
if is_dataclass(内容) and not isinstance(内容, type):
内容 = asdict(内容)
try:
字节 = json.dumps(
内容, ensure_ascii=False, sort_keys=True, separators=(",", ":"), allow_nan=False
).encode()
except (ValueError, TypeError):
raise 变更错误("INVALID_REQUEST", "变更请求必须是可固定哈希的类型化内容") from None
return hashlib.sha256(字节).hexdigest()
class 作者动作(StrEnum):
人工保存 = "manual_save"
采纳 = "adopt"
拒绝 = "reject"
暂缓 = "defer"
@dataclass(frozen=True, slots=True, order=True)
class 依赖引用:
种类: str
对象ID: str
版本: str
@dataclass(frozen=True, slots=True)
class 审阅内容:
目标: str
数据版本: int
候选ID: str | None
候选版本: int | None
候选哈希: str | None
差异哈希: str
有效结构哈希: str | None
投影版本: str | None
来源摘要哈希: str
变更项: tuple[str, ...]
依赖: tuple[依赖引用, ...]
@dataclass(frozen=True, slots=True)
class 变更命令:
命令ID: str
入口: str
目标: str
动作: 作者动作
预期数据版本: int
业务请求: object
审阅ID: str | None = None
审阅哈希: str | None = None
批准变更: tuple[str, ...] = ()
def __post_init__(self):
if not isinstance(self.动作, 作者动作):
raise 变更错误("INVALID_REQUEST", "作者动作必须使用已登记的类型")
if not self.命令ID or not self.入口 or not self.目标 or self.预期数据版本 < 0:
raise 变更错误("INVALID_REQUEST", "变更命令需要身份、目标和明确的预期版本")
if len(self.批准变更) != len(set(self.批准变更)):
raise 变更错误("INVALID_REQUEST", "批准清单不能重复")
@dataclass(frozen=True, slots=True)
class 参与操作:
参与者: str
请求: object
@dataclass(frozen=True, slots=True)
class 修改影响:
对象: str
原因: str
原版本: str
新版本: str
@dataclass(frozen=True, slots=True)
class 提交计划:
操作: tuple[参与操作, ...]
影响: tuple[修改影响, ...] = ()
class 变更入口(Protocol):
身份: str
请求类型: type
def 查看(self, 连: psycopg.Connection, 身份: 调用身份, 命令: 变更命令) -> 审阅内容:
"""核对作者与目标范围,从业务权威读取当前版本和依赖。"""
...
def 准备(
self, 连: psycopg.Connection, 身份: 调用身份, 命令: 变更命令, 当前: 审阅内容
) -> 提交计划:
"""再次执行业务前置;采纳和人工保存不能绕过该入口的验证。"""
...
class 事务参与者(Protocol):
身份: str
请求类型: type
def 提交(self, 连: psycopg.Connection, 请求: object, /) -> dict:
"""只操作本模块数据,沿用S01事务,返回可回查身份。"""
...
class 依赖维护者(Protocol):
种类: str
def 锁定当前(self, 连: psycopg.Connection, 对象ID: str) -> str:
"""锁定可变的当前指针行直到事务结束,返回当前有效版本。"""
...

View File

@ -0,0 +1,17 @@
"""目标、候选与被审结构的版本比较,不写业务内容。"""
from muse.正式变更.模型 import 作者动作, 变更命令, 变更错误, 固定哈希, 审阅内容
def 检查目标(命令: 变更命令, 当前: 审阅内容) -> None:
if 命令.目标 != 当前.目标:
raise 变更错误("SCOPE_DENIED", "业务目标与命令不一致")
if 命令.预期数据版本 != 当前.数据版本:
raise 变更错误("REVISION_CONFLICT", "目标内容已经改变,保留输入并比较当前版本")
def 检查审阅(命令: 变更命令, 当前: 审阅内容, 已审哈希: str) -> None:
if 命令.审阅哈希 != 已审哈希 or 固定哈希(当前) != 已审哈希:
raise 变更错误("REVIEW_MISMATCH", "候选、差异、结构或来源与作者审阅时不一致")
if 命令.动作 is 作者动作.采纳 and set(命令.批准变更) != set(当前.变更项):
raise 变更错误("REVIEW_MISMATCH", "批准清单与本候选不一致;部分采纳先形成新候选版本")

View File

@ -0,0 +1,11 @@
"""作者决定与检查结果独立;具体采纳质量条件由业务owner判断。"""
from muse.正式变更.模型 import 作者动作, 变更错误
def 检查决定(原决定: str, 动作: 作者动作) -> str:
if 原决定 not in {"undecided", "defer"}:
raise 变更错误("REVISION_CONFLICT", "该候选版本已有最终决定,请查看回执")
if 动作 is 作者动作.人工保存:
raise 变更错误("INVALID_REQUEST", "人工保存不改写候选的作者决定")
return 动作.value

View File

@ -0,0 +1,332 @@
"""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

View File

@ -0,0 +1,38 @@
-- S01只保存命令、审阅、决定与影响;业务实例表由各owner迁移定义。
CREATE TABLE public.muse_change_command (
author_id text NOT NULL,
run_purpose public.muse_purpose NOT NULL CHECK (run_purpose='production'),
command_id text NOT NULL,
request_hash text NOT NULL,
receipt jsonb,
created_at timestamptz NOT NULL DEFAULT clock_timestamp(),
PRIMARY KEY(author_id,run_purpose,command_id)
);
CREATE TABLE public.muse_author_review (
review_id uuid PRIMARY KEY,
author_id text NOT NULL,
entry_id text NOT NULL,
target_ref text NOT NULL,
content_hash text NOT NULL,
snapshot jsonb NOT NULL,
created_at timestamptz NOT NULL DEFAULT clock_timestamp()
);
CREATE TABLE public.muse_candidate_decision (
candidate_id text NOT NULL,
candidate_revision bigint NOT NULL,
decision text NOT NULL CHECK (decision IN ('adopt','reject','defer')),
author_id text NOT NULL,
review_id uuid NOT NULL REFERENCES public.muse_author_review(review_id),
command_id text NOT NULL,
PRIMARY KEY(candidate_id,candidate_revision)
);
CREATE TABLE public.muse_change_impact (
impact_id uuid PRIMARY KEY,
receipt_id uuid NOT NULL,
target_ref text NOT NULL,
reason text NOT NULL,
before_version text NOT NULL,
after_version text NOT NULL,
state text NOT NULL DEFAULT 'pending' CHECK (state IN ('pending','handled')),
UNIQUE(receipt_id,target_ref,reason)
);