332 lines
14 KiB
Python
332 lines
14 KiB
Python
"""真实 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 psycopg import sql
|
|
|
|
from muse.任务运行.接口 import 任务服务, 任务请求, 步骤处理器, 步骤结果, 步骤计划
|
|
from muse.任务运行.模型 import 内容哈希
|
|
from muse.任务运行.配置版本 import (
|
|
凭据引用,
|
|
提供方配置,
|
|
运行配置内容,
|
|
配置版本管理,
|
|
配置版本错误,
|
|
配置验证证据,
|
|
)
|
|
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_budget_{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_BUDGET_{声明.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 新任务(工厂: 数据库工厂):
|
|
登记 = 流程登记()
|
|
登记.登记处理器(步骤处理器("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 管理
|
|
|
|
|
|
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")
|
|
|
|
|
|
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
|
|
|
|
|
|
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.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.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"
|
|
|
|
|
|
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")
|
|
|
|
|
|
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.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",
|
|
)
|
|
|
|
|
|
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"
|
|
|
|
|
|
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")
|