muse-agent-example/tests/集成/test_原文归档与恢复.py
zizi d909d1bd1b 后端实现与用例身份:19 包集成落地并修复收尾缺陷
实现侧:
- 上下文:任务范围拆分为 范围校验/范围授权;索引按可发现口径重建、索引新鲜度改对称差;依赖校验统一快照漂移说明。
- 知识方法:方法与材料读取口径统一;超限方法材料按可选省略,核对路径不再二次计费;删除无合同的读时重算。
- 任务运行:新增 context.usage/tool.denied 事件类型;连接池常驻并在装配生命周期内开关;调用结算与核对分列。
- 效果评测/审校修订/交付连载/作者经验/作品规划:凭据冻结、标定消费、导出补证、事实引文核对等收尾修复。
- 资源加载:能力正文不再夹带索引用的导航注记(该注记此前进入角色与技能的模型提示)。
- 元数据:受保护骨架与代码保护属性对齐;字段校验与内置结构口径同步。
- 基础设施:环境预检进入装配生命周期;数据库连接运行期字段不参与相等比较;索引指纹归一化 jsonb 浮点。
- 删除被替代实现:7 份旧提示词模板与空壳 资料来源 读取器。

用例侧:
- 用例身份与导航元信息迁移;夹具补生命周期、同库暴露与模板封存;
- 本轮定向修复:方法材料省略、事实引文、迁移回执、额度与暂停用例、慢用例超时预算等。
2026-09-18 01:15:00 +08:00

728 lines
32 KiB
Python
Raw Permalink 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.

"""已授权的合成原文;独立PG验证先租约、完整归档与中断恢复。"""
import hashlib
from datetime import UTC, datetime, timedelta
import pytest
from muse.任务运行.接口 import (
任务服务,
任务请求,
原文服务,
原文错误,
步骤处理器,
步骤结果,
步骤计划,
证据服务,
)
from muse.共享.时间 import 假时钟
from muse.共享.调用身份 import 内容用途, 用途
from muse.基础设施.受控文件 import 受控文件, 受控文件错误
from muse.编排.接口 import 流程定义, 流程服务, 流程登记
pytestmark = pytest.mark.数据库
def _断言租约状态(状态, 租约, 预期, *, 归档=None, 已清理=False):
"""仅比较已读取值,各失败断点仍由独立用例触发。"""
assert 状态 == {
"lease_id": 租约,
"state": 预期,
"archive_id": 归档,
"cleanup_complete": 已清理,
}
def _清理快照(服务, 任务, 租约, 库):
with 库[用途.生产].连接() as 连:
回执 = 连.execute(
"SELECT cleanup_receipt, failure_code FROM muse_raw_lease WHERE lease_id=%s",
(租约,),
).fetchone()
文件集合 = {str(p.relative_to(服务.文件.根)) for p in 服务.文件.根.rglob("*")}
return 服务.状态(任务, 租约), 文件集合, 回执
@pytest.mark.case_id(
"TC-raw-orphan-scope",
environment="真实隔离PostgreSQL、CLI与受控目录;假时钟只固定活跃租约",
given="本命名空间无租约对象、已登记租约、未知资料和另一个命名空间根",
when="经真实CLI盘点无租约暂存并再次盘点",
then=[
"无租约对象只登记为待核对、不写意图不删字节;已登记对象、未知资料与其他命名空间不变,不伪造清理回执"
],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
def test_raw_orphan_scope__a82112(原文环境, tmp_path):
"""当前用例 a82112;历史 LC-d6e3f58513a1 来源未知,不作为覆盖证据。"""
import json
import subprocess
import sys
import uuid
服务, 任务, 授权, 数据, 哈希, 时钟, 库 = 原文环境
租约 = 服务.创建租约(任务, 授权, "kept", 时钟.现在() + timedelta(hours=1), 最少剩余秒=30)
服务.写入(任务, 租约, 哈希, 数据)
孤儿 = str(uuid.uuid4())
原孤儿字节 = b"secret" # 旧例的合成字节,仅模拟缺租约残留。
孤儿哈希 = hashlib.sha256(原孤儿字节).hexdigest()
服务.文件.写入(孤儿, 孤儿哈希, 原孤儿字节)
外根 = 受控文件(tmp_path / "other-namespace")
外根.绑定命名空间(str(uuid.uuid4()), "production")
外根.写入(孤儿, 孤儿哈希, 原孤儿字节)
未知 = 服务.文件.根 / "合成资料.md"
未知.write_text("不属于原文对象的合成资料")
应用 = tmp_path / "cleanup.toml"
应用.write_text(
'["数据库"]\n"取值方式"="受控存储"\n"位置"='
+ json.dumps(库[用途.生产].引用.位置, ensure_ascii=False)
+ '\n["资源"]\n"发布身份"="test"\n["文件"]\n"原文暂存"='
+ json.dumps(str(服务.文件.根), ensure_ascii=False)
+ "\n"
)
结果 = subprocess.run(
[sys.executable, "-I", "-m", "muse", "管理", str(应用), "清理无租约暂存"],
cwd=tmp_path,
capture_output=True,
text=True,
)
assert 结果.returncode == 0, 结果.stderr
报告 = json.loads(结果.stdout)
assert 报告["preview_only"] is True and 报告["reason"] == "unmapped_lease_requires_review"
assert 报告["needs_recovery"] == [孤儿]
assert 报告["cleaned_orphan_ids"] == [] and 报告["receipts"] == []
assert 报告["unknown_entries"] == 1
assert (服务.文件.根 / 孤儿).exists() and 服务.读取(任务, 租约, 哈希) == 数据
assert (
外根.读取(孤儿, 孤儿哈希) == 原孤儿字节 and 未知.read_text() == "不属于原文对象的合成资料"
)
with pytest.raises(原文错误, match="没有此对象的清理回执"):
服务.读取孤儿清理回执(孤儿)
再次 = 服务.清理无租约原文()
assert 再次["needs_recovery"] == [孤儿] and 再次["cleaned_orphan_ids"] == []
assert (服务.文件.根 / 孤儿).exists()
with pytest.raises(受控文件错误, match="其他数据库"):
原文服务(库[用途.生产], 外根).清理无租约原文()
assert 外根.读取(孤儿, 孤儿哈希) == 原孤儿字节
@pytest.mark.case_id(
"NC-raw-orphan-cleanup-recovery",
environment="真实隔离PG与受控目录",
given="本命名空间无租约对象与一个有效租约",
when="经旧维护入口盘点无租约暂存并重复盘点",
then=[
"只登记待人工核对对象:不写清理意图、不删字节、不伪造清理回执;重复盘点结果一致,活跃原文保留"
],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
def test_孤儿清理先保存意图且确认失败可幂等恢复__a82111(原文环境):
"""旧入口现只盘点;具名到期清理的删除与幂等回执见 tests/单元/test_原文到期回收.py。"""
import uuid
服务, 任务, 授权, 数据, 哈希, 时钟, 库 = 原文环境
租约 = 服务.创建租约(
任务, 授权, "kept-for-namespace", 时钟.现在() + timedelta(hours=1), 最少剩余秒=0
)
服务.写入(任务, 租约, 哈希, 数据)
孤儿 = str(uuid.uuid4())
服务.文件.写入(孤儿, 哈希, 数据)
首次 = 服务.清理无租约原文()
assert 首次["preview_only"] is True and 首次["needs_recovery"] == [孤儿]
assert 首次["cleaned_orphan_ids"] == [] and 首次["receipts"] == []
assert (服务.文件.根 / 孤儿).exists()
with pytest.raises(原文错误):
服务.读取孤儿清理回执(孤儿)
再次 = 服务.清理无租约原文()
assert 再次 == 首次
assert (服务.文件.根 / 孤儿).exists()
assert 服务.读取(任务, 租约, 哈希) == 数据
def 建证据(环境):
原文, 任务, _, 数据, 哈希, _, 库 = 环境
登记 = 流程登记()
登记.登记处理器(步骤处理器("raw-test", "1", lambda _: 步骤结果({}), "1", "1"))
运行 = 任务服务(库[用途.生产], 登记)
领取 = 运行.领取步骤("evidence-worker", ["raw-test"])
assert 领取 is not None and 领取.任务ID == 任务
return 证据服务(库[用途.生产]), 领取
@pytest.mark.case_id(
"NC-raw-failed-task-retention",
environment="真实隔离PostgreSQL与受控文件,假时钟固定期限",
given="任务失败,但原文租约仍有效且字节已保存",
when="恢复原文状态",
then=["保留有效租约及字节;failed不等于取消或完成,不提前删除恢复输入"],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
def test_失败任务仍可恢复所以保留有效原文租约__a82113(原文环境):
from muse.任务运行.接口 import 任务状态
服务, 任务, 授权, 数据, 哈希, 时钟, 库 = 原文环境
租约 = 服务.创建租约(
任务, 授权, "recoverable-failure", 时钟.现在() + timedelta(hours=1), 最少剩余秒=30
)
服务.写入(任务, 租约, 哈希, 数据)
登记 = 流程登记()
登记.登记处理器(步骤处理器("raw-test", "1", lambda _: 步骤结果({}), "1", "1"))
运行 = 任务服务(库[用途.生产], 登记)
领取 = 运行.领取步骤("failed-worker", ["raw-test"])
assert 领取 is not None
运行.失败步骤(领取, "synthetic_failure")
assert 运行.读取任务(任务).状态 == 任务状态.已失败
assert 服务.恢复任务原文(任务, "author") == [{"lease_id": 租约, "action": "retained"}]
assert 服务.读取(任务, 租约, 哈希) == 数据
@pytest.mark.case_id(
"NC-evidence-same-hash-backfill",
environment="隔离 PostgreSQL 与合成原文",
given="真实任务尝试已保存失败响应的哈希",
when="持久授权后按同一证据ID补交字节",
then=["拒绝临时授权和不同字节;补交幂等且结果仍失败,不建立新尝试"],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
def test_仅哈希补交完整证据保持失败与同一身份__a82007(原文环境) -> None:
原文, 任务, 临时授权, 数据, 哈希, 时钟, 库 = 原文环境
证据, 领取 = 建证据(原文环境)
元信息 = {"call_id": "call-failed", "error_code": "MODEL_INCOMPLETE", "output_tokens": None}
回执 = 证据.保存(任务, 领取.尝试ID, "model_response", 哈希, "failed", 元信息)
assert 回执["retention"] == "hash_only" and 回执["outcome"] == "failed"
with pytest.raises(原文错误, match="^完整运行证据缺少对应的持久保留授权$"):
证据.补交原文(任务, 回执["evidence_id"], 数据, 临时授权)
授权 = 原文.批准保留(
任务,
"author",
"persistent",
来源版本="source-v1",
哈希=(哈希,),
内容用途="detection",
方式="persistent",
有效期=时钟.现在() + timedelta(days=1),
)
with pytest.raises(原文错误, match="^补交字节与原证据哈希不一致$"):
证据.补交原文(任务, 回执["evidence_id"], 数据 + b"changed", 授权)
补交 = 证据.补交原文(任务, 回执["evidence_id"], 数据, 授权)
assert 补交 == {**回执, "retention": "full", "revision": 2}
assert 证据.补交原文(任务, 回执["evidence_id"], 数据, 授权) == 补交
with 库[用途.生产].连接() as 连:
assert (
连.execute("SELECT count(*) FROM muse_attempt WHERE task_id=%s", (任务,)).fetchone()[0]
== 1
)
assert (
连.execute(
"SELECT content FROM muse_runtime_evidence WHERE evidence_id=%s",
(回执["evidence_id"],),
).fetchone()[0]
== 数据
)
@pytest.mark.case_id(
"NC-evidence-outcome-immutable",
environment="隔离 PostgreSQL 与合成原文",
given="真实任务的partial证据",
when="尝试覆写成功状态或夹带正文元信息",
then=["拒绝且原回执不变"],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
def test_证据不能覆盖结果或写入未登记正文元信息__a82008(原文环境) -> None:
_, 任务, _, _, 哈希, _, _ = 原文环境
证据, 领取 = 建证据(原文环境)
元信息 = {"call_id": "call-partial"}
回执 = 证据.保存(任务, 领取.尝试ID, "model_response", 哈希, "partial", 元信息)
with pytest.raises(原文错误):
证据.保存(任务, 领取.尝试ID, "model_response", 哈希, "completed", 元信息)
with pytest.raises(原文错误):
证据.保存(任务, 领取.尝试ID, "failure", 哈希, "failed", {"source_refs": [{"raw": "正文"}]})
assert 证据.读取回执(任务, 回执["evidence_id"]) == 回执
@pytest.fixture
def 原文环境(应用测试库, tmp_path):
登记 = 流程登记()
登记.登记处理器(步骤处理器("raw-test", "1", lambda _: 步骤结果({}), "1", "1"))
登记.登记类型("raw-test", 必需保护=())
任务 = 任务服务(应用测试库[用途.生产], 登记)
流程 = 流程服务(任务, 登记)
流程.发布(流程定义("raw-test", "1", (步骤计划("raw", "raw-test", "1"),)))
id_ = 任务.创建任务(
任务请求(
"raw-test",
"cmd",
"author",
用途.生产,
内容用途.检测,
{},
"role-1",
"res-1",
{
"source_scope": {},
"schema_versions": {},
"authorization": "grant",
"budget": {},
"stop_conditions": [],
},
),
"raw-test",
"1",
)
时钟 = 假时钟(datetime.now(UTC))
文件 = 受控文件(tmp_path / "私人原文")
服务 = 原文服务(应用测试库[用途.生产], 文件, 时钟)
数据 = "合成原文,保留每一个字节。".encode()
哈希 = hashlib.sha256(数据).hexdigest()
授权 = 服务.批准保留(
id_,
"author",
"approve",
来源版本="source-v1",
哈希=(哈希,),
内容用途="detection",
方式="temporary",
有效期=时钟.现在() + timedelta(hours=2),
)
return 服务, id_, 授权, 数据, 哈希, 时钟, 应用测试库
@pytest.mark.case_id(
"NC-raw-lease-expiry",
environment="隔离 PostgreSQL、临时合成原文、必要时真实ASGI会话",
given="确切授权哈希和假时钟",
when="创建租约并从独立PG连接核对,再写入、推进至到期、重复清理",
then=["租约先提交,未授权时没有字节;过期不能读写,清理可重放"],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
def test_租约先提交才写字节且过期仍可清理__a82001(原文环境) -> None:
服务, 任务, 授权, 数据, 哈希, 时钟, 库 = 原文环境
with pytest.raises(原文错误):
服务.创建租约(任务, 授权, "too-long", 时钟.现在() + timedelta(hours=25), 最少剩余秒=60)
assert list(服务.文件.根.iterdir()) == []
租约 = 服务.创建租约(任务, 授权, "lease", 时钟.现在() + timedelta(hours=1), 最少剩余秒=60)
with 库[用途.生产].连接() as 连:
assert (
连.execute("SELECT state FROM muse_raw_lease WHERE lease_id=%s", (租约,)).fetchone()[0]
== "open"
)
assert list(服务.文件.根.iterdir()) == []
服务.写入(任务, 租约, 哈希, 数据)
assert 服务.读取(任务, 租约, 哈希) == 数据
时钟.推进(timedelta(hours=1))
with pytest.raises(原文错误):
服务.读取(任务, 租约, 哈希)
with pytest.raises(原文错误):
服务.写入(任务, 租约, 哈希, 数据)
assert 服务.清理(任务, 租约) is True
首次 = _清理快照(服务, 任务, 租约, 库)
_断言租约状态(首次[0], 租约, "closed", 已清理=True)
assert 首次[2][0] is not None
assert not (服务.文件.根 / 租约).exists()
assert 服务.清理(任务, 租约) is False
assert _清理快照(服务, 任务, 租约, 库) == 首次
@pytest.mark.case_id(
"NC-raw-pg-archive",
environment="隔离 PostgreSQL、临时合成原文、必要时真实ASGI会话",
given="有临时原文与独立归档授权",
when="两次按同租约归档",
then=["完整字节在PG,回执相同,临时副本删除而归档保留"],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
def test_完整PG归档可重复对账并清理临时副本__a82002(原文环境) -> None:
服务, 任务, 授权, 数据, 哈希, 时钟, 库 = 原文环境
租约 = 服务.创建租约(任务, 授权, "lease", 时钟.现在() + timedelta(hours=1), 最少剩余秒=60)
服务.写入(任务, 租约, 哈希, 数据)
归档授权 = 服务.批准保留(
任务,
"author",
"archive",
来源版本="source-v1",
哈希=(哈希,),
内容用途="detection",
方式="archive",
有效期=时钟.现在() + timedelta(days=1),
)
回执 = 服务.归档(任务, 租约, 归档授权)
assert 回执["entry_count"] == 1 and 回执["total_bytes"] == len(数据)
assert 服务.归档(任务, 租约, 归档授权) == 回执
_断言租约状态(服务.状态(任务, 租约), 租约, "migrated", 归档=回执["archive_id"], 已清理=True)
assert not (服务.文件.根 / 租约).exists()
assert 服务.读取归档(任务, 回执["archive_id"], 哈希) == 数据
assert 原文服务(库[用途.生产], 时钟实现=时钟).读取归档(任务, 回执["archive_id"], 哈希) == 数据
时钟.推进(timedelta(hours=2))
首次 = _清理快照(服务, 任务, 租约, 库)
assert 首次[2][0] is not None
assert 服务.清理(任务, 租约) is False
assert _清理快照(服务, 任务, 租约, 库) == 首次
assert 服务.读取归档(任务, 回执["archive_id"], 哈希) == 数据
with 库[用途.生产].连接() as 连:
assert (
连.execute(
"SELECT content FROM muse_raw_archive_item WHERE archive_id=%s",
(回执["archive_id"],),
).fetchone()[0]
== 数据
)
@pytest.mark.case_id(
"NC-raw-missing-migrating",
environment="隔离 PostgreSQL、临时合成原文、必要时真实ASGI会话",
given="已授权租约缺少完整字节",
when="归档失败后推进到期并清理",
then=["migrating保留,缺字节不成功,普通清理拒绝"],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
def test_缺字节归档保留migrating不被过期删除__a82003(原文环境) -> None:
服务, 任务, 授权, 数据, 哈希, 时钟, _ = 原文环境
租约 = 服务.创建租约(任务, 授权, "lease", 时钟.现在() + timedelta(hours=1), 最少剩余秒=60)
归档授权 = 服务.批准保留(
任务,
"author",
"archive",
来源版本="source-v1",
哈希=(哈希,),
内容用途="detection",
方式="archive",
有效期=时钟.现在() + timedelta(days=1),
)
with pytest.raises(受控文件错误):
服务.归档(任务, 租约, 归档授权)
原状态 = 服务.状态(任务, 租约)
assert 原状态["archive_id"] is not None
_断言租约状态(原状态, 租约, "migrating", 归档=原状态["archive_id"])
assert 服务.恢复任务原文(任务, "author")[0]["action"] == "needs_recovery"
时钟.推进(timedelta(hours=2))
with pytest.raises(原文错误, match="^迁移中的原文只能先对账恢复$"):
服务.清理(任务, 租约)
assert 服务.状态(任务, 租约) == 原状态
@pytest.mark.case_id(
"NC-raw-transaction-recovery",
environment="隔离 PostgreSQL、临时合成原文、必要时真实ASGI会话",
given="真实PG注入归档条目写入失败",
when="归档后解除故障并恢复",
then=["header和bytea同事务回滚,源字节保留,原archive_id恢复"],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
def test_归档事务失败不留半归档且同ID恢复__a82004(原文环境) -> None:
服务, 任务, 授权, 数据, 哈希, 时钟, 库 = 原文环境
租约 = 服务.创建租约(任务, 授权, "lease", 时钟.现在() + timedelta(hours=1), 最少剩余秒=60)
服务.写入(任务, 租约, 哈希, 数据)
归档授权 = 服务.批准保留(
任务,
"author",
"archive",
来源版本="source-v1",
哈希=(哈希,),
内容用途="detection",
方式="archive",
有效期=时钟.现在() + timedelta(days=1),
)
with 库[用途.维护].连接() as 连:
连.execute("""CREATE FUNCTION reject_raw() RETURNS trigger LANGUAGE plpgsql AS $$
BEGIN RAISE EXCEPTION 'synthetic archive failure'; END $$;
CREATE TRIGGER reject_raw BEFORE INSERT ON muse_raw_archive_item
FOR EACH ROW EXECUTE FUNCTION reject_raw();""")
with pytest.raises(原文错误):
服务.归档(任务, 租约, 归档授权)
原状态 = 服务.状态(任务, 租约)
assert 原状态["archive_id"] is not None
_断言租约状态(原状态, 租约, "migrating", 归档=原状态["archive_id"])
assert 服务.文件.读取(租约, 哈希) == 数据
with 库[用途.维护].连接() as 连:
assert 连.execute("SELECT count(*) FROM muse_raw_archive").fetchone()[0] == 0
assert 连.execute("SELECT count(*) FROM muse_raw_archive_item").fetchone()[0] == 0
连.execute("DROP TRIGGER reject_raw ON muse_raw_archive_item")
assert 服务.归档(任务, 租约, 归档授权)["archive_id"] == 原状态["archive_id"]
@pytest.mark.case_id(
"NC-raw-cleanup-author",
environment="隔离 PostgreSQL、临时合成原文、必要时真实ASGI会话",
given="任务作者与伪作者,受控目录权限漂移",
when="批准及清理后恢复权限重试",
then=["伪作者不能批准;清理失败不closed,不产生清理成功回执"],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
def test_清理失败不关闭且错误作者不能批准__a82005(原文环境) -> None:
服务, 任务, 授权, 数据, 哈希, 时钟, _ = 原文环境
with pytest.raises(原文错误):
服务.批准保留(
任务,
"创始人",
"forged",
来源版本="source-v1",
哈希=(哈希,),
内容用途="detection",
方式="archive",
有效期=时钟.现在() + timedelta(days=1),
)
租约 = 服务.创建租约(任务, 授权, "lease", 时钟.现在() + timedelta(hours=1), 最少剩余秒=60)
服务.写入(任务, 租约, 哈希, 数据)
目录 = 服务.文件.根 / 租约
目录.chmod(0o755)
try:
with pytest.raises(受控文件错误):
服务.清理(任务, 租约, 作者="author")
_断言租约状态(服务.状态(任务, 租约), 租约, "open")
finally:
目录.chmod(0o700)
服务.清理(任务, 租约, 作者="author")
assert 服务.状态(任务, 租约)["state"] == "closed"
@pytest.mark.case_id(
"NC-raw-authenticated-http",
environment="隔离 PostgreSQL、临时合成原文、必要时真实ASGI会话",
given="真实ASGI会话与隔离PG任务",
when="认证、批准临时保存、创建租约、写入字节、单独批准归档",
then=["未认证或自报approved_by拒绝;完整授权链归档和清理真实生效"],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
def test_HTTP作者授权到原文写入归档使用真实会话__a82006(原文环境, tmp_path) -> None:
from fastapi.testclient import TestClient
from muse.接入.http.应用 import 创建应用
from muse.配置 import 应用配置, 服务配置
_, 任务, _, 数据, 哈希, _, 库 = 原文环境
口令 = tmp_path / "口令.txt"
口令.write_text("synthetic-raw-author")
配置 = 应用配置(
库[用途.生产].引用,
"test",
HTTP=服务配置(
str(口令),
作者ID="author",
公开地址="http://testserver",
允许来源=("http://testserver",),
),
原文暂存=str(tmp_path / "HTTP原文"),
)
base = f"/api/v1/tasks/{任务}"
body = {
"command_id": "http-approve",
"source_version": "source-v1",
"content_hashes": [哈希],
"content_purpose": "detection",
"retention_mode": "temporary",
"valid_until": (datetime.now(UTC) + timedelta(hours=2)).isoformat(),
}
with TestClient(创建应用(配置)) as client:
client.headers["Origin"] = "http://testserver"
assert client.post(base + "/raw-authorizations", json=body).status_code == 401
assert (
client.post("/api/v1/session", json={"password": "synthetic-raw-author"}).status_code
== 200
)
assert (
client.post(
base + "/raw-authorizations", json={**body, "approved_by": "创始人"}
).status_code
== 422
)
授权 = client.post(base + "/raw-authorizations", json=body)
assert 授权.status_code == 200, 授权.text
租约 = client.post(
base + "/raw-leases",
json={
"authorization_id": 授权.json()["authorization_id"],
"command_id": "http-lease",
"retain_until": (datetime.now(UTC) + timedelta(hours=1)).isoformat(),
"min_remaining": 60,
},
)
assert 租约.status_code == 200, 租约.text
lease = 租约.json()["lease_id"]
write = client.put(
base + f"/raw-leases/{lease}/content/{哈希}",
content=数据,
headers={"Content-Type": "application/octet-stream"},
)
assert write.status_code == 200 and write.json()["bytes"] == len(数据)
archive = client.post(
base + "/raw-authorizations",
json={**body, "command_id": "http-archive", "retention_mode": "archive"},
)
result = client.post(
base + f"/raw-leases/{lease}/archive",
json={"authorization_id": archive.json()["authorization_id"]},
)
assert result.status_code == 200 and result.json()["entry_count"] == 1
assert not (tmp_path / "HTTP原文" / lease).exists()
@pytest.mark.case_id(
"NC-raw-legacy-import",
environment="隔离 PostgreSQL 与合成旧文件归档",
given="旧格式回执和三份原文件,其中两份内容相同;作者明确批准完整哈希与来源",
when="完整导入PG、重试,再篡改旧来源",
then=["保留旧载体、全部文件映射与原回执,按字节去重但不丢名称;同源重试回执相同;改动来源拒绝"],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
def test_历史归档完整导入PG且保留旧载体__a82009(原文环境, tmp_path):
import json
from muse.基础设施.受控文件 import 受控文件错误
服务, 任务, _, _, _, 时钟, 库 = 原文环境
目录 = tmp_path / "历史归档"
目录.mkdir()
字节 = {
"input.txt": "旧输入。".encode(),
"output.txt": "旧输出。".encode(),
"same-output.txt": "旧输出。".encode(),
}
清单 = []
for name, data in sorted(字节.items()):
(目录 / name).write_bytes(data)
清单.append(
{
"relativePath": name,
"contentSha256": "sha256:" + hashlib.sha256(data).hexdigest(),
"sizeBytes": len(data),
}
)
tree = hashlib.sha256(
json.dumps(清单, ensure_ascii=False, sort_keys=True, separators=(",", ":")).encode()
).hexdigest()
old = {
"schemaVersion": "raw-vault-migration-receipt-v1",
"archiveId": "legacy-1",
"status": "migrated",
"archiveSha256": "sha256:" + tree,
"fileCount": 3,
"totalBytes": sum(map(len, 字节.values())),
}
(目录 / ".migration-receipt.json").write_text(json.dumps(old, ensure_ascii=False))
哈希 = tuple(sorted({hashlib.sha256(v).hexdigest() for v in 字节.values()}))
授权 = 服务.批准保留(
任务,
"author",
"legacy-approve",
来源版本="legacy-source-1",
哈希=哈希,
内容用途="detection",
方式="archive",
有效期=时钟.现在() + timedelta(days=1),
)
回执 = 服务.导入历史归档(任务, 授权, "legacy-source-1", 目录)
assert 回执 == 服务.导入历史归档(任务, 授权, "legacy-source-1", 目录)
assert 回执["entry_count"] == 2
with 库[用途.生产].连接() as 连:
assert (
连.execute(
"SELECT legacy_manifest FROM muse_raw_archive WHERE archive_id=%s",
(回执["archive_id"],),
).fetchone()[0]["entries"]
== 清单
)
for name, data in 字节.items():
assert (目录 / name).read_bytes() == data
assert 服务.读取归档(任务, 回执["archive_id"], hashlib.sha256(data).hexdigest()) == data
(目录 / "output.txt").write_bytes(b"changed")
with pytest.raises(受控文件错误):
服务.导入历史归档(任务, 授权, "legacy-source-1", 目录)
@pytest.mark.case_id(
"TC-raw-unknown-commit-reconcile",
environment="隔离 PostgreSQL 与临时目录,假时钟",
given="PG 已提交归档但客户端未收到结果",
when="同 archive_id 恢复",
then=["核对完整字节哈希和原回执,返回同一结果,不重复归档、不删唯一副本"],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
def test_PG已提交但响应丢失以原归档身份恢复__a82010(原文环境):
import psycopg
服务, 任务, 授权, 数据, 哈希, 时钟, 库 = 原文环境
租约 = 服务.创建租约(任务, 授权, "lease", 时钟.现在() + timedelta(hours=1), 最少剩余秒=60)
服务.写入(任务, 租约, 哈希, 数据)
归档授权 = 服务.批准保留(
任务,
"author",
"archive",
来源版本="source-v1",
哈希=(哈希,),
内容用途="detection",
方式="archive",
有效期=时钟.现在() + timedelta(days=1),
)
原工厂 = 服务.数据库
class 丢失确认工厂:
用途 = 原工厂.用途
次数 = 0
def 连接(self):
self.次数 += 1
次数 = self.次数
实际 = 原工厂.连接()
class 连接上下文:
def __enter__(self):
return 实际.__enter__()
def __exit__(self, *args):
result = 实际.__exit__(*args)
if 次数 == 2 and args[0] is None:
raise psycopg.OperationalError(
"synthetic acknowledgement loss after commit"
)
return result
return 连接上下文()
服务.数据库 = 丢失确认工厂()
try:
with pytest.raises(原文错误):
服务.归档(任务, 租约, 归档授权)
finally:
服务.数据库 = 原工厂
assert 服务.文件.读取(租约, 哈希) == 数据
旧状态 = 服务.状态(任务, 租约)
_断言租约状态(旧状态, 租约, "migrated", 归档=旧状态["archive_id"])
恢复 = 服务.归档(任务, 租约, 归档授权)
assert 恢复["archive_id"] == 旧状态["archive_id"]
assert 服务.状态(任务, 租约)["cleanup_complete"]
with 库[用途.生产].连接() as 连:
assert 连.execute("SELECT count(*) FROM muse_raw_archive").fetchone()[0] == 1
@pytest.mark.case_id(
"TC-raw-open-recovery",
environment="隔离 PostgreSQL 与临时目录,假时钟",
given="终止任务分别有遗留 open 字节和缺失字节",
when="执行恢复",
then=["只清理本管理器未归档的终止临时原文;分别记录已清理或缺失,不声称证据仍完整"],
contract="docs/系统架构/新版设计/接口契约/原文证据生命周期.md",
)
@pytest.mark.parametrize("缺失", [False, True], ids=["retained-bytes", "missing-bytes"])
def test_终止任务恢复清理尚未到期原文__a82011(原文环境, 缺失):
from muse.任务运行.接口 import 任务状态
服务, 任务, 授权, 数据, 哈希, 时钟, 库 = 原文环境
租约 = 服务.创建租约(任务, 授权, "lease", 时钟.现在() + timedelta(hours=1), 最少剩余秒=60)
服务.写入(任务, 租约, 哈希, 数据)
if 缺失:
服务.文件.删除对象(租约)
运行 = 任务服务(库[用途.生产], 流程登记())
运行.控制任务(任务, "author", 任务状态.待运行, "取消", 命令ID="cancel")
assert 服务.恢复任务原文(任务, "author") == [
{"lease_id": 租约, "action": "missing" if 缺失 else "cleaned"}
]
assert 服务.状态(任务, 租约)["state"] == "closed"
assert not (服务.文件.根 / 租约).exists()
with 库[用途.生产].连接() as 连:
失败码 = 连.execute(
"SELECT failure_code FROM muse_raw_lease WHERE lease_id=%s", (租约,)
).fetchone()[0]
assert 失败码 == ("RAW_CONTENT_MISSING" if 缺失 else None)