W08 原文证据、临时租约与归档:原文证据生命周期、临时租约、调用授权与归档恢复。

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

View File

@ -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"],
}

View File

@ -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 回执

View File

@ -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 回执, 条目, 字节

View File

@ -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(受控文件错误):
管理器.读取("../其他资料", 哈希)

View File

@ -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.上下文 == {}

View File

@ -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)