diff --git a/src/muse/任务运行/原文生命周期.py b/src/muse/任务运行/原文生命周期.py new file mode 100644 index 0000000..9cc5b1c --- /dev/null +++ b/src/muse/任务运行/原文生命周期.py @@ -0,0 +1,880 @@ +"""原文保留授权、最长24小时租约与完整PG归档;不接管正式正文。""" + +from __future__ import annotations + +import hashlib +import re +import uuid +from contextlib import contextmanager +from datetime import datetime +from pathlib import Path +from typing import LiteralString + +import psycopg +from psycopg import sql +from psycopg.rows import dict_row + +from muse.任务运行.任务领取 import 校验领取 +from muse.任务运行.存储 import 任务存储 +from muse.任务运行.模型 import 内容哈希, 领取凭证 +from muse.共享.时间 import 时钟, 真实时钟 +from muse.共享.调用身份 import 用途 +from muse.共享.错误 import Muse错误 +from muse.基础设施.受控文件 import 受控文件, 读取历史归档 +from muse.基础设施.数据库.连接 import 数据库工厂 + + +class 原文错误(Muse错误): + 错误码 = "RAW_LIFECYCLE" + + +def 校验原文字节(数据: bytes) -> None: + """拒绝凭据赋值和私钥,普通token篇幅说明保留;错误不引用命中正文。""" + 模式 = ( + r"(?i)\b(?:api[_-]?key|apikey|secret|password|passwd|token)[\"']?\s*[:=]\s*[\"']?[^\s,\"'}]+", + r"(?i)authorization\s*:\s*bearer\s+\S+", + r"(?i)\bsk-[a-z0-9_-]{16,}\b", + r"(?i)-----begin private key-----", + ) + # Unicode词边界可区分“最大输出token”等正常字段;只检查,不改变原始字节。 + 文本 = 数据.decode("utf-8", errors="replace") + if any(re.search(p, 文本) for p in 模式): + raise 原文错误("原文含凭据赋值或私钥,拒绝保存") + + +def 校验期限(创建: datetime, 截止: datetime, 最少剩余秒: float) -> None: + if 创建.utcoffset() is None or 截止.utcoffset() is None: + raise 原文错误("原文时间必须带时区") + 剩余 = (截止 - 创建).total_seconds() + if not 0 < 剩余 <= 24 * 3600 or not 0 <= 最少剩余秒 <= 剩余: + raise 原文错误("原文期限超过24小时或不足本次执行与清理余量") + + +class 原文服务: + def __init__( + self, 数据库: 数据库工厂, 文件: 受控文件 | None = None, 时钟实现: 时钟 | None = None + ): + self.数据库 = 数据库 + self.文件 = 文件 + self.时钟 = 时钟实现 or 真实时钟() + self.schema = "evaluation" if 数据库.用途 is 用途.评测 else "public" + + def _查询(self, 连, 语句: LiteralString, 参数=()): + return 连.cursor(row_factory=dict_row).execute( + sql.SQL(语句).format(s=sql.Identifier(self.schema)), 参数 + ) + + @contextmanager + def _事务(self): + try: + with self.数据库.连接() as 连, 连.transaction(): + yield 连 + except psycopg.Error as exc: + raise 原文错误( + "原文账本操作未确认成功,保留原文等待恢复", 上下文={"sqlstate": exc.sqlstate} + ) from None + + def _绑定文件(self, 连) -> 受控文件: + if self.文件 is None: + raise 原文错误("临时原文操作需要明确的受控暂存根") + 行 = self._查询( + 连, "SELECT namespace_id FROM {s}.muse_raw_namespace WHERE singleton" + ).fetchone() + if 行 is None: + raise 原文错误("数据库尚未建立原文命名空间") + self.文件.绑定命名空间(str(行["namespace_id"]), self.数据库.用途.value) + return self.文件 + + def 批准保留( + self, + 任务ID: str, + 作者: str, + 命令ID: str, + *, + 来源版本: str, + 哈希: tuple[str, ...], + 内容用途: str, + 方式: str, + 有效期: datetime, + 调用ID: str | None = None, + 调用请求哈希: str | None = None, + ) -> str: + """作者由已认证接入取得;批准后保存不可变范围,不接受角色自报approved_by。""" + 当前 = self.时钟.现在() + if ( + not 命令ID + or not 来源版本 + or not 内容用途 + or not 哈希 + or any(not re.fullmatch(r"[0-9a-f]{64}", h) for h in 哈希) + ): + raise 原文错误("原文批准需要确切来源版本、用途和非空哈希集合") + if ( + 方式 not in {"temporary", "archive", "persistent"} + or 有效期.utcoffset() is None + or 有效期 <= 当前 + ): + raise 原文错误("原文保存方式或授权有效期无效") + 范围 = { + "source": 来源版本, + "hashes": sorted(set(哈希)), + "purpose": 内容用途, + "mode": 方式, + "until": 有效期.isoformat(), + "call_id": 调用ID, + "call_request_hash": 调用请求哈希, + } + if (调用ID is None) != (调用请求哈希 is None) or ( + 调用ID is not None and (not 调用ID or 调用请求哈希 not in 哈希) + ): + raise 原文错误("响应保留必须绑定确切调用与已批准的请求哈希") + with self._事务() as 连: + 任务 = 任务存储(连, self.数据库.用途).任务行(任务ID, 锁定=True) + if 任务["author_id"] != 作者: + raise 原文错误("原文批准者与任务作者不一致") + if 任务["frozen_input"]["内容用途"] != 内容用途: + raise 原文错误("原文授权不能改变任务冻结的内容用途") + 身份 = str(uuid.uuid4()) + self._查询( + 连, + ( + "INSERT INTO {s}.muse_raw_authorization " + "(authorization_id,task_id,command_id,approved_by,source_version," + "content_hashes,purpose,retention_mode,approved_at,valid_until,request_hash," + "call_id,call_request_hash) " + "VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s) ON CONFLICT DO NOTHING" + ), + ( + 身份, + 任务ID, + 命令ID, + 作者, + 来源版本, + 范围["hashes"], + 内容用途, + 方式, + 当前, + 有效期, + 内容哈希(范围), + 调用ID, + 调用请求哈希, + ), + ) + 行 = self._查询( + 连, + "SELECT * FROM {s}.muse_raw_authorization WHERE task_id=%s AND command_id=%s", + (任务ID, 命令ID), + ).fetchone() + if 行 is None or 行["request_hash"] != 内容哈希(范围): + raise 原文错误("原文批准命令对应其他内容范围") + return str(行["authorization_id"]) + + def _授权(self, 连, 授权ID: str, 任务ID: str, 方式: str, *, 恢复=False): + 行 = self._查询( + 连, + "SELECT * FROM {s}.muse_raw_authorization " + "WHERE authorization_id=%s AND task_id=%s FOR SHARE", + (授权ID, 任务ID), + ).fetchone() + if 行 is None or 行["revoked"] or 行["retention_mode"] != 方式: + raise 原文错误("缺少对应任务和保存方式的原文授权") + if 行["parent_authorization_id"]: + self._授权(连, str(行["parent_authorization_id"]), 任务ID, 方式, 恢复=恢复) + if not 恢复 and 行["valid_until"] <= self.时钟.现在(): + raise 原文错误("原文授权已过期") + if 行["response_hash"]: + 行 = {**行, "content_hashes": list(set(行["content_hashes"]) | {行["response_hash"]})} + return 行 + + def 检查调用授权( + self, 任务ID: str, 授权ID: str, 调用ID: str, 请求哈希: str, 方式: str, 所需秒: float + ) -> None: + with self._事务() as 连: + 行 = self._授权(连, 授权ID, 任务ID, 方式) + if (行["call_id"], 行["call_request_hash"]) != (调用ID, 请求哈希): + raise 原文错误("模型调用不在本次批准的固定请求范围") + if (行["valid_until"] - self.时钟.现在()).total_seconds() < 所需秒: + raise 原文错误("原文授权不足本次调用与保存余量") + + def 批准会话( + self, + 任务ID: str, + 作者: str, + 命令ID: str, + *, + 步骤ID: str, + 阶段: str, + 初始请求哈希: str, + 最大模型调用: int, + 最大工具调用: int, + 有效期: datetime, + ) -> str: + """由已认证作者批准有限角色会话;执行者只能派生该范围内的回合。""" + if ( + type(最大模型调用) is not int + or 最大模型调用 < 1 + or type(最大工具调用) is not int + or 最大工具调用 < 0 + ): + raise 原文错误("会话必须固定正数模型次数和非负工具次数") + with self._事务() as 连: + 任务 = 任务存储(连, self.数据库.用途).任务行(任务ID) + 步骤 = next((s for s in 任务["flow_snapshot"]["步骤"] if s["步骤ID"] == 步骤ID), None) + if 步骤 is None or not 步骤["角色"] or 任务["author_id"] != 作者: + raise 原文错误("会话必须属于本作者任务的具名角色步骤") + 固定哈希 = 内容哈希(任务["frozen_input"]) + 内容用途 = 任务["frozen_input"]["内容用途"] + 策略 = { + "step": 步骤ID, + "stage": 阶段, + "initial": 初始请求哈希, + "frozen": 固定哈希, + "model_calls": 最大模型调用, + "tool_calls": 最大工具调用, + } + # 普通批准未绑定call_id,只有下方会话记录成功后才能派生响应保留权。 + 授权 = self.批准保留( + 任务ID, + 作者, + 命令ID, + 来源版本=固定哈希, + 哈希=(初始请求哈希,), + 内容用途=内容用途, + 方式="persistent", + 有效期=有效期, + ) + with self._事务() as 连: + self._查询( + 连, + "INSERT INTO {s}.muse_role_session_authorization " + "(authorization_id,step_id,stage,initial_request_hash,frozen_input_hash," + "max_model_calls,max_tool_calls,policy_hash) VALUES (%s,%s,%s,%s,%s,%s,%s,%s) " + "ON CONFLICT DO NOTHING", + ( + 授权, + 步骤ID, + 阶段, + 初始请求哈希, + 固定哈希, + 最大模型调用, + 最大工具调用, + 内容哈希(策略), + ), + ) + 行 = self._查询( + 连, + "SELECT policy_hash FROM {s}.muse_role_session_authorization " + "WHERE authorization_id=%s", + (授权,), + ).fetchone() + if 行 is None or 行["policy_hash"] != 内容哈希(策略): + raise 原文错误("会话批准命令不能改变已批准次数、阶段或范围") + return 授权 + + def 核对会话(self, 领取: 领取凭证, 会话ID: str, 初始请求哈希: str, 阶段: str) -> dict: + with self._事务() as 连: + 任务, _ = 校验领取(任务存储(连, self.数据库.用途), 领取) + 授权 = self._授权(连, 会话ID, 领取.任务ID, "persistent") + 会话 = self._查询( + 连, + "SELECT * FROM {s}.muse_role_session_authorization WHERE authorization_id=%s", + (会话ID,), + ).fetchone() + if 会话 is None or ( + 会话["step_id"], + 会话["stage"], + 会话["initial_request_hash"], + 会话["frozen_input_hash"], + ) != (领取.步骤ID, 阶段, 初始请求哈希, 内容哈希(任务["frozen_input"])): + raise 原文错误("角色会话与批准的输入、阶段或冻结范围不一致") + return {**会话, "valid_until": 授权["valid_until"]} + + def 派生会话保留( + self, + 领取: 领取凭证, + 会话ID: str, + 引用: str, + 哈希: str, + *, + 初始请求哈希: str, + 阶段: str, + 种类: str, + 前序证据: tuple[str, ...], + ) -> str: + """S02内部调用:实际字节只由已固定初始输入及本会话模型/工具回执组装。""" + self.核对会话(领取, 会话ID, 初始请求哈希, 阶段) + if 种类 not in {"model", "tool"} or not re.fullmatch(r"[0-9a-f]{64}", 哈希): + raise 原文错误("派生种类或内容哈希无效") + with self._事务() as 连: + 校验领取(任务存储(连, self.数据库.用途), 领取) + # 父行锁固定整个会话的次数,原批准身份和时间由继承取得。 + 父 = self._查询( + 连, + "SELECT * FROM {s}.muse_raw_authorization WHERE authorization_id=%s FOR UPDATE", + (会话ID,), + ).fetchone() + self._授权(连, 会话ID, 领取.任务ID, "persistent") + 会话 = self._查询( + 连, + "SELECT * FROM {s}.muse_role_session_authorization WHERE authorization_id=%s", + (会话ID,), + ).fetchone() + assert 父 is not None and 会话 is not None # 已核对的父行在本事务中持锁。 + 旧 = self._查询( + 连, + "SELECT * FROM {s}.muse_raw_authorization " + "WHERE parent_authorization_id=%s AND command_id=%s", + (会话ID, 引用), + ).fetchone() + if 旧: + if ( + 旧["call_request_hash"], + 旧["derivation_kind"], + str(旧["attempt_id"]), + tuple(map(str, 旧["evidence_refs"])), + ) != (哈希, 种类, 领取.尝试ID, 前序证据): + raise 原文错误("派生引用已对应其他回合或内容") + return str(旧["authorization_id"]) + 计数行 = self._查询( + 连, + "SELECT count(*) AS n FROM {s}.muse_raw_authorization " + "WHERE parent_authorization_id=%s AND derivation_kind=%s", + (会话ID, 种类), + ).fetchone() + assert 计数行 is not None # 聚合查询恒返回一行。 + 次数 = 计数行["n"] + if 次数 >= 会话["max_model_calls" if 种类 == "model" else "max_tool_calls"]: + raise 原文错误("角色会话已达到批准的调用次数") + if 种类 == "model" and 次数 == 0: + if 哈希 != 会话["initial_request_hash"] or 前序证据: + raise 原文错误("首回合必须使用明确批准的初始请求") + elif not 前序证据: + raise 原文错误("后续回合必须关联本会话已保存的前序证据") + for 证据ID in 前序证据: + 行 = self._查询( + 连, + "SELECT e.evidence_id FROM {s}.muse_runtime_evidence e " + "JOIN {s}.muse_raw_authorization a ON a.authorization_id=e.authorization_id " + "WHERE e.evidence_id=%s AND e.task_id=%s " + "AND e.content IS NOT NULL AND e.outcome='completed' " + "AND a.parent_authorization_id=%s", + (证据ID, 领取.任务ID, 会话ID), + ).fetchone() + if 行 is None: + raise 原文错误("派生依据不是本任务已经保存的会话证据") + 身份 = str(uuid.uuid4()) + self._查询( + 连, + "INSERT INTO {s}.muse_raw_authorization " + "(authorization_id,task_id,command_id,approved_by,source_version,content_hashes," + "purpose,retention_mode,approved_at,valid_until,request_hash,call_id,call_request_hash," + "parent_authorization_id,attempt_id,derivation_kind,evidence_refs) " + "VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)", + ( + 身份, + 领取.任务ID, + 引用, + 父["approved_by"], + 父["source_version"], + [哈希], + 父["purpose"], + 父["retention_mode"], + 父["approved_at"], + 父["valid_until"], + 哈希, + 引用, + 哈希, + 会话ID, + 领取.尝试ID, + 种类, + list(前序证据), + ), + ) + return 身份 + + def 读取会话证据( + self, 领取: 领取凭证, 会话ID: str, 初始请求哈希: str, 阶段: str + ) -> dict[str, dict]: + """历史证据保留旧尝试归属;新尝试只复用已结算完整结果。""" + self.核对会话(领取, 会话ID, 初始请求哈希, 阶段) + with self._事务() as 连: + 行集 = self._查询( + 连, + "SELECT a.authorization_id,a.revoked,a.call_id,a.call_request_hash," + "a.derivation_kind,a.response_hash," + "e.evidence_id,e.content_hash,e.content,e.outcome,e.revision," + "b.state AS budget_state,b.over_budget FROM {s}.muse_raw_authorization a " + "LEFT JOIN {s}.muse_runtime_evidence e ON e.authorization_id=a.authorization_id " + "AND e.kind=CASE a.derivation_kind WHEN 'model' THEN 'model_response' " + "ELSE 'tool_result' END " + "LEFT JOIN {s}.muse_budget_reservation b ON b.task_id=a.task_id " + "AND b.call_id=a.call_id " + "WHERE a.parent_authorization_id=%s AND a.task_id=%s", + (会话ID, 领取.任务ID), + ).fetchall() + 结果 = {} + for 行 in 行集: + if 行["revoked"]: + raise 原文错误("会话历史回合的保留批准已撤销") + if 行["content"] is None or 行["outcome"] != "completed": + raise 原文错误("会话已有回合缺少可恢复的完整成功证据") + if 行["derivation_kind"] == "model" and ( + 行["budget_state"] != "settled" or 行["over_budget"] + ): + raise 原文错误("会话历史调用仍需费用对账或预算处理") + if hashlib.sha256(行["content"]).hexdigest() != 行["content_hash"]: + raise 原文错误("会话历史字节与固定哈希不符") + 结果[行["call_id"]] = { + **行, + "evidence_id": str(行["evidence_id"]), + "authorization_id": str(行["authorization_id"]), + } + return 结果 + + def 绑定调用响应(self, 任务ID: str, 授权ID: str, 调用ID: str, 请求哈希: str, 哈希: str) -> None: + """由模型运行 owner 计算响应哈希;不作为作者或工具的任意扩权入口。""" + if not re.fullmatch(r"[0-9a-f]{64}", 哈希): + raise 原文错误("响应哈希无效") + with self._事务() as 连: + 任务存储(连, self.数据库.用途).任务行(任务ID, 锁定=True) + # 派生回合同时受父批准撤销和截止约束。 + 授权方式 = self._查询( + 连, + "SELECT retention_mode FROM {s}.muse_raw_authorization " + "WHERE authorization_id=%s AND task_id=%s", + (授权ID, 任务ID), + ).fetchone() + if 授权方式 is None: + raise 原文错误("本任务没有对应调用授权") + self._授权(连, 授权ID, 任务ID, 授权方式["retention_mode"]) + 行 = self._查询( + 连, + "SELECT * FROM {s}.muse_raw_authorization " + "WHERE task_id=%s AND authorization_id=%s FOR UPDATE", + (任务ID, 授权ID), + ).fetchone() + if ( + 行 is None + or 行["revoked"] + or 行["valid_until"] <= self.时钟.现在() + or (行["call_id"], 行["call_request_hash"]) != (调用ID, 请求哈希) + or 行["response_hash"] not in {None, 哈希} + ): + raise 原文错误("响应与原调用授权或既有响应不一致") + 调用 = self._查询( + 连, + "SELECT call_id FROM {s}.muse_budget_reservation " + "WHERE task_id=%s AND call_id=%s AND sent_at IS NOT NULL", + (任务ID, 调用ID), + ).fetchone() + if 调用 is None: + raise 原文错误("尚未发出的调用不能登记响应") + self._查询( + 连, + "UPDATE {s}.muse_raw_authorization SET response_hash=%s WHERE authorization_id=%s", + (哈希, 授权ID), + ) + + def 撤销批准(self, 任务ID: str, 作者: str, 授权ID: str) -> None: + with self._事务() as 连: + 任务 = 任务存储(连, self.数据库.用途).任务行(任务ID, 锁定=True) + if 任务["author_id"] != 作者: + raise 原文错误("只能由本任务作者撤销原文批准") + 行 = self._查询( + 连, + "UPDATE {s}.muse_raw_authorization SET revoked=true " + "WHERE task_id=%s AND authorization_id=%s RETURNING authorization_id", + (任务ID, 授权ID), + ).fetchone() + if 行 is None: + raise 原文错误("本任务没有对应原文批准") + + def 创建租约( + self, 任务ID: str, 授权ID: str, 命令ID: str, 截止: datetime, *, 最少剩余秒: float + ) -> str: + 当前 = self.时钟.现在() + 校验期限(当前, 截止, 最少剩余秒) + with self._事务() as 连: + 授权 = self._授权(连, 授权ID, 任务ID, "temporary") + if 截止 > 授权["valid_until"]: + raise 原文错误("租约超过作者批准的期限") + 身份 = str(uuid.uuid4()) + self._查询( + 连, + ( + "INSERT INTO {s}.muse_raw_lease " + "(lease_id,task_id,authorization_id,command_id,created_at,retain_until,state) " + "VALUES (%s,%s,%s,%s,%s,%s,'open') ON CONFLICT DO NOTHING" + ), + (身份, 任务ID, 授权ID, 命令ID, 当前, 截止), + ) + 行 = self._查询( + 连, + "SELECT * FROM {s}.muse_raw_lease WHERE task_id=%s AND command_id=%s", + (任务ID, 命令ID), + ).fetchone() + if 行 is None or str(行["authorization_id"]) != 授权ID or 行["retain_until"] != 截止: + raise 原文错误("租约命令对应不同授权或期限") + # 离开事务后才返回身份;此时才允许调用文件写入。 + return str(行["lease_id"]) + + def _租约(self, 连, 租约ID: str, 任务ID: str): + 行 = self._查询( + 连, + "SELECT * FROM {s}.muse_raw_lease WHERE lease_id=%s AND task_id=%s FOR UPDATE", + (租约ID, 任务ID), + ).fetchone() + if 行 is None: + raise 原文错误("本任务没有对应原文租约") + return 行 + + def _可消费(self, 连, 行): + if 行["state"] != "open" or 行["retain_until"] <= self.时钟.现在(): + raise 原文错误("原文租约当前不可消费") + return self._授权(连, str(行["authorization_id"]), str(行["task_id"]), "temporary") + + def 写入(self, 任务ID: str, 租约ID: str, 哈希: str, 数据: bytes) -> None: + 校验原文字节(数据) + with self._事务() as 连: + 行 = self._租约(连, 租约ID, 任务ID) + 授权 = self._可消费(连, 行) + if 哈希 not in 授权["content_hashes"]: + raise 原文错误("原文字节不在批准范围") + self._绑定文件(连).写入(租约ID, 哈希, 数据) + + def 检查执行余量( + self, 任务ID: str, 租约ID: str, 所需秒: float, *, 授权ID: str | None = None + ) -> None: + """在预算预留和模型启动之前调用;创建租约时的检查不能代替本次检查。""" + with self._事务() as 连: + 行 = self._租约(连, 租约ID, 任务ID) + if 授权ID is not None and str(行["authorization_id"]) != 授权ID: + raise 原文错误("临时租约与本次调用批准不一致") + self._可消费(连, 行) + 校验期限(self.时钟.现在(), 行["retain_until"], 所需秒) + + def 读取(self, 任务ID: str, 租约ID: str, 哈希: str) -> bytes: + with self._事务() as 连: + 行 = self._租约(连, 租约ID, 任务ID) + 授权 = self._可消费(连, 行) + if 哈希 not in 授权["content_hashes"]: + raise 原文错误("请求原文不在批准范围") + return self._绑定文件(连).读取(租约ID, 哈希) + + def 状态(self, 任务ID: str, 租约ID: str) -> dict: + with self._事务() as 连: + 行 = self._租约(连, 租约ID, 任务ID) + return { + "lease_id": 租约ID, + "state": 行["state"], + "archive_id": str(行["archive_id"]) if 行["archive_id"] else None, + "cleanup_complete": 行["cleanup_receipt"] is not None, + } + + def 读取归档(self, 任务ID: str, 归档ID: str, 哈希: str) -> bytes: + """归档保留与消费权限分开;普通读取仍需本任务当前有效授权。""" + with self._事务() as 连: + 行 = self._查询( + 连, + "SELECT a.authorization_id FROM {s}.muse_raw_archive a " + "JOIN {s}.muse_raw_authorization g USING(authorization_id) " + "WHERE a.archive_id=%s AND g.task_id=%s", + (归档ID, 任务ID), + ).fetchone() + if 行 is None: + raise 原文错误("本任务没有对应的完整归档") + 头 = self._查询( + 连, + "SELECT authorization_id FROM {s}.muse_raw_archive WHERE archive_id=%s", + (归档ID,), + ).fetchone() + if 头 is None: + raise 原文错误("归档清单缺失") + 授权 = self._授权(连, str(头["authorization_id"]), 任务ID, "archive") + if 哈希 not in 授权["content_hashes"]: + raise 原文错误("归档读取超出批准哈希范围") + self._核对归档(连, 归档ID) + 字节 = self._查询( + 连, + "SELECT content FROM {s}.muse_raw_archive_item " + "WHERE archive_id=%s AND content_hash=%s", + (归档ID, 哈希), + ).fetchone() + if 字节 is None: + raise 原文错误("归档中没有指定原文") + return 字节["content"] + + def 导入历史归档(self, 任务ID: str, 授权ID: str, 旧引用: str, 目录: Path) -> dict: + """迁移接入必须提供已批准的旧来源身份;旧字节和原回执保留不动。""" + if not 旧引用: + raise 原文错误("历史归档导入需要稳定来源引用") + 旧回执, 旧清单, 内容 = 读取历史归档(目录) + for 数据 in 内容.values(): + 校验原文字节(数据) + 旧元信息 = {"receipt": 旧回执, "entries": 旧清单} + with self._事务() as 连: + 授权 = self._授权(连, 授权ID, 任务ID, "archive") + if 授权["source_version"] != 旧引用 or set(授权["content_hashes"]) != set(内容): + raise 原文错误("历史归档来源或完整哈希集合不在批准范围") + 身份 = str(uuid.uuid4()) + 清单 = [{"hash": h, "bytes": len(数据)} for h, 数据 in sorted(内容.items())] + from psycopg.types.json import Jsonb + + 插入 = self._查询( + 连, + "INSERT INTO {s}.muse_raw_archive " + "(archive_id,legacy_ref,legacy_manifest,authorization_id,tree_hash,entry_count," + "total_bytes,receipt_id,created_at) VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s) " + "ON CONFLICT (legacy_ref) DO NOTHING RETURNING archive_id", + ( + 身份, + 旧引用, + Jsonb(旧元信息), + 授权ID, + 内容哈希(清单), + len(内容), + sum(map(len, 内容.values())), + str(uuid.uuid4()), + self.时钟.现在(), + ), + ).fetchone() + 行 = self._查询( + 连, + "SELECT * FROM {s}.muse_raw_archive WHERE legacy_ref=%s FOR UPDATE", + (旧引用,), + ).fetchone() + if ( + 行 is None + or str(行["authorization_id"]) != 授权ID + or 行["legacy_manifest"] != 旧元信息 + or 行["tree_hash"] != 内容哈希(清单) + ): + raise 原文错误("同一旧归档身份已经绑定其他内容或授权") + if 插入: + for 哈希, 数据 in 内容.items(): + self._查询( + 连, + "INSERT INTO {s}.muse_raw_archive_item " + "(archive_id,content_hash,content) VALUES (%s,%s,%s)", + (身份, 哈希, 数据), + ) + return self._核对归档(连, str(行["archive_id"])) + + def 清理无租约原文(self) -> dict: + """先确定文件候选,再在新事务中核对租约;清理意图可靠保存后才删除。""" + with self._事务() as 连: + 文件 = self._绑定文件(连) + 候选, 未知 = 文件.列出托管对象() + with self._事务() as 连: + 已登记 = { + str(r["lease_id"]) + for r in self._查询( + 连, + "SELECT l.lease_id FROM {s}.muse_raw_lease l JOIN {s}.muse_task t " + "ON t.task_id=l.task_id WHERE t.run_purpose=%s", + (self.数据库.用途.value,), + ).fetchall() + } + for 身份 in set(候选) - 已登记: + self._查询( + 连, + "INSERT INTO {s}.muse_raw_orphan_cleanup " + "(run_purpose,object_id,state) VALUES (%s,%s,'pending') " + "ON CONFLICT (run_purpose,object_id) DO UPDATE SET state='pending'", + (self.数据库.用途.value, 身份), + ) + 待清理 = self._查询( + 连, + "SELECT object_id,receipt_id FROM {s}.muse_raw_orphan_cleanup " + "WHERE run_purpose=%s AND state='pending'", + (self.数据库.用途.value,), + ).fetchall() + 结果 = { + "cleaned_orphan_ids": [], + "receipts": [], + "needs_recovery": [], + "unknown_entries": 未知, + } + for 行 in 待清理: + 身份 = str(行["object_id"]) + if 身份 in 已登记: + 结果["needs_recovery"].append(身份) + continue + try: + 删除 = 文件.删除对象(身份) + with self._事务() as 连: + self._查询( + 连, + "UPDATE {s}.muse_raw_orphan_cleanup SET state='completed',removed=%s," + "completed_at=clock_timestamp() WHERE run_purpose=%s AND object_id=%s", + (删除, self.数据库.用途.value, 身份), + ) + 结果["cleaned_orphan_ids"].append(身份) + 结果["receipts"].append(str(行["receipt_id"])) + except Muse错误: + 结果["needs_recovery"].append(身份) + return 结果 + + def 读取孤儿清理回执(self, 对象ID: str) -> dict: + with self._事务() as 连: + 行 = self._查询( + 连, + "SELECT object_id,receipt_id,state,removed,completed_at " + "FROM {s}.muse_raw_orphan_cleanup WHERE run_purpose=%s AND object_id=%s", + (self.数据库.用途.value, 对象ID), + ).fetchone() + if 行 is None: + raise 原文错误("当前用途没有此对象的清理回执") + return { + **行, + "object_id": str(行["object_id"]), + "receipt_id": str(行["receipt_id"]), + "completed_at": 行["completed_at"].isoformat() if 行["completed_at"] else None, + } + + def 恢复任务原文(self, 任务ID: str, 作者: str) -> list[dict]: + """先读PG再按状态恢复;仍有效的open租约保持可用,不顺手删除。""" + with self._事务() as 连: + 任务 = 任务存储(连, self.数据库.用途).任务行(任务ID) + if 任务["author_id"] != 作者: + raise 原文错误("任务原文恢复需要对应作者身份") + # failed仍允许作者恢复;不能因一次失败删除其尚有效的原文。 + 已终止 = 任务["state"] in {"cancelled", "completed"} + 租约 = self._查询( + 连, + "SELECT * FROM {s}.muse_raw_lease WHERE task_id=%s ORDER BY created_at,lease_id", + (任务ID,), + ).fetchall() + 结果 = [] + for 行 in 租约: + 身份 = str(行["lease_id"]) + try: + if 行["state"] == "migrating": + 回执 = self.归档(任务ID, 身份, str(行["archive_authorization_id"])) + 结果.append({"lease_id": 身份, "action": "archived", "receipt": 回执}) + elif 行["state"] == "open" and 行["retain_until"] > self.时钟.现在() and not 已终止: + 结果.append({"lease_id": 身份, "action": "retained"}) + else: + 已删除 = self.清理(任务ID, 身份, 作者=作者) + 结果.append({"lease_id": 身份, "action": "cleaned" if 已删除 else "missing"}) + except Muse错误 as exc: + 结果.append({"lease_id": 身份, "action": "needs_recovery", "code": exc.错误码}) + return 结果 + + def 清理(self, 任务ID: str, 租约ID: str, *, 作者: str | None = None) -> bool: + with self._事务() as 连: + 行 = self._租约(连, 租约ID, 任务ID) + if 行["state"] == "migrating": + raise 原文错误("迁移中的原文只能先对账恢复") + if 行["state"] == "open" and 行["retain_until"] > self.时钟.现在(): + if 任务存储(连, self.数据库.用途).任务行(任务ID)["author_id"] != 作者: + raise 原文错误("未到期原文清理需要任务作者动作") + if 行["state"] == "migrated": + self._核对归档(连, str(行["archive_id"])) + 已删除 = self._绑定文件(连).删除对象(租约ID) + self._查询( + 连, + ( + "UPDATE {s}.muse_raw_lease SET state=CASE " + "WHEN state='migrated' THEN state ELSE " + "'closed' END,cleanup_receipt=COALESCE(cleanup_receipt,%s)," + "failure_code=CASE WHEN state='open' AND NOT %s THEN 'RAW_CONTENT_MISSING' " + "ELSE failure_code END WHERE lease_id=%s" + ), + (str(uuid.uuid4()), 已删除, 租约ID), + ) + return 已删除 + + def 归档(self, 任务ID: str, 租约ID: str, 归档授权ID: str) -> dict: + with self._事务() as 连: + 行 = self._租约(连, 租约ID, 任务ID) + if 行["state"] == "closed": + raise 原文错误("已关闭租约不能补造归档") + 临时授权 = self._授权(连, str(行["authorization_id"]), 任务ID, "temporary", 恢复=True) + 授权 = self._授权( + 连, 归档授权ID, 任务ID, "archive", 恢复=行["state"] in {"migrating", "migrated"} + ) + if ( + set(临时授权["content_hashes"]) != set(授权["content_hashes"]) + or 临时授权["source_version"] != 授权["source_version"] + ): + raise 原文错误("归档授权没有覆盖完整来源及哈希") + if 行["state"] == "open": + self._可消费(连, 行) + self._查询( + 连, + "UPDATE {s}.muse_raw_lease SET state='migrating',archive_id=%s," + "archive_authorization_id=%s WHERE lease_id=%s", + (str(uuid.uuid4()), 归档授权ID, 租约ID), + ) + elif str(行["archive_authorization_id"]) != 归档授权ID: + raise 原文错误("归档恢复不能改变原授权") + with self._事务() as 连: + 行 = self._租约(连, 租约ID, 任务ID) + archive_id = str(行["archive_id"]) + self._授权(连, 归档授权ID, 任务ID, "archive", 恢复=True) + if 行["state"] == "migrating": + 文件 = self._绑定文件(连) + 内容 = {h: 文件.读取(租约ID, h) for h in sorted(授权["content_hashes"])} + 清单 = [{"hash": h, "bytes": len(v)} for h, v in 内容.items()] + self._查询( + 连, + ( + "INSERT INTO {s}.muse_raw_archive " + "(archive_id,lease_id,authorization_id,tree_hash,entry_count," + "total_bytes,receipt_id,created_at) " + "VALUES (%s,%s,%s,%s,%s,%s,%s,%s)" + ), + ( + archive_id, + 租约ID, + 归档授权ID, + 内容哈希(清单), + len(内容), + sum(map(len, 内容.values())), + str(uuid.uuid4()), + self.时钟.现在(), + ), + ) + for h, v in 内容.items(): + self._查询( + 连, + "INSERT INTO {s}.muse_raw_archive_item (archive_id,content_hash,content) " + "VALUES (%s,%s,%s)", + (archive_id, h, v), + ) + self._查询( + 连, + "UPDATE {s}.muse_raw_lease SET state='migrated' WHERE lease_id=%s", + (租约ID,), + ) + 回执 = self._核对归档(连, archive_id) + self.清理(任务ID, 租约ID) + return 回执 + + def _核对归档(self, 连, 归档ID: str) -> dict: + 头 = self._查询( + 连, "SELECT * FROM {s}.muse_raw_archive WHERE archive_id=%s", (归档ID,) + ).fetchone() + 行 = self._查询( + 连, + "SELECT content_hash,content FROM {s}.muse_raw_archive_item " + "WHERE archive_id=%s ORDER BY content_hash", + (归档ID,), + ).fetchall() + 清单 = [{"hash": r["content_hash"], "bytes": len(r["content"])} for r in 行] + if ( + 头 is None + or len(行) != 头["entry_count"] + or sum(r["bytes"] for r in 清单) != 头["total_bytes"] + or 内容哈希(清单) != 头["tree_hash"] + or any(hashlib.sha256(r["content"]).hexdigest() != r["content_hash"] for r in 行) + ): + raise 原文错误("归档字节和完整清单不一致") + return { + "archive_id": 归档ID, + "receipt_id": str(头["receipt_id"]), + "entry_count": 头["entry_count"], + "total_bytes": 头["total_bytes"], + "tree_hash": 头["tree_hash"], + } diff --git a/src/muse/任务运行/运行证据.py b/src/muse/任务运行/运行证据.py new file mode 100644 index 0000000..d0a745d --- /dev/null +++ b/src/muse/任务运行/运行证据.py @@ -0,0 +1,193 @@ +"""证据关联与同哈希补交;仅哈希和完整原文分别记录,不代替模型成功裁决。""" + +import hashlib +import re +import uuid +from contextlib import contextmanager +from typing import LiteralString + +import psycopg +from psycopg import sql +from psycopg.rows import dict_row +from psycopg.types.json import Jsonb + +from muse.任务运行.原文生命周期 import 原文错误, 校验原文字节 +from muse.任务运行.存储 import 任务存储 +from muse.共享.调用身份 import 用途 +from muse.基础设施.数据库.连接 import 数据库工厂 + +_元信息 = { + "call_id", + "model", + "provider", + "input_tokens", + "output_tokens", + "policy_version", + "pricing_version", + "runtime_config_hash", + "raw_lease_id", + "resource_release_id", + "error_code", + "source_refs", +} + + +class 证据服务: + def __init__(self, 数据库: 数据库工厂): + self.数据库 = 数据库 + self.schema = "evaluation" if 数据库.用途 is 用途.评测 else "public" + + def _查询(self, 连, 语句: LiteralString, 参数=()): + return 连.cursor(row_factory=dict_row).execute( + sql.SQL(语句).format(s=sql.Identifier(self.schema)), 参数 + ) + + @contextmanager + def _事务(self): + try: + with self.数据库.连接() as 连, 连.transaction(): + yield 连 + except psycopg.Error as exc: + raise 原文错误("证据没有确认保存成功", 上下文={"sqlstate": exc.sqlstate}) from None + + def 保存( + self, + 任务ID: str, + 尝试ID: str, + 种类: str, + 哈希: str, + 结果状态: str, + 元信息: dict, + *, + 数据: bytes | None = None, + 授权ID: str | None = None, + ) -> dict: + if 种类 not in { + "model_input", + "model_response", + "tool_result", + "failure", + } or 结果状态 not in {"completed", "failed", "partial"}: + raise 原文错误("证据种类或结果状态无效") + if not re.fullmatch(r"[0-9a-f]{64}", 哈希) or set(元信息) - _元信息: + raise 原文错误("证据只接受确切哈希与登记的技术元信息") + for 名称, 值 in 元信息.items(): + if 名称 in {"input_tokens", "output_tokens"}: + 合法 = 值 is None or (type(值) is int and 值 >= 0) + elif 名称 == "source_refs": + 合法 = isinstance(值, list) and all( + isinstance(项, dict) + and set(项) == {"source_id", "version", "content_hash"} + and all(isinstance(v, str) and v for v in 项.values()) + for 项 in 值 + ) + else: + 合法 = 值 is None or isinstance(值, str) + if not 合法: + raise 原文错误("证据元信息不符合登记的字段结构") + if 数据 is not None and hashlib.sha256(数据).hexdigest() != 哈希: + raise 原文错误("补交字节与原证据哈希不一致") + if 数据 is not None: + 校验原文字节(数据) + # 调用身份区分事件,内容哈希校验字节;相同输出不能合并不同调用。 + 引用 = f"call:{元信息['call_id']}" if 元信息.get("call_id") else f"content:{哈希}" + with self._事务() as 连: + 任务存储(连, self.数据库.用途).任务行(任务ID, 锁定=True) + 尝试 = self._查询( + 连, + "SELECT attempt_id FROM {s}.muse_attempt WHERE task_id=%s AND attempt_id=%s", + (任务ID, 尝试ID), + ).fetchone() + if 尝试 is None: + raise 原文错误("证据不能伪造任务或尝试关联") + if 数据 is not None: + 授权 = self._查询( + 连, + "SELECT a.content_hashes,a.response_hash FROM {s}.muse_raw_authorization a " + "WHERE a.task_id=%s AND a.authorization_id=%s " + "AND a.retention_mode='persistent' " + "AND NOT a.revoked AND a.valid_until>clock_timestamp() " + "AND (a.parent_authorization_id IS NULL OR EXISTS " + "(SELECT 1 FROM {s}.muse_raw_authorization p " + "WHERE p.authorization_id=a.parent_authorization_id " + "AND NOT p.revoked AND p.valid_until>clock_timestamp())) FOR SHARE", + (任务ID, 授权ID), + ).fetchone() + if 授权 is None or 哈希 not in {*授权["content_hashes"], 授权["response_hash"]}: + raise 原文错误("完整运行证据缺少对应的持久保留授权") + 身份 = str(uuid.uuid4()) + self._查询( + 连, + "INSERT INTO {s}.muse_runtime_evidence " + "(evidence_id,task_id,attempt_id,kind,reference_id,content_hash,content," + "authorization_id,outcome,metadata) " + "VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s) ON CONFLICT DO NOTHING", + (身份, 任务ID, 尝试ID, 种类, 引用, 哈希, 数据, 授权ID, 结果状态, Jsonb(元信息)), + ) + 行 = self._查询( + 连, + "SELECT * FROM {s}.muse_runtime_evidence WHERE attempt_id=%s AND kind=%s " + "AND reference_id=%s FOR UPDATE", + (尝试ID, 种类, 引用), + ).fetchone() + if ( + 行 is None + or 行["content_hash"] != 哈希 + or 行["outcome"] != 结果状态 + or 行["metadata"] != 元信息 + ): + raise 原文错误("同一证据身份不能覆盖原结果或元信息") + if 行["content"] is None and 数据 is not None: + self._查询( + 连, + "UPDATE {s}.muse_runtime_evidence SET content=%s,authorization_id=%s," + "revision=revision+1 WHERE evidence_id=%s", + (数据, 授权ID, 行["evidence_id"]), + ) + 行 = {**行, "content": 数据, "revision": 行["revision"] + 1} + return self._回执(行) + + def 读取回执(self, 任务ID: str, 证据ID: str) -> dict: + with self._事务() as 连: + 行 = self._查询( + 连, + "SELECT * FROM {s}.muse_runtime_evidence WHERE task_id=%s AND evidence_id=%s", + (任务ID, 证据ID), + ).fetchone() + if 行 is None: + raise 原文错误("本任务没有对应证据") + return self._回执(行) + + def 补交原文(self, 任务ID: str, 证据ID: str, 数据: bytes, 授权ID: str) -> dict: + """只承接现有证据的确切字节;调用者不能改结果、尝试或技术元信息。""" + with self._事务() as 连: + 行 = self._查询( + 连, + "SELECT * FROM {s}.muse_runtime_evidence WHERE task_id=%s AND evidence_id=%s", + (任务ID, 证据ID), + ).fetchone() + if 行 is None: + raise 原文错误("本任务没有对应证据") + return self.保存( + 任务ID, + str(行["attempt_id"]), + 行["kind"], + 行["content_hash"], + 行["outcome"], + 行["metadata"], + 数据=数据, + 授权ID=授权ID, + ) + + @staticmethod + def _回执(行) -> dict: + 回执 = { + "evidence_id": str(行["evidence_id"]), + "content_hash": 行["content_hash"], + "outcome": 行["outcome"], + "retention": "full" if 行["content"] is not None else "hash_only", + "revision": 行["revision"], + } + if 行["metadata"].get("raw_lease_id"): + 回执["raw_lease_id"] = 行["metadata"]["raw_lease_id"] + return 回执 diff --git a/src/muse/基础设施/受控文件.py b/src/muse/基础设施/受控文件.py new file mode 100644 index 0000000..b326122 --- /dev/null +++ b/src/muse/基础设施/受控文件.py @@ -0,0 +1,248 @@ +"""只通过不透明租约ID和内容哈希访问私人暂存;状态许可由S02执行。""" + +import hashlib +import json +import os +import re +import stat +import uuid +from contextlib import contextmanager +from pathlib import Path + +from muse.共享.错误 import Muse错误 + + +class 受控文件错误(Muse错误): + 错误码 = "CONTROLLED_FILE_INVALID" + + +class 受控文件: + def __init__(self, 根: Path): + 根.mkdir(mode=0o700, parents=True, exist_ok=True) + self.根 = 根 + self._命名空间 = None + with self._根句柄(): + pass + + def 绑定命名空间(self, 标识: str, 运行用途: str) -> None: + """根只认当前库和用途;未标记的非空目录不能被自动接管。""" + self._校验对象ID(标识) + 内容 = {"format": 1, "namespace_id": 标识, "run_purpose": 运行用途} + 名称 = ".muse-namespace.json" + with self._根句柄() as fd: + try: + 已有 = os.open(名称, os.O_RDONLY | os.O_NOFOLLOW, dir_fd=fd) + except FileNotFoundError: + if os.listdir(fd): + raise 受控文件错误("非空暂存根缺少命名空间标记,须先核对并迁移") from None + 暂名 = ".namespace-" + uuid.uuid4().hex + out = os.open(暂名, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600, dir_fd=fd) + try: + with os.fdopen(out, "wb") as 流: + 流.write(json.dumps(内容, sort_keys=True).encode()) + 流.flush() + os.fsync(流.fileno()) + try: + os.link(暂名, 名称, src_dir_fd=fd, dst_dir_fd=fd, follow_symlinks=False) + except FileExistsError: + pass + os.fsync(fd) + finally: + os.unlink(暂名, dir_fd=fd) + 已有 = os.open(名称, os.O_RDONLY | os.O_NOFOLLOW, dir_fd=fd) + try: + with os.fdopen(已有, "rb") as 流: + 信息 = os.fstat(流.fileno()) + if not stat.S_ISREG(信息.st_mode) or stat.S_IMODE(信息.st_mode) != 0o600: + raise 受控文件错误("暂存命名空间标记类型或权限无效") + if json.load(流) != 内容: + raise 受控文件错误("暂存根属于其他数据库或运行用途") + except (ValueError, UnicodeError): + raise 受控文件错误("暂存根命名空间标记损坏,禁止自动清理") from None + self._命名空间 = 内容 + + def 列出托管对象(self) -> tuple[list[str], int]: + if self._命名空间 is None: + raise 受控文件错误("扫描前必须从当前数据库核对暂存根命名空间") + self.绑定命名空间(self._命名空间["namespace_id"], self._命名空间["run_purpose"]) + 对象, 未知 = [], 0 + with self._根句柄() as fd: + for 名 in os.listdir(fd): + if 名 == ".muse-namespace.json": + continue + try: + self._校验对象ID(名) + if not stat.S_ISDIR(os.stat(名, dir_fd=fd, follow_symlinks=False).st_mode): + raise 受控文件错误("非对象目录") + 对象.append(名) + except 受控文件错误: + 未知 += 1 + return sorted(对象), 未知 + + @contextmanager + def _根句柄(self): + try: + fd = os.open(self.根, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW) + try: + if stat.S_IMODE(os.fstat(fd).st_mode) != 0o700: + raise 受控文件错误("受控暂存根权限必须为0700") + yield fd + finally: + os.close(fd) + except OSError: + raise 受控文件错误("受控暂存不可访问") from None + + @staticmethod + def _校验对象ID(对象ID: str) -> None: + try: + if str(uuid.UUID(对象ID)) != 对象ID: + raise ValueError() + except ValueError: + raise 受控文件错误("暂存只接受规范不透明对象身份") from None + + @contextmanager + def _对象目录(self, 对象ID: str, *, 创建: bool = False): + self._校验对象ID(对象ID) + with self._根句柄() as root: + if 创建: + try: + os.mkdir(对象ID, mode=0o700, dir_fd=root) + except FileExistsError: + pass + fd = os.open(对象ID, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW, dir_fd=root) + try: + if stat.S_IMODE(os.fstat(fd).st_mode) != 0o700: + raise 受控文件错误("暂存对象权限必须为0700") + yield fd + finally: + os.close(fd) + + @staticmethod + def _文件名(哈希: str) -> str: + if not re.fullmatch(r"[0-9a-f]{64}", 哈希): + raise 受控文件错误("内容哈希格式无效") + return 哈希 + ".bin" + + def 写入(self, 对象ID: str, 哈希: str, 数据: bytes) -> None: + 名称 = self._文件名(哈希) + if hashlib.sha256(数据).hexdigest() != 哈希: + raise 受控文件错误("输入字节与授权哈希不同") + with self._对象目录(对象ID, 创建=True) as fd: + 暂名 = ".part-" + uuid.uuid4().hex + out = os.open( + 暂名, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600, dir_fd=fd + ) + try: + with os.fdopen(out, "wb") as stream: + stream.write(数据) + stream.flush() + os.fsync(stream.fileno()) + try: + os.link(暂名, 名称, src_dir_fd=fd, dst_dir_fd=fd, follow_symlinks=False) + except FileExistsError: + if self.读取(对象ID, 哈希) != 数据: + raise 受控文件错误("现有暂存内容不一致") from None + os.fsync(fd) + finally: + os.unlink(暂名, dir_fd=fd) + + def 读取(self, 对象ID: str, 哈希: str) -> bytes: + with self._对象目录(对象ID) as fd: + file = os.open(self._文件名(哈希), os.O_RDONLY | os.O_NOFOLLOW, dir_fd=fd) + with os.fdopen(file, "rb") as stream: + info = os.fstat(stream.fileno()) + if not stat.S_ISREG(info.st_mode) or stat.S_IMODE(info.st_mode) != 0o600: + raise 受控文件错误("暂存文件类型或权限无效") + 数据 = stream.read() + if hashlib.sha256(数据).hexdigest() != 哈希: + raise 受控文件错误("暂存字节已发生变化") + return 数据 + + def 删除对象(self, 对象ID: str) -> bool: + self._校验对象ID(对象ID) + with self._根句柄() as root: + try: + os.stat(对象ID, dir_fd=root, follow_symlinks=False) + except FileNotFoundError: + return False + with self._对象目录(对象ID) as fd: + 文件 = os.listdir(fd) + for 名 in 文件: + if not re.fullmatch(r"(?:[0-9a-f]{64}\.bin|\.part-[0-9a-f]{32})", 名): + raise 受控文件错误("对象目录存在未知资料,停止自动清理") + if not stat.S_ISREG(os.stat(名, dir_fd=fd, follow_symlinks=False).st_mode): + raise 受控文件错误("对象目录存在非普通文件,停止自动清理") + for 名 in 文件: + os.unlink(名, dir_fd=fd) + os.fsync(fd) + with self._根句柄() as root: + os.rmdir(对象ID, dir_fd=root) + os.fsync(root) + return True + + +def 读取历史归档(目录: Path) -> tuple[dict, list[dict], dict[str, bytes]]: + """只读旧raw归档及其已封存回执;保留文件名映射,不改旧目录和权限。""" + import json + + 条目: list[dict] = [] + 字节: dict[str, bytes] = {} + + def 读取(fd: int, 名称: str) -> bytes: + file = os.open(名称, os.O_RDONLY | os.O_NOFOLLOW, dir_fd=fd) + with os.fdopen(file, "rb") as 流: + info = os.fstat(流.fileno()) + if not stat.S_ISREG(info.st_mode): + raise 受控文件错误("历史归档含非普通文件") + return 流.read() + + def 扫描(fd: int, 前缀: str = "") -> None: + for 名称 in sorted(os.listdir(fd)): + if not 前缀 and 名称 == ".migration-receipt.json": + continue + info = os.stat(名称, dir_fd=fd, follow_symlinks=False) + if stat.S_ISDIR(info.st_mode): + sub = os.open(名称, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW, dir_fd=fd) + try: + 扫描(sub, 前缀 + 名称 + "/") + finally: + os.close(sub) + else: + if not stat.S_ISREG(info.st_mode): + raise 受控文件错误("历史归档含链接或非普通文件") + 数据 = 读取(fd, 名称) + 哈希 = hashlib.sha256(数据).hexdigest() + 条目.append( + { + "relativePath": 前缀 + 名称, + "contentSha256": "sha256:" + 哈希, + "sizeBytes": len(数据), + } + ) + 字节[哈希] = 数据 + + try: + fd = os.open(目录, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW) + try: + 回执 = json.loads(读取(fd, ".migration-receipt.json")) + 扫描(fd) + finally: + os.close(fd) + except (OSError, ValueError): + raise 受控文件错误("历史归档或迁移回执无法读取") from None + if not isinstance(回执, dict): + raise 受控文件错误("历史归档迁移回执必须是对象") + 清单字节 = json.dumps( + 条目, ensure_ascii=False, sort_keys=True, separators=(",", ":"), allow_nan=False + ).encode() + if ( + 回执.get("schemaVersion") != "raw-vault-migration-receipt-v1" + or 回执.get("status") != "migrated" + or not 回执.get("archiveId") + or 回执.get("archiveSha256") != "sha256:" + hashlib.sha256(清单字节).hexdigest() + or 回执.get("fileCount") != len(条目) + or 回执.get("totalBytes") != sum(v["sizeBytes"] for v in 条目) + or not 条目 + ): + raise 受控文件错误("历史归档的完整清单与回执不一致") + return 回执, 条目, 字节 diff --git a/tests/单元/test_受控文件CAS与恢复.py b/tests/单元/test_受控文件CAS与恢复.py new file mode 100644 index 0000000..04bee52 --- /dev/null +++ b/tests/单元/test_受控文件CAS与恢复.py @@ -0,0 +1,43 @@ +"""合成私人暂存的字节、权限与清理边界。""" + +import hashlib +import os +import uuid + +import pytest + +from muse.基础设施.受控文件 import 受控文件, 受控文件错误 + + +def test_内容哈希排他发布和重复写入保持字节__a80001(tmp_path) -> None: + 管理器 = 受控文件(tmp_path / "私人暂存") + 身份 = str(uuid.uuid4()) + 原文 = '中文"引号"\n换行\\反斜杠'.encode() + 哈希 = hashlib.sha256(原文).hexdigest() + 管理器.写入(身份, 哈希, 原文) + 管理器.写入(身份, 哈希, 原文) + assert 管理器.读取(身份, 哈希) == 原文 + assert (管理器.根 / 身份 / (哈希 + ".bin")).stat().st_mode & 0o777 == 0o600 + with pytest.raises(受控文件错误): + 管理器.写入(身份, 哈希, b"different") + 管理器.删除对象(身份) + assert not (管理器.根 / 身份).exists() + + +def test_软链接和未知资料不会被清理或跟随__a80002(tmp_path) -> None: + 管理器 = 受控文件(tmp_path / "私人暂存") + 身份 = str(uuid.uuid4()) + 原文 = b"synthetic" + 哈希 = hashlib.sha256(原文).hexdigest() + 管理器.写入(身份, 哈希, 原文) + 外部 = tmp_path / "其他资料" + 外部.write_bytes(b"preserve") + (管理器.根 / 身份 / (哈希 + ".bin")).unlink() + os.symlink(外部, 管理器.根 / 身份 / (哈希 + ".bin")) + with pytest.raises(受控文件错误): + 管理器.读取(身份, 哈希) + with pytest.raises(受控文件错误): + 管理器.删除对象(身份) + assert 外部.read_bytes() == b"preserve" + with pytest.raises(受控文件错误): + 管理器.读取("../其他资料", 哈希) diff --git a/tests/单元/test_证据内容校验.py b/tests/单元/test_证据内容校验.py new file mode 100644 index 0000000..08153a3 --- /dev/null +++ b/tests/单元/test_证据内容校验.py @@ -0,0 +1,28 @@ +"""逐例承接旧凭据过滤;字节校验独立于数据库,无原文改写。""" + +import pytest + +from muse.任务运行.接口 import 原文错误, 校验原文字节 + + +def test_普通token篇幅说明允许原样保留__42d154() -> None: + 数据 = "输出不超过 4000 字;token budget 只用于篇幅估算。".encode() + assert 校验原文字节(数据) is None + assert 校验原文字节('{"最大输出token":64}'.encode()) is None + + +@pytest.mark.parametrize( + "数据", + [ + b"api_key=sk-example-secret-value", + b"token: abcdefghijklmnop", + b"Authorization: Bearer abcdefghijklmnop", + b"-----BEGIN PRIVATE KEY-----", + ], + ids=["api-key", "token-assignment", "bearer", "private-key"], +) +def test_凭据赋值拒绝且错误不泄露命中值__44009d(数据: bytes) -> None: + with pytest.raises(原文错误) as 捕获: + 校验原文字节(数据) + assert 数据.decode() not in str(捕获.value) + assert 捕获.value.上下文 == {} diff --git a/tests/集成/test_原文归档与恢复.py b/tests/集成/test_原文归档与恢复.py new file mode 100644 index 0000000..2d634aa --- /dev/null +++ b/tests/集成/test_原文归档与恢复.py @@ -0,0 +1,590 @@ +"""已授权的合成原文;独立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 test_raw_orphan_scope__a82112(原文环境, tmp_path): + """保留LC-d6e3f58513a1;租约替身升级为真实PG,原占位孤儿改用规范UUID。""" + 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 报告["cleaned_orphan_ids"] == [孤儿] and 报告["unknown_entries"] == 1 + assert not (服务.文件.根 / 孤儿).exists() + assert 服务.读取(任务, 租约, 哈希) == 数据 + assert ( + 外根.读取(孤儿, 孤儿哈希) == 原孤儿字节 and 未知.read_text() == "不属于原文对象的合成资料" + ) + 回执 = 服务.读取孤儿清理回执(孤儿) + assert 回执["receipt_id"] in 报告["receipts"] and 回执["state"] == "completed" + assert 服务.清理无租约原文()["cleaned_orphan_ids"] == [] + assert 服务.读取孤儿清理回执(孤儿) == 回执 + with pytest.raises(受控文件错误, match="其他数据库"): + 原文服务(库[用途.生产], 外根).清理无租约原文() + assert 外根.读取(孤儿, 孤儿哈希) == 原孤儿字节 + + +def test_孤儿清理先保存意图且确认失败可幂等恢复__a82111(原文环境): + import uuid + + 服务, 任务, 授权, 数据, 哈希, 时钟, 库 = 原文环境 + 租约 = 服务.创建租约( + 任务, 授权, "kept-for-namespace", 时钟.现在() + timedelta(hours=1), 最少剩余秒=0 + ) + 服务.写入(任务, 租约, 哈希, 数据) + 孤儿 = str(uuid.uuid4()) + 服务.文件.写入(孤儿, 哈希, 数据) + with 库[用途.维护].连接() as 连: + 连.execute("""CREATE FUNCTION reject_cleanup() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN RAISE EXCEPTION 'synthetic cleanup confirmation failure'; END $$; + CREATE TRIGGER reject_cleanup BEFORE UPDATE ON muse_raw_orphan_cleanup + FOR EACH ROW EXECUTE FUNCTION reject_cleanup();""") + 首次 = 服务.清理无租约原文() + assert 首次["needs_recovery"] == [孤儿] and not (服务.文件.根 / 孤儿).exists() + 待续 = 服务.读取孤儿清理回执(孤儿) + assert 待续["state"] == "pending" + with 库[用途.维护].连接() as 连: + 连.execute("DROP TRIGGER reject_cleanup ON muse_raw_orphan_cleanup") + 再次 = 服务.清理无租约原文() + assert 再次["cleaned_orphan_ids"] == [孤儿] + 完成 = 服务.读取孤儿清理回执(孤儿) + assert 完成["receipt_id"] == 待续["receipt_id"] and 完成["state"] == "completed" + assert 完成["removed"] is False # 此次确认已不存在,不补造首次删除成功记录。 + assert 服务.读取(任务, 租约, 哈希) == 数据 + + +def 建证据(环境): + 原文, 任务, _, 数据, 哈希, _, 库 = 环境 + 登记 = 流程登记() + 登记.登记处理器(步骤处理器("raw-test", "1", lambda _: 步骤结果({}), "1", "1")) + 运行 = 任务服务(库[用途.生产], 登记) + 领取 = 运行.领取步骤("evidence-worker", ["raw-test"]) + assert 领取 is not None and 领取.任务ID == 任务 + return 证据服务(库[用途.生产]), 领取 + + +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 服务.读取(任务, 租约, 哈希) == 数据 + + +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(原文错误): + 证据.补交原文(任务, 回执["evidence_id"], 数据, 临时授权) + 授权 = 原文.批准保留( + 任务, + "author", + "persistent", + 来源版本="source-v1", + 哈希=(哈希,), + 内容用途="detection", + 方式="persistent", + 有效期=时钟.现在() + timedelta(days=1), + ) + with pytest.raises(原文错误): + 证据.补交原文(任务, 回执["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] + == 数据 + ) + + +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_, 授权, 数据, 哈希, 时钟, 应用测试库 + + +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 服务.状态(任务, 租约)["state"] == "closed" + 服务.清理(任务, 租约) + + +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 服务.归档(任务, 租约, 归档授权) == 回执 + assert 服务.状态(任务, 租约) == { + "lease_id": 租约, + "state": "migrated", + "archive_id": 回执["archive_id"], + "cleanup_complete": True, + } + assert not (服务.文件.根 / 租约).exists() + assert 服务.读取归档(任务, 回执["archive_id"], 哈希) == 数据 + assert 原文服务(库[用途.生产], 时钟实现=时钟).读取归档(任务, 回执["archive_id"], 哈希) == 数据 + 时钟.推进(timedelta(hours=2)) + 服务.清理(任务, 租约) + assert 服务.读取归档(任务, 回执["archive_id"], 哈希) == 数据 + with 库[用途.生产].连接() as 连: + assert ( + 连.execute( + "SELECT content FROM muse_raw_archive_item WHERE archive_id=%s", + (回执["archive_id"],), + ).fetchone()[0] + == 数据 + ) + + +def test_缺字节归档保留migrating不被过期删除__a82003(原文环境) -> None: + 服务, 任务, 授权, 数据, 哈希, 时钟, _ = 原文环境 + 租约 = 服务.创建租约(任务, 授权, "lease", 时钟.现在() + timedelta(hours=1), 最少剩余秒=60) + 归档授权 = 服务.批准保留( + 任务, + "author", + "archive", + 来源版本="source-v1", + 哈希=(哈希,), + 内容用途="detection", + 方式="archive", + 有效期=时钟.现在() + timedelta(days=1), + ) + with pytest.raises(受控文件错误): + 服务.归档(任务, 租约, 归档授权) + assert 服务.状态(任务, 租约)["state"] == "migrating" + assert 服务.恢复任务原文(任务, "author")[0]["action"] == "needs_recovery" + 时钟.推进(timedelta(hours=2)) + with pytest.raises(原文错误): + 服务.清理(任务, 租约) + assert 服务.状态(任务, 租约)["state"] == "migrating" + + +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 原状态["state"] == "migrating" + 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"] + + +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") + assert 服务.状态(任务, 租约)["state"] == "open" + assert not 服务.状态(任务, 租约)["cleanup_complete"] + finally: + 目录.chmod(0o700) + 服务.清理(任务, 租约, 作者="author") + assert 服务.状态(任务, 租约)["state"] == "closed" + + +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() + + +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", 目录) + + +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 服务.文件.读取(租约, 哈希) == 数据 + 旧状态 = 服务.状态(任务, 租约) + assert 旧状态["state"] == "migrated" and not 旧状态["cleanup_complete"] + 恢复 = 服务.归档(任务, 租约, 归档授权) + assert 恢复["archive_id"] == 旧状态["archive_id"] + assert 服务.状态(任务, 租约)["cleanup_complete"] + with 库[用途.生产].连接() as 连: + assert 连.execute("SELECT count(*) FROM muse_raw_archive").fetchone()[0] == 1 + + +@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)