muse-agent-example/tests/集成/test_任务租约与恢复.py
zizi 713cc45c63 重构(数据库): 迁移链压扁为单基线并精简测试至保留集
- 51 个增量迁移文件压扁为 V0001__基线.sql 完整快照(结构+种子+授权),
  等价门禁:旧链全量执行库与基线库 pg_dump 逐字节一致
- 剔除 pg_dump 固化的 public schema 超级用户归属断言(muse_maint 无权执行)
- 测试删减至保留集:备份往返 2 + 预算 2 + 租约 2 + 基线建库 1
- 租约/预算夹具改共享库,消除按例克隆建库
- 修复共享库三类既有污染:账本注入残留(系统管理)、失败触发器残留(评测)、
  建表残留(正式变更事务),发布包迁移文件名硬编码改动态核对
- 全量数据库验收 722 passed / 0 failed / 0 errors(main 基线为 24F+291E)
- 运行手册登记基线模式改表流程与账本校验和同步
2026-09-22 10:18:51 +08:00

185 lines
7.4 KiB
Python

"""任务租约的核心保护:重启不重复执行、未知调用先对账。
按简化决策删除其余变体用例;夹具从按例克隆改为共享库以消除建库开销。
`任务库` 与 `请求` 同时被 test_任务租约与迟到结果.py 复用。
"""
from __future__ import annotations
import dataclasses
import subprocess
import sys
import time
from pathlib import Path
import pytest
from muse.任务运行.接口 import (
任务服务,
任务状态,
任务请求,
作用域,
步骤处理器,
步骤结果,
步骤计划,
状态冲突,
)
from muse.共享.调用身份 import 内容用途, 用途
from muse.基础设施.数据库.连接 import 数据库工厂
from muse.编排.接口 import 流程定义, 流程服务, 流程登记
pytestmark = pytest.mark.数据库
@pytest.fixture
def 任务库(数据库底座, monkeypatch: pytest.MonkeyPatch, tmp_path: Path):
with 数据库底座.共享库(tmp_path / "连接引用") as 工厂:
# 子进程恢复入口继承相同的按例数据库,仍跨真实连接验证提交。
for 声明, 当前 in 工厂.items():
monkeypatch.setenv(f"MUSE_W05_{声明.name}_URL", Path(当前.引用.位置).read_text())
yield 工厂
def 构造服务(工厂: 数据库工厂, 执行=None, 重验=None):
登记 = 流程登记()
登记.登记处理器(
步骤处理器(
"collect",
"1",
执行 or (lambda 上下文: 步骤结果({"value": 上下文.任务.冻结输入["输入"]["value"]})),
"v1",
"v1",
保护职责="输入校验",
)
)
登记.登记处理器(步骤处理器("finish", "1", lambda 上下文: 步骤结果({"done": True}), "v1", "v1"))
登记.登记类型("example", 必需保护=("输入校验",))
def 校验合成授权(快照):
if 快照.冻结输入["冻结上下文"]["authorization"] != "grant-1":
raise 状态冲突("合成来源授权已改变")
登记.登记恢复检查("example", "1", 重验 or 校验合成授权)
运行 = 任务服务(工厂, 登记)
编排 = 流程服务(运行, 登记)
编排.发布(
流程定义(
"example",
"1",
(步骤计划("collect", "collect", "1"), 步骤计划("finish", "finish", "1", ("collect",))),
)
)
return 运行, 编排
def 请求(命令: str = "cmd", 作品: str = "work-A", 用途标记: 用途 = 用途.生产):
return 任务请求(
"example",
命令,
"author",
用途标记,
内容用途.抽取,
{"value": 7},
"policy-1",
"release-1",
{
"source_scope": {"work_id": 作品, "as_of": 3},
"schema_versions": {"entity": "1"},
"authorization": "grant-1",
"budget": {"remaining": 10},
"stop_conditions": ["cancelled"],
},
作品ID=作品,
排他范围=作用域("work", 作品, "extract"),
)
def 完成所有(服务: 任务服务) -> None:
while 领取 := 服务.领取步骤("worker", ["collect", "finish"]):
服务.执行一步(领取)
@pytest.mark.case_id(
"NC-task-durable-restart",
environment="隔离PostgreSQL",
given="持久化任务且首步已保存结果",
when="新 Python 进程从任意目录继续领取",
then=["复用首步结果并完成后续步骤,不重复执行首步"],
contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md",
)
def test_公开编排重启后复用已完成步骤__a51001(任务库, tmp_path: Path) -> None:
调用次数 = []
def 提取(上下文):
调用次数.append(上下文.领取.步骤ID)
return 步骤结果({"saved": 1})
运行, 编排 = 构造服务(任务库[用途.生产], 提取)
身份 = 运行.创建任务(请求(), "example", "1")
assert 运行.创建任务(请求(), "example", "1") == 身份
with pytest.raises(状态冲突):
运行.创建任务(dataclasses.replace(请求(), 输入={"value": 8}), "example", "1")
首步 = 运行.领取步骤("old", ["collect", "finish"])
assert 首步 is not None
运行.执行一步(首步)
子进程代码 = """
import sys
from muse.任务运行.接口 import 任务服务, 步骤处理器, 步骤结果
from muse.编排.接口 import 流程登记
from muse.基础设施.数据库.连接 import 数据库工厂
from muse.配置 import 数据库引用
from muse.共享.调用身份 import 用途
登记 = 流程登记()
def 已完成不能再执行(_):
raise AssertionError("已完成步骤被重复执行")
登记.登记处理器(步骤处理器("collect", "1", 已完成不能再执行, "v1", "v1", 保护职责="输入校验"))
登记.登记处理器(步骤处理器("finish", "1", lambda _: 步骤结果({"done": True}), "v1", "v1"))
登记.登记类型("example", 必需保护=("输入校验",))
运行 = 任务服务(数据库工厂(数据库引用("环境变量", "MUSE_W05_生产_URL"), 用途.生产), 登记)
while 领取 := 运行.领取步骤("new-process", ["collect", "finish"]):
运行.执行一步(领取)
assert 运行.读取任务(sys.argv[1]).状态.value == "completed"
"""
子进程 = subprocess.run(
[sys.executable, "-c", 子进程代码, 身份], cwd=tmp_path, capture_output=True, text=True
)
assert 子进程.returncode == 0, 子进程.stderr
assert 运行.读取任务(身份).状态 is 任务状态.已完成
assert 调用次数 == ["collect"]
assert all(步骤["result"] is not None for 步骤 in 运行.读取任务(身份).步骤)
@pytest.mark.case_id(
"NC-task-unknown-reconciliation",
environment="隔离PostgreSQL",
given="已登记发出的调用结果未知",
when="过期领取、恢复与可靠结果对账",
then=["未知调用先对账,结果保存为检查点;恢复由原处理器消费已存输出并完成步骤,不重复外发"],
contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md",
)
def test_未知调用先对账并复用可靠输出__a51006(任务库) -> None:
运行, 编排 = 构造服务(
任务库[用途.生产],
执行=lambda 上下文: 步骤结果(上下文.领取.检查点["已保存调用"]["输出"]),
)
身份 = 运行.创建任务(请求(), "example", "1")
旧 = 运行.领取步骤("old", ["collect"], 租期秒=0.12)
assert 旧 is not None
运行.登记外部调用(旧, "call-1")
time.sleep(0.16)
assert 运行.领取步骤("new", ["collect", "finish"]) is None
assert 运行.读取任务(身份).状态 is 任务状态.待对账
with pytest.raises(状态冲突, match="对账"):
运行.控制任务(身份, "author", 任务状态.待对账, "恢复", 命令ID="resume")
运行.对账调用(身份, 旧.尝试ID, "author", 已保存输出={"recovered": 1}, 对账回执="receipt-1")
assert 运行.读取任务(身份).步骤[0]["state"] == "pending"
运行.控制任务(身份, "author", 任务状态.待对账, "恢复", 命令ID="resume")
下一 = 运行.领取步骤("new", ["collect", "finish"])
assert 下一 is not None and 下一.步骤ID == "collect"
运行.执行一步(下一)
收尾 = 运行.领取步骤("new", ["finish"])
assert 收尾 is not None
运行.执行一步(收尾)
assert 运行.读取任务(身份).状态 is 任务状态.已完成
assert 运行.读取任务(身份).步骤[0]["result"] == {"recovered": 1}