"""公开任务接口的持久状态、独立连接竞争与恢复合同。""" 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 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-scope-competition", environment="隔离PostgreSQL", given="同作品两任务与另一作品任务", when="两个独立连接并发领取", then=["同范围仅一个持有者,另一范围可并行"], contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md", ) 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})) @pytest.mark.case_id( "NC-task-lease-fencing", environment="隔离PostgreSQL", given="持久检查点及已过期租约", when="新执行者接管,旧执行者提交或续租", then=["旧身份被拒,新持有代次递增且检查点保留"], contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md", ) 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} @pytest.mark.case_id( "NC-task-cancel-fencing", environment="隔离PostgreSQL", given="正在执行且持有作品范围的任务", when="作者取消后旧结果到达", then=["拒绝迟到写入并允许新任务取得范围"], contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md", ) 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 任务状态.已取消 @pytest.mark.case_id( "NC-task-resume-revalidation", environment="隔离PostgreSQL", given="暂停任务及已保存检查点", when="重建服务并提供业务恢复校验", then=["校验失败仍暂停,合法恢复保留检查点"], contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md", ) 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} @pytest.mark.case_id( "NC-task-registered-revalidation", environment="隔离 PostgreSQL、具名接入;合成 owner 检查", given="暂停任务和登记的合成 owner 来源、结构、授权与预算检查", when="从 HTTP/CLI 共用入口恢复并重放命令", then=["当前来源不匹配时保持暂停", "重验通过后排队,重复命令返回相同回执"], contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md", ) 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", 控制) == 回执 @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} @pytest.mark.case_id( "NC-task-event-reconnect", environment="隔离PostgreSQL", given="有序事件及固定事件身份", when="重复提交、同 ID 不同内容及实际删除中间事件", then=["重复幂等、冲突拒绝、游标缺口要求任务快照"], contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md", ) 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 运行.续接事件(身份, 事件.序号).事件 == () @pytest.mark.case_id( "NC-task-failure-stop", environment="隔离PostgreSQL", given="处理器失败和待执行任务", when="记录失败并停止当前进程领取", then=["失败不变成功,重建服务可领取持久任务"], contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md", ) 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 @pytest.mark.case_id( "NC-task-purpose-isolation", environment="隔离PostgreSQL", given="相同命令的生产和评测任务", when="各自数据库角色领取和查询", then=["分别持久执行且不能跨用途读取任务"], contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md", ) 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) @pytest.mark.case_id( "NC-task-short-transaction", environment="隔离PostgreSQL", given="已领取任务的真实处理器", when="隔离集群管理员检查 pg_stat_activity", then=["处理器期间目标库不存在悬挂事务"], contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md", ) 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} @pytest.mark.case_id( "NC-task-control-idempotent", environment="隔离PostgreSQL", given="重复或异参数的任务控制请求", when="控制同一持久化任务", then=["同命令返回原回执且不重复事件,异参数拒绝"], contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md", ) 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") @pytest.mark.case_id( "NC-task-list-scope", environment="隔离PostgreSQL", given="不同作者和作品的已登记任务", when="列出当前作者和作品任务", then=["只返回匹配的真实任务"], contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md", ) 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") == [] @pytest.mark.case_id( "NC-task-step-order", environment="隔离PostgreSQL", given="不同于字典序的冻结流程", when="读取任务步骤", then=["按冻结流程顺序显示"], contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md", ) 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"] @pytest.mark.case_id( "NC-worker-heartbeat-stop", environment="隔离PostgreSQL", given="执行时间超过原租期的处理器", when="运行执行器并停止", then=["自动续租完成当前步骤,停止后不领取后续"], contract="docs/系统架构/新版设计/接口契约/任务工具与事件.md", ) 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"