muse-agent-example/tests/集成/test_任务租约与恢复.py
zizi f76c3cd04a docs(架构): 归档R2历史留痕并清理过程计划与实施分期标签
- 归档并收敛 R2 历史留痕至 docs/实现回顾/R2改造历史留痕.md,删除已退出生命周期的改造计划、旧文件处置及验证过程文档
- 现行测试用例身份收敛至 tests/用例清单.json(1358条扁平登记),适配 conftest、索引维护与测试执行验证
- 清理源码、配置、SQL迁移头、测试夹具及规则文档中的 Wxx 实施分期标签与过时过程描述
- 同步重新打包构建资源清单,通过全量离线测试、模块边界、类型检查与索引双向强核验
2026-09-16 19:09:14 +08:00

440 lines
19 KiB
Python

"""公开任务接口的持久状态、独立连接竞争与恢复合同。"""
from __future__ import annotations
import dataclasses
import subprocess
import sys
import time
import uuid
from concurrent.futures import ThreadPoolExecutor
from pathlib import Path
from threading import Barrier
import psycopg
import pytest
from psycopg import sql
from muse.任务运行.接口 import (
事件类型,
任务服务,
任务状态,
任务请求,
任务错误,
作用域,
步骤处理器,
步骤结果,
步骤计划,
状态冲突,
租约失效,
)
from muse.共享.调用身份 import 内容用途, 用途
from muse.基础设施.数据库.迁移 import 执行迁移
from muse.基础设施.数据库.连接 import 数据库工厂
from muse.编排.接口 import 流程定义, 流程服务, 流程登记
from muse.配置 import 数据库引用
pytestmark = pytest.mark.数据库
根 = Path(__file__).resolve().parents[2]
@pytest.fixture
def 任务库(隔离数据库URL: str, monkeypatch: pytest.MonkeyPatch, tmp_path: Path):
库名 = f"muse_task_{uuid.uuid4().hex[:12]}"
with psycopg.connect(隔离数据库URL, autocommit=True) as 连:
连.execute((根 / "数据库/初始化/用途角色.sql").read_text())
连.execute(
sql.SQL("CREATE DATABASE {} OWNER muse_maint TEMPLATE template0").format(
sql.Identifier(库名)
)
)
工厂 = {}
try:
for 声明, 角色 in [
(用途.维护, "muse_maint"),
(用途.生产, "muse_app"),
(用途.评测, "muse_eval"),
]:
参数 = psycopg.conninfo.conninfo_to_dict(隔离数据库URL)
参数.update(dbname=库名, user=角色)
引用名 = f"MUSE_W05_{声明.name}_URL"
monkeypatch.setenv(引用名, psycopg.conninfo.make_conninfo(**参数))
工厂[声明] = 数据库工厂(数据库引用("环境变量", 引用名), 声明)
for 名称 in [
"V0001__共享标识与版本.sql",
"V0004__任务运行与证据.sql",
"V0016__流程定义与版本.sql",
]:
(tmp_path / 名称).write_bytes((根 / "数据库/迁移" / 名称).read_bytes())
with 工厂[用途.维护].连接() as 连:
执行迁移(连, tmp_path)
yield 工厂
finally:
with psycopg.connect(隔离数据库URL, autocommit=True) as 连:
连.execute(sql.SQL("DROP DATABASE {} WITH (FORCE)").format(sql.Identifier(库名)))
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"]):
服务.执行一步(领取)
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 运行.读取任务(身份).步骤)
def test_独立连接竞争同范围而不同范围可并行__a51002(任务库) -> None:
运行, 编排 = 构造服务(任务库[用途.生产])
A = 编排.发起(请求("a"), "example", "1")
B = 编排.发起(请求("b"), "example", "1")
C = 编排.发起(请求("c", "work-C"), "example", "1")
起点 = Barrier(2)
def 领取(名称):
起点.wait()
return 运行.领取步骤(名称, ["collect"])
with ThreadPoolExecutor(max_workers=2) as 线程:
结果 = list(线程.map(领取, ["one", "two"]))
有效 = [项 for 项 in 结果 if 项]
# 同范围争用者可能先让出;下一次领取应直接命中独立范围。
if len(有效) == 1:
下一 = 运行.领取步骤("third", ["collect"])
assert 下一 is not None
有效.append(下一)
assert len(有效) == 2
assert C in {项.任务ID for 项 in 有效}
assert len({项.任务ID for 项 in 有效} & {A, B}) == 1
assert 运行.领取步骤("blocked", ["collect"]) is None
for 项 in 有效:
运行.完成步骤(项, 步骤结果({"done": 1}))
def test_过期接管拒绝旧步骤和作用域代次__a51003(任务库) -> None:
运行, 编排 = 构造服务(任务库[用途.生产])
身份 = 编排.发起(请求(), "example", "1")
旧 = 运行.领取步骤("old", ["collect"], 租期秒=0.12)
assert 旧 is not None
运行.保存检查点(旧, {"offset": 4})
time.sleep(0.16)
新 = 运行.领取步骤("new", ["collect"])
assert 新 is not None and 新.任务ID == 身份
assert 新.尝试ID != 旧.尝试ID and 新.作用域代次 > 旧.作用域代次
assert 新.检查点 == {"offset": 4}
for 操作 in [lambda: 运行.完成步骤(旧, 步骤结果({"late": 1})), lambda: 运行.续租(旧)]:
with pytest.raises(租约失效):
操作()
with pytest.raises(租约失效):
运行.执行一步(dataclasses.replace(新, 处理器ID="finish"))
运行.完成步骤(新, 步骤结果({"current": 1}))
assert 运行.读取任务(身份).步骤[0]["result"] == {"current": 1}
def test_取消后拒绝旧结果并释放作品范围__a51004(任务库) -> None:
运行, 编排 = 构造服务(任务库[用途.生产])
身份 = 编排.发起(请求("old"), "example", "1")
旧 = 运行.领取步骤("old", ["collect"])
assert 旧 is not None
运行.控制任务(身份, "author", 任务状态.运行中, "取消", 命令ID="cancel")
with pytest.raises(租约失效):
运行.完成步骤(旧, 步骤结果({"late": 1}))
新任务 = 编排.发起(请求("new"), "example", "1")
新 = 运行.领取步骤("new", ["collect"])
assert 新 is not None and 新.任务ID == 新任务
assert 运行.读取任务(身份).状态 is 任务状态.已取消
def test_暂停恢复重验且检查点持久保留__a51005(任务库) -> None:
运行, 编排 = 构造服务(任务库[用途.生产])
身份 = 编排.发起(请求(), "example", "1")
旧 = 运行.领取步骤("old", ["collect"])
assert 旧 is not None
运行.保存检查点(旧, {"part": 2})
运行.控制任务(身份, "author", 任务状态.运行中, "暂停", 命令ID="pause")
def 拒绝过期来源(快照):
raise 状态冲突("来源已变化")
重启, _ = 构造服务(任务库[用途.生产], 重验=拒绝过期来源)
with pytest.raises(状态冲突, match="来源"):
重启.控制任务(身份, "author", 任务状态.已暂停, "恢复", 命令ID="resume")
assert 重启.读取任务(身份).状态 is 任务状态.已暂停
def 校验当前来源(快照):
assert 快照.冻结输入["冻结上下文"]["authorization"] == "grant-1"
重启, _ = 构造服务(任务库[用途.生产], 重验=校验当前来源)
重启.控制任务(身份, "author", 任务状态.已暂停, "恢复", 命令ID="resume")
新 = 重启.领取步骤("new", ["collect"])
assert 新 is not None and 新.检查点 == {"part": 2}
def test_具名接入恢复由登记owner重验__a5100f(任务库) -> None:
"""HTTP/CLI 共用入口不能自报重验;服务按冻结流程找到当前业务检查。"""
from muse.接入.主会话 import 任务控制请求, 执行任务控制
当前来源 = {"版本": 4}
def 重验合成任务(快照):
上下文 = 快照.冻结输入["冻结上下文"]
if 上下文["source_scope"]["as_of"] != 当前来源["版本"]:
raise 状态冲突("来源版本已变化")
assert 上下文["schema_versions"] == {"entity": "1"}
assert 上下文["authorization"] == "grant-1"
assert 上下文["budget"]["remaining"] > 0
运行, 编排 = 构造服务(任务库[用途.生产], 重验=重验合成任务)
身份 = 编排.发起(请求(), "example", "1")
运行.控制任务(身份, "author", 任务状态.待运行, "暂停", 命令ID="pause")
控制 = 任务控制请求(
command_id="resume", target_ref=身份, expected_state=任务状态.已暂停, action="恢复"
)
with pytest.raises(状态冲突, match="来源版本"):
执行任务控制(运行, "author", 控制)
assert 运行.读取任务(身份).状态 is 任务状态.已暂停
当前来源["版本"] = 3
回执 = 执行任务控制(运行, "author", 控制)
assert 回执.state is 任务状态.待运行
assert 执行任务控制(运行, "author", 控制) == 回执
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}
def test_事件重复幂等且游标缺口要求快照__a51007(任务库) -> None:
运行, 编排 = 构造服务(任务库[用途.生产])
身份 = 编排.发起(请求(), "example", "1")
领取 = 运行.领取步骤("one", ["collect"])
assert 领取 is not None
事件ID = str(uuid.uuid4())
事件 = 运行.追加事件(领取, 事件类型.文本片段, {"文本": "片段"}, 事件ID=事件ID)
assert 运行.追加事件(领取, 事件类型.文本片段, {"文本": "片段"}, 事件ID=事件ID) == 事件
with pytest.raises(状态冲突):
运行.追加事件(领取, 事件类型.文本片段, {"文本": "不同"}, 事件ID=事件ID)
续接 = 运行.续接事件(身份)
assert [e.序号 for e in 续接.事件] == list(range(1, 续接.最后序号 + 1))
assert not 续接.需要快照
with 任务库[用途.维护].连接() as 连:
连.execute("DELETE FROM public.muse_task_event WHERE task_id=%s AND sequence=2", (身份,))
assert 运行.续接事件(身份).需要快照
assert 运行.续接事件(身份, 事件.序号).事件 == ()
def test_失败如实保留且停止领取不丢状态__a51008(任务库) -> None:
def 失败处理器(_):
raise RuntimeError("合成失败")
运行, 编排 = 构造服务(任务库[用途.生产], 失败处理器)
身份 = 编排.发起(请求(), "example", "1")
领取 = 运行.领取步骤("one", ["collect"])
assert 领取 is not None
with pytest.raises(RuntimeError):
运行.执行一步(领取)
assert 运行.读取任务(身份).状态 is 任务状态.已失败
编排.发起(请求("next", "another"), "example", "1")
运行.停止领取()
assert 运行.领取步骤("stop", ["collect"]) is None
重启, _ = 构造服务(任务库[用途.生产])
assert 重启.领取步骤("restart", ["collect"]) is not None
def test_生产评测任务使用独立权限与实例__a51009(任务库) -> None:
生产, 生产编排 = 构造服务(任务库[用途.生产])
评测, 评测编排 = 构造服务(任务库[用途.评测])
A = 生产编排.发起(请求(), "example", "1")
B = 评测编排.发起(请求(用途标记=用途.评测), "example", "1")
assert A != B
assert 生产.领取步骤("prod", ["collect"]) is not None
assert 评测.领取步骤("eval", ["collect"]) is not None
with pytest.raises(任务错误, match="不存在|用途"):
评测.读取任务(A)
def test_处理器调用期间不持有数据库事务__a5100a(任务库, 隔离数据库URL: str) -> None:
def 处理器(上下文):
with 任务库[用途.维护].连接() as 连:
目标库 = 连.info.dbname
with psycopg.connect(隔离数据库URL) as 管理员:
数量 = 管理员.execute(
"SELECT count(*) FROM pg_stat_activity WHERE datname=%s "
"AND pid<>pg_backend_pid() AND state='idle in transaction'",
(目标库,),
).fetchone()[0]
assert 数量 == 0
return 步骤结果({"done": True})
运行, 编排 = 构造服务(任务库[用途.生产], 处理器)
编排.发起(请求(), "example", "1")
领取 = 运行.领取步骤("one", ["collect"])
assert 领取 is not None
assert 运行.执行一步(领取).输出 == {"done": True}
def test_控制命令重放不重复事件且异参数拒绝__a5100b(任务库) -> None:
运行, 编排 = 构造服务(任务库[用途.生产])
身份 = 编排.发起(请求(), "example", "1")
首次 = 运行.控制任务(身份, "author", 任务状态.待运行, "取消", 命令ID="cancel-one")
重放 = 运行.控制任务(身份, "author", 任务状态.待运行, "取消", 命令ID="cancel-one")
assert 重放 == 首次
assert 运行.读取任务(身份).最后序号 == 首次["last_sequence"]
with pytest.raises(任务错误):
运行.控制任务(身份, "author", 任务状态.待运行, "暂停", 命令ID="cancel-one")
def test_任务列表按作者作品及用途隔离__a5100c(任务库) -> None:
运行, 编排 = 构造服务(任务库[用途.生产])
A = 编排.发起(请求("a", "work-A"), "example", "1")
编排.发起(请求("b", "work-B"), "example", "1")
编排.发起(dataclasses.replace(请求("c", "work-A"), 作者="other"), "example", "1")
assert [x.任务ID for x in 运行.列出任务("author", 作品ID="work-A")] == [A]
assert 运行.列出任务("missing") == []
def test_任务步骤按冻结流程顺序显示__a5100d(任务库) -> None:
运行, 编排 = 构造服务(任务库[用途.生产])
编排.发布(
流程定义(
"example",
"2",
(
步骤计划("z-first", "collect", "1"),
步骤计划("a-last", "finish", "1", ("z-first",)),
),
)
)
身份 = 编排.发起(请求("ordered"), "example", "2")
assert [step["step_id"] for step in 运行.读取任务(身份).步骤] == ["z-first", "a-last"]
def test_执行器续租并在停止后不领取新步骤__a5100e(任务库) -> None:
from muse.接入.执行器 import 执行器
def 长步骤(_):
time.sleep(0.45)
return 步骤结果({"done": True})
运行, 编排 = 构造服务(任务库[用途.生产], 长步骤)
身份 = 编排.发起(请求(), "example", "1")
worker = 执行器(运行, "durable-worker", ["collect", "finish"], 租期秒=0.3)
assert worker.运行一次()
worker.停止()
assert not worker.运行一次()
assert 运行.读取任务(身份).步骤[0]["state"] == "completed"
assert 运行.读取任务(身份).步骤[1]["state"] == "pending"