From a6f4f14a8140f8c7085d39a76c78d2f170e3bf53 Mon Sep 17 00:00:00 2001 From: zizi Date: Thu, 10 Sep 2026 19:25:41 +0800 Subject: [PATCH] =?UTF-8?q?W09=20=E6=AD=A3=E5=BC=8F=E5=8F=98=E6=9B=B4?= =?UTF-8?q?=E3=80=81=E5=AE=A1=E9=98=85=E3=80=81CAS=20=E4=B8=8E=E5=8E=9F?= =?UTF-8?q?=E5=AD=90=E6=8F=90=E4=BA=A4=EF=BC=9AS01=20=E6=AD=A3=E5=BC=8F?= =?UTF-8?q?=E5=8F=98=E6=9B=B4=E3=80=81=E4=BD=9C=E8=80=85=E5=AE=A1=E9=98=85?= =?UTF-8?q?=E3=80=81CAS=20=E5=B9=B6=E5=8F=91=E4=B8=8E=E8=B7=A8=E6=A8=A1?= =?UTF-8?q?=E5=9D=97=E5=8E=9F=E5=AD=90=E6=8F=90=E4=BA=A4=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 按 R2 串行阶段整理提交;包内文件为该阶段交付(含后续小增量),状态以工作包清单为准。 --- src/muse/正式变更/__init__.py | 1 + src/muse/正式变更/作者审阅.py | 27 +++ src/muse/正式变更/依赖锁.py | 18 ++ src/muse/正式变更/修改影响.py | 11 ++ src/muse/正式变更/存储.py | 111 +++++++++++ src/muse/正式变更/幂等命令.py | 9 + src/muse/正式变更/接口.py | 131 +++++++++++++ src/muse/正式变更/提交事务.py | 42 ++++ src/muse/正式变更/模型.py | 134 +++++++++++++ src/muse/正式变更/版本检查.py | 17 ++ src/muse/正式变更/状态迁移.py | 11 ++ tests/集成/test_正式变更事务.py | 332 ++++++++++++++++++++++++++++++++ 数据库/迁移/V0003__正式变更.sql | 38 ++++ 13 files changed, 882 insertions(+) create mode 100644 src/muse/正式变更/__init__.py create mode 100644 src/muse/正式变更/作者审阅.py create mode 100644 src/muse/正式变更/依赖锁.py create mode 100644 src/muse/正式变更/修改影响.py create mode 100644 src/muse/正式变更/存储.py create mode 100644 src/muse/正式变更/幂等命令.py create mode 100644 src/muse/正式变更/接口.py create mode 100644 src/muse/正式变更/提交事务.py create mode 100644 src/muse/正式变更/模型.py create mode 100644 src/muse/正式变更/版本检查.py create mode 100644 src/muse/正式变更/状态迁移.py create mode 100644 tests/集成/test_正式变更事务.py create mode 100644 数据库/迁移/V0003__正式变更.sql diff --git a/src/muse/正式变更/__init__.py b/src/muse/正式变更/__init__.py new file mode 100644 index 0000000..79cc27c --- /dev/null +++ b/src/muse/正式变更/__init__.py @@ -0,0 +1 @@ +"""正式变更模块;公开合同由接口.py提供。""" diff --git a/src/muse/正式变更/作者审阅.py b/src/muse/正式变更/作者审阅.py new file mode 100644 index 0000000..f98e055 --- /dev/null +++ b/src/muse/正式变更/作者审阅.py @@ -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", 命令.动作) diff --git a/src/muse/正式变更/依赖锁.py b/src/muse/正式变更/依赖锁.py new file mode 100644 index 0000000..8718876 --- /dev/null +++ b/src/muse/正式变更/依赖锁.py @@ -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 变更错误(代码, "当前来源、结构或投影已改变,需要重新核对") diff --git a/src/muse/正式变更/修改影响.py b/src/muse/正式变更/修改影响.py new file mode 100644 index 0000000..ba1edf3 --- /dev/null +++ b/src/muse/正式变更/修改影响.py @@ -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 影响] diff --git a/src/muse/正式变更/存储.py b/src/muse/正式变更/存储.py new file mode 100644 index 0000000..9769127 --- /dev/null +++ b/src/muse/正式变更/存储.py @@ -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 身份 diff --git a/src/muse/正式变更/幂等命令.py b/src/muse/正式变更/幂等命令.py new file mode 100644 index 0000000..7efb999 --- /dev/null +++ b/src/muse/正式变更/幂等命令.py @@ -0,0 +1,9 @@ +"""调用身份先于幂等回读校验;模型任务不能冒充作者确认。""" + +from muse.共享.调用身份 import 用途, 调用身份 +from muse.正式变更.模型 import 变更错误 + + +def 核对正式身份(身份: 调用身份, 数据库用途: 用途) -> None: + if not 身份.允许写正式内容 or 数据库用途 is not 用途.生产: + raise 变更错误("SCOPE_DENIED", "正式变更只接受生产用途的已认证作者动作") diff --git a/src/muse/正式变更/接口.py b/src/muse/正式变更/接口.py new file mode 100644 index 0000000..82d02cc --- /dev/null +++ b/src/muse/正式变更/接口.py @@ -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__ = [ + "正式变更服务", + "参与者目录", + "作者动作", + "依赖引用", + "依赖维护者", + "修改影响", + "参与操作", + "变更入口", + "变更命令", + "变更错误", + "固定哈希", + "审阅内容", + "提交计划", + "事务参与者", + "检查决定", + "读取作者决定", +] diff --git a/src/muse/正式变更/提交事务.py b/src/muse/正式变更/提交事务.py new file mode 100644 index 0000000..920f440 --- /dev/null +++ b/src/muse/正式变更/提交事务.py @@ -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 计划.操作] diff --git a/src/muse/正式变更/模型.py b/src/muse/正式变更/模型.py new file mode 100644 index 0000000..5ed836a --- /dev/null +++ b/src/muse/正式变更/模型.py @@ -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: + """锁定可变的当前指针行直到事务结束,返回当前有效版本。""" + ... diff --git a/src/muse/正式变更/版本检查.py b/src/muse/正式变更/版本检查.py new file mode 100644 index 0000000..3579241 --- /dev/null +++ b/src/muse/正式变更/版本检查.py @@ -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", "批准清单与本候选不一致;部分采纳先形成新候选版本") diff --git a/src/muse/正式变更/状态迁移.py b/src/muse/正式变更/状态迁移.py new file mode 100644 index 0000000..e04a290 --- /dev/null +++ b/src/muse/正式变更/状态迁移.py @@ -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 diff --git a/tests/集成/test_正式变更事务.py b/tests/集成/test_正式变更事务.py new file mode 100644 index 0000000..9fb09d1 --- /dev/null +++ b/tests/集成/test_正式变更事务.py @@ -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 diff --git a/数据库/迁移/V0003__正式变更.sql b/数据库/迁移/V0003__正式变更.sql new file mode 100644 index 0000000..cf90063 --- /dev/null +++ b/数据库/迁移/V0003__正式变更.sql @@ -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) +);