"""任务租约的核心保护:重启不重复执行、未知调用先对账。 按简化决策删除其余变体用例;夹具从按例克隆改为共享库以消除建库开销。 `任务库` 与 `请求` 同时被 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}