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