"""真实 PG 的预算原子预留、成本状态与配置版本;不调用外部模型。""" from __future__ import annotations import dataclasses import time import uuid from concurrent.futures import ThreadPoolExecutor from datetime import UTC, datetime, timedelta from decimal import Decimal from pathlib import Path from threading import Barrier import psycopg import pytest from muse.任务运行.接口 import 任务服务, 任务请求, 步骤处理器, 步骤结果, 步骤计划 from muse.任务运行.模型 import 内容哈希 from muse.任务运行.配置版本 import ( 凭据引用, 提供方配置, 运行配置内容, 配置版本管理, 配置版本错误, 配置验证证据, ) 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_BUDGET_{声明.name}_URL", Path(当前.引用.位置).read_text()) yield 工厂 def 新任务(工厂: 数据库工厂): 登记 = 流程登记() 登记.登记处理器(步骤处理器("call", "1", lambda _: 步骤结果({}), "v1", "v1")) 登记.登记类型("budget-test", 必需保护=()) 运行 = 任务服务(工厂, 登记) 编排 = 流程服务(运行, 登记) 编排.发布(流程定义("budget-test", "1", (步骤计划("call", "call", "1"),))) 请求 = 任务请求( "budget-test", uuid.uuid4().hex, "author", 工厂.用途, 内容用途.生成, {}, "policy-1", "release-1", { "source_scope": {}, "schema_versions": {}, "authorization": "grant", "budget": {"approved": True}, "stop_conditions": ["cancelled"], }, ) 身份 = 运行.创建任务(请求, "budget-test", "1") 领取 = 运行.领取步骤("worker", ["call"]) assert 领取 is not None and 领取.任务ID == 身份 return 身份, 领取 def 建预算(工厂: 数据库工厂, 任务ID: str, *, 窗口金额="24", 单次="6", 次数=3, 窗口次数=6000): 管理 = 预算管理(工厂, "quota") 管理.登记策略(额度策略("quota", "1", Decimal(窗口金额), 窗口次数)) 管理.登记任务预算( 任务ID, 任务预算计划( Decimal(单次) * 次数, (角色预算("writer", 次数, 次数, Decimal(单次)),), "approval", datetime.now(UTC) + timedelta(minutes=10), ), ) return 管理 @pytest.mark.case_id( "NC-budget-concurrent-reserve", environment="隔离 PostgreSQL;配置验证器为显式协议替身", given="两个真实任务共用10额度账户", when="两个独立连接同时预留6", then=["只有一笔成功且余额为4"], contract="docs/系统架构/新版设计/模块设计/S02-任务运行.md", ) def test_并发预留共享窗口且不超额__a61001(预算环境) -> None: A, 领取A = 新任务(预算环境[用途.生产]) B, 领取B = 新任务(预算环境[用途.生产]) 预算A = 建预算(预算环境[用途.生产], A, 窗口金额="10") 预算B = 建预算(预算环境[用途.生产], B, 窗口金额="10") 栅栏 = Barrier(2) def 预留(输入): 管理, 领取 = 输入 栅栏.wait() try: return 管理.预留(领取, 领取.任务ID, "writer") except 预算不足: return None with ThreadPoolExecutor(max_workers=2) as 线程: 结果 = list(线程.map(预留, [(预算A, 领取A), (预算B, 领取B)])) assert sum(r is not None for r in 结果) == 1 assert 预算A.窗口余额()["在途预留"] == Decimal("6") assert 预算A.窗口余额()["可用金额"] == Decimal("4") @pytest.mark.case_id( "NC-budget-idempotent-settlement", environment="隔离 PostgreSQL;配置验证器为显式协议替身", given="同调用身份和失败调用可信成本回执", when="重复预留、外发申请与结算", then=["只授权外发一次、只计费一次,冲突成本拒绝"], contract="docs/系统架构/新版设计/模块设计/S02-任务运行.md", ) def test_重复预留发送和结算不重复计费__a61002(预算环境) -> None: 身份, 领取 = 新任务(预算环境[用途.生产]) 预算 = 建预算(预算环境[用途.生产], 身份) 初次 = 预算.预留(领取, "call-1", "writer") assert 预算.预留(领取, "call-1", "writer") == 初次 assert 预算.标记已发送(领取, "call-1", 最长秒=30).允许外发 assert not 预算.标记已发送(领取, "call-1", 最长秒=30).允许外发 结果 = 预算.结算("call-1", Decimal("0.125"), 回执ID="failed-call-receipt") assert 预算.结算("call-1", Decimal("0.125"), 回执ID="failed-call-receipt") == 结果 with pytest.raises(预算状态冲突): 预算.结算("call-1", Decimal("0.2"), 回执ID="changed") assert 预算.窗口余额()["已知成本"] == Decimal("0.125") assert 预算.窗口余额()["已占次数"] == 1 @pytest.mark.case_id( "NC-budget-unknown-across-window", environment="隔离 PostgreSQL;配置验证器为显式协议替身", given="已发送无成本回执的调用且注入旧窗口", when="尝试新调用并进行可信对账", then=["未知不计零且跨窗阻断,对账后恢复"], contract="docs/系统架构/新版设计/模块设计/S02-任务运行.md", ) def test_未知成本跨窗口仍阻止后续调用直至对账__a61003(预算环境) -> None: 身份, 领取 = 新任务(预算环境[用途.生产]) 预算 = 建预算(预算环境[用途.生产], 身份) 预算.预留(领取, "unknown", "writer") 预算.标记已发送(领取, "unknown", 最长秒=30) 未知 = 预算.结算("unknown", None, 回执ID="missing-cost") assert 未知.实际金额 is None and 未知.状态 == "unknown" with 预算环境[用途.维护].连接() as 连: 连.execute( "UPDATE public.muse_budget_reservation SET window_start=window_start-interval '1 " "day',window_end=window_end-interval '1 day' WHERE call_id='unknown'" ) assert 预算.窗口余额()["未知调用数"] == 1 with pytest.raises(预算不足, match="未知成本"): 预算.预留(领取, "next", "writer") 预算.结算("unknown", Decimal("0.2"), 回执ID="resolved-receipt") assert 预算.预留(领取, "next", "writer").状态 == "reserved" @pytest.mark.case_id( "NC-budget-cancellation", environment="隔离 PostgreSQL;配置验证器为显式协议替身", given="未发送或已发送的预留", when="关闭任务预算", then=["未发送释放、已发送未知且任务不再外发"], contract="docs/系统架构/新版设计/模块设计/S02-任务运行.md", ) @pytest.mark.parametrize("已发出", [False, True], ids=["before-send", "after-send"]) def test_取消保留已发送调用的未知成本__a61004(预算环境, 已发出: bool) -> None: 身份, 领取 = 新任务(预算环境[用途.生产]) 预算 = 建预算(预算环境[用途.生产], 身份) 预算.预留(领取, "cancel", "writer") if 已发出: 预算.标记已发送(领取, "cancel", 最长秒=30) 预算.关闭任务预算(身份) assert 预算.读取("cancel").状态 == ("unknown" if 已发出 else "released") with pytest.raises(预算不足): 预算.预留(领取, "blocked", "writer") @pytest.mark.case_id( "NC-budget-expiration", environment="隔离 PostgreSQL;配置验证器为显式协议替身", given="未发送或在途预留期限到达", when="真实时间过期后回收", then=["未发送释放,在途保留未知成本"], contract="docs/系统架构/新版设计/模块设计/S02-任务运行.md", ) @pytest.mark.parametrize("已发出", [False, True], ids=["reserved", "in-flight"]) def test_过期只释放尚未发送的预留__a61005(预算环境, 已发出: bool) -> None: 身份, 领取 = 新任务(预算环境[用途.生产]) 预算 = 建预算(预算环境[用途.生产], 身份) 预算.预留(领取, "expire", "writer", 有效秒=0.08 if not 已发出 else 30) if 已发出: 预算.标记已发送(领取, "expire", 最长秒=0.08) time.sleep(0.12) 预算.回收过期() assert 预算.读取("expire").状态 == ("unknown" if 已发出 else "released") if not 已发出: assert 预算.窗口余额()["在途预留"] == 0 assert 预算.预留(领取, "fresh", "writer").状态 == "reserved" @pytest.mark.case_id( "NC-budget-actual-over-cap", environment="隔离 PostgreSQL;配置验证器为显式协议替身", given="实际成本超过单次预留金额", when="用真实回执结算", then=["保存真实金额和超额标记,停止任务预算"], contract="docs/系统架构/新版设计/模块设计/S02-任务运行.md", ) def test_超单次成本先保存事实并关闭任务预算__a61006(预算环境) -> None: 身份, 领取 = 新任务(预算环境[用途.生产]) 预算 = 建预算(预算环境[用途.生产], 身份) 预算.预留(领取, "over", "writer") 预算.标记已发送(领取, "over", 最长秒=30) 已记 = 预算.结算("over", Decimal("7"), 回执ID="actual-over") assert 已记.超预算 and 已记.实际金额 == Decimal("7") assert 预算.窗口余额()["已知成本"] == Decimal("7") with pytest.raises(预算不足): 预算.预留(领取, "next", "writer") @pytest.mark.case_id( "NC-budget-settlement-rollback", environment="隔离 PostgreSQL;配置验证器为显式协议替身", given="调用在途及PG注入结算失败触发器", when="结算回滚后移除故障再结算", then=["故障保留在途标记和未知实际金额,可恢复结算"], contract="docs/系统架构/新版设计/模块设计/S02-任务运行.md", ) def test_结算失败回滚仍保留在途标记__a61007(预算环境) -> None: 身份, 领取 = 新任务(预算环境[用途.生产]) 预算 = 建预算(预算环境[用途.生产], 身份) 预算.预留(领取, "atomic", "writer") 预算.标记已发送(领取, "atomic", 最长秒=30) with 预算环境[用途.维护].连接() as 连: 连.execute( "CREATE FUNCTION public.reject_settlement() RETURNS trigger LANGUAGE plpgsql AS $$ " "BEGIN IF NEW.state='settled' THEN RAISE EXCEPTION 'synthetic failure'; END IF; " "RETURN NEW; END $$" ) 连.execute( "CREATE TRIGGER reject_settlement BEFORE UPDATE ON public.muse_budget_reservation " "FOR EACH ROW EXECUTE FUNCTION public.reject_settlement()" ) with pytest.raises(psycopg.errors.RaiseException): 预算.结算("atomic", Decimal("1"), 回执ID="receipt") assert 预算.读取("atomic").状态 == "in_flight" assert 预算.读取("atomic").实际金额 is None with 预算环境[用途.维护].连接() as 连: 连.execute("DROP TRIGGER reject_settlement ON public.muse_budget_reservation") assert 预算.结算("atomic", Decimal("1"), 回执ID="receipt").状态 == "settled" @pytest.mark.case_id( "NC-budget-call-limits", environment="隔离 PostgreSQL;配置验证器为显式协议替身", given="窗口次数或任务计划只剩一次", when="结算零成本调用后尝试下一次", then=["两类次数限制都独立生效"], contract="docs/系统架构/新版设计/模块设计/S02-任务运行.md", ) @pytest.mark.parametrize("模式", ["window", "task"], ids=["window", "task"]) def test_窗口次数和任务计划分别限制调用__a61008(预算环境, 模式: str) -> None: 身份, 领取 = 新任务(预算环境[用途.生产]) 预算 = 建预算( 预算环境[用途.生产], 身份, 窗口次数=1 if 模式 == "window" else 6000, 次数=3 if 模式 == "window" else 1, ) 预算.预留(领取, "first", "writer") 预算.标记已发送(领取, "first", 最长秒=30) 预算.结算("first", Decimal("0"), 回执ID="free-but-called") with pytest.raises(预算不足, match="次数"): 预算.预留(领取, "second", "writer") class 合成配置验证器: 身份 = "synthetic-contract-check-v1" def 验证(self, 内容: 运行配置内容, 执行用途: 用途) -> 配置验证证据: assert 内容.角色配置["writer"]["model"] == "model" return 配置验证证据( 内容哈希(内容.冻结()), 内容.角色策略版本, 内容.资源发布身份, 执行用途, ("synthetic-protocol-result",), "offline_contract", ) def 配置内容(版本="1"): return 运行配置内容( "direct", 版本, "policy-1", "release-1", "quota-1", {"writer": {"provider": "provider", "model": "model", "thinking": "high"}}, (凭据引用("api", "环境变量", "MUSE_API_KEY"),), (提供方配置("provider", "responses", "https://provider.invalid/v1/responses", "api"),), "price-1", ) @pytest.mark.case_id( "NC-config-version-freeze", environment="隔离 PostgreSQL;配置验证器为显式协议替身", given="显式合成验证器和两版配置", when="验证启用、任务冻结、升级和停用", then=["版本不可覆盖,CAS与回执匹配,旧任务副本保留"], contract="docs/系统架构/新版设计/模块设计/S02-任务运行.md", ) def test_配置验证启用与任务冻结版本互不覆盖__a61009(预算环境) -> None: 管理 = 配置版本管理(预算环境[用途.评测], 合成配置验证器()) A, _ = 新任务(预算环境[用途.评测]) 原 = 管理.保存草案("config", "1", 配置内容()) with pytest.raises(配置版本错误): 管理.冻结到任务(A, "config") with pytest.raises(配置版本错误): 管理.保存草案("config", "1", 配置内容("changed")) R1 = 管理.验证版本("config", "1") assert 管理.启用("config", "1", 验证回执=R1, 批准引用="approve-1", 预期代次=0) == 1 assert 管理.启用("config", "1", 验证回执=R1, 批准引用="approve-1", 预期代次=0) == 1 冻结 = 管理.冻结到任务(A, "config") 管理.保存草案("config", "2", 配置内容("2")) R2 = 管理.验证版本("config", "2") with pytest.raises(配置版本错误): 管理.启用("config", "2", 验证回执=R1, 批准引用="approve-2", 预期代次=1) assert 管理.启用("config", "2", 验证回执=R2, 批准引用="approve-2", 预期代次=1) == 2 assert 管理.冻结到任务(A, "config") == 冻结 == 原 B, _ = 新任务(预算环境[用途.评测]) assert 管理.冻结到任务(B, "config").版本 == "2" 管理.停用("config", 预期代次=2) C, _ = 新任务(预算环境[用途.评测]) with pytest.raises(配置版本错误): 管理.冻结到任务(C, "config") assert 管理.冻结到任务(A, "config").版本 == "1" @pytest.mark.case_id( "NC-config-validation-binding", environment="隔离 PostgreSQL;配置验证器为显式协议替身", given="陈旧绑定或离线证据用于生产", when="验证配置版本", then=["绑定不符和离线生产都拒绝"], contract="docs/系统架构/新版设计/模块设计/S02-任务运行.md", ) def test_配置不接受陈旧验证或离线证据启用生产__a6100a(预算环境) -> None: class 陈旧验证器(合成配置验证器): def 验证(self, 内容, 执行用途): return dataclasses.replace(super().验证(内容, 执行用途), 配置哈希="old-hash") 管理 = 配置版本管理(预算环境[用途.评测], 陈旧验证器()) 管理.保存草案("config", "1", 配置内容()) with pytest.raises(配置版本错误, match="绑定"): 管理.验证版本("config", "1") 生产 = 配置版本管理(预算环境[用途.生产], 合成配置验证器()) 生产.保存草案("config", "1", 配置内容()) with pytest.raises(配置版本错误, match="离线"): 生产.验证版本("config", "1")