"""到期发现仅元数据;具名清理以合成租约/时钟/文件验证拒绝与重试。""" import copy import hashlib import uuid from contextlib import contextmanager from datetime import UTC, datetime, timedelta from types import SimpleNamespace import pytest from muse.任务运行 import 原文生命周期 as 模块 from muse.任务运行.原文生命周期 import 原文回收判定, 原文服务, 原文错误 from muse.共享.调用身份 import 用途 from muse.基础设施.受控文件 import 受控文件, 受控文件错误 当前 = datetime(2026, 9, 17, tzinfo=UTC) def 元信息(**修改): return { "task_id": str(uuid.uuid4()), "lease_id": str(uuid.uuid4()), "authorization_id": str(uuid.uuid4()), "namespace_id": str(uuid.uuid4()), "created_at": 当前 - timedelta(hours=2), "retain_until": 当前 - timedelta(hours=1), "content_hashes": [hashlib.sha256(b"synthetic").hexdigest()], "response_hash": None, "state": "open", "task_state": "cancelled", "author_id": "author", "archive_id": None, "cleanup_receipt": None, "pending_cost": False, "active_attempt": False, **修改, } class 合成原文服务(原文服务): def __init__(self, 行, 文件=None): super().__init__(SimpleNamespace(用途=用途.生产), 文件, SimpleNamespace(现在=lambda: 当前)) self.行 = 行 self.查询次数 = 0 self.归档核对次数 = 0 @contextmanager def _事务(self): yield object() def _租约(self, 连, 租约ID, 任务ID): assert (租约ID, 任务ID) == (self.行["lease_id"], self.行["task_id"]) return self.行 def _回收行(self, 连, 任务ID, 租约ID): return self._租约(连, 租约ID, 任务ID) def _查询(self, 连, SQL, 参数=()): self.查询次数 += 1 if SQL.startswith("UPDATE"): if self.行["state"] != "migrated": self.行["state"] = "closed" self.行["cleanup_receipt"] = self.行["cleanup_receipt"] or 参数[0] return SimpleNamespace(fetchone=lambda: (1,) if self.行["pending_cost"] else None) def _绑定文件(self, 连): assert self.文件 is not None self.文件.绑定命名空间(self.行["namespace_id"], "production") return self.文件 def _核对归档(self, 连, 归档ID): self.归档核对次数 += 1 assert 归档ID == self.行["archive_id"] return {"archive_id": 归档ID} @pytest.fixture def 环境(tmp_path, monkeypatch): 行 = 元信息() 文件 = 受控文件(tmp_path / "raw") 文件.绑定命名空间(行["namespace_id"], "production") 文件.写入(行["lease_id"], 行["content_hashes"][0], b"synthetic") 服务 = 合成原文服务(行, 文件) monkeypatch.setattr( 模块, "任务存储", lambda *a: SimpleNamespace(任务行=lambda *a, **k: {"author_id": "author"}) ) return 服务, 行, 文件 @pytest.mark.case_id("NC-raw-retention-policy") @pytest.mark.parametrize( "修改,decision,reason", [ ({}, "candidate", "expired_temporary"), ({"retain_until": 当前 + timedelta(seconds=1)}, "keep", "not_expired"), ({"state": "migrating"}, "keep", "archive_reconciliation_required"), ({"namespace_id": None}, "unknown", "namespace_missing"), ({"task_state": "failed"}, "keep", "task_can_resume"), ({"pending_cost": True}, "keep", "execution_or_accounting_pending"), ({"active_attempt": True}, "keep", "execution_or_accounting_pending"), ({"state": "migrated"}, "unknown", "archive_identity_missing"), ], ) def test_期限不替代租约引用与费用保护(修改, decision, reason): 行 = 元信息(**修改) 原 = copy.deepcopy(行) 判定 = 原文回收判定(行, 当前, "production") assert (判定["decision"], 判定["reason"]) == (decision, reason) assert 行 == 原 @pytest.mark.case_id("NC-raw-retention-idempotent") def test_具名执行删除临时文件并返回稳定回执(环境): 服务, 行, 文件 = 环境 计划 = 原文回收判定(行, 当前, "production") 首次 = 服务.执行到期清理("author", 计划) assert 首次["already_cleaned"] is False and 首次["receipt_id"] assert not (文件.根 / 行["lease_id"]).exists() 再次 = 服务.执行到期清理("author", 计划) assert 再次 == {**首次, "already_cleaned": True} assert 行["state"] == "closed" @pytest.mark.case_id("NC-raw-retention-drift") @pytest.mark.parametrize("改变", ["pending_cost", "authorization", "bytes", "author"]) def test_执行前状态身份字节或作者变化一律保留(环境, 改变): 服务, 行, 文件 = 环境 计划 = 原文回收判定(行, 当前, "production") 作者 = "author" if 改变 == "pending_cost": 行["pending_cost"] = True if 改变 == "authorization": 行["authorization_id"] = str(uuid.uuid4()) if 改变 == "bytes": (文件.根 / 行["lease_id"] / (行["content_hashes"][0] + ".bin")).write_bytes(b"changed") if 改变 == "author": 作者 = "another-author" with pytest.raises((原文错误, 受控文件错误)): 服务.执行到期清理(作者, 计划) assert (文件.根 / 行["lease_id"]).is_dir() assert 行["cleanup_receipt"] is None @pytest.mark.case_id("NC-raw-retention-archive") def test_已归档先核归档只删临时目录(环境): 服务, 行, 文件 = 环境 行.update(state="migrated", archive_id=str(uuid.uuid4())) 计划 = 原文回收判定(行, 当前, "production") 服务.执行到期清理("author", 计划) assert 服务.归档核对次数 == 1 assert 行["state"] == "migrated" and 行["archive_id"] assert not (文件.根 / 行["lease_id"]).exists() @pytest.mark.case_id("NC-raw-retention-paging") def test_分页发现单次查询不访问文件且删除后游标不漏(环境): 服务, 行, 文件 = 环境 行集 = sorted([元信息() for _ in range(5)], key=lambda r: r["lease_id"]) 查询次数 = [] def 查询(连, SQL, 参数): 查询次数.append(SQL) 作者, 用途值, 截止, 游标时间, 游标ID, 上限 = 参数 result = [ r for r in 行集 if (r["retain_until"], r["lease_id"]) > (游标时间, 游标ID) and r["retain_until"] <= 截止 and r["cleanup_receipt"] is None ] return SimpleNamespace(fetchall=lambda: result[:上限]) 服务._查询 = 查询 服务._绑定文件 = lambda _: pytest.fail("发现不访问文件") 第一页 = 服务.发现到期清理候选("author", 截至=当前, 上限=2) assert len(查询次数) == 1 and len(第一页["items"]) == 2 行集[0]["cleanup_receipt"] = "completed-after-preview" 第二页 = 服务.发现到期清理候选("author", 截至=当前, 上限=2, 游标=第一页["next_cursor"]) 第三页 = 服务.发现到期清理候选("author", 截至=当前, 上限=2, 游标=第二页["next_cursor"]) assert [r["lease_id"] for page in [第一页, 第二页, 第三页] for r in page["items"]] == [ r["lease_id"] for r in 行集 ] assert 第三页["next_cursor"] is None @pytest.mark.case_id("NC-raw-retention-direct-cost-guard") def test_显式清理也不能绕过未知费用保护(环境): 服务, 行, 文件 = 环境 行["pending_cost"] = True with pytest.raises(原文错误, match="未知费用"): 服务.清理(行["task_id"], 行["lease_id"], 作者="author") assert 文件.读取(行["lease_id"], 行["content_hashes"][0]) == b"synthetic"