muse-agent-example/tests/单元/test_连接池配置与退避.py

218 lines
7.3 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""连接池配置解耦、校验及 SSE 空轮询退避机制验证。"""
import asyncio
import pytest
from muse.共享.错误 import 配置错误
from muse.配置 import 数据库引用, 解析配置, 连接池配置
@pytest.mark.case_id("TC-POOL-CONFIG-VALIDATION")
def test_连接池配置字段强校验():
# 正常值
cfg = 连接池配置(最大连接=8, 等待秒=5.0)
assert cfg.最大连接 == 8
assert cfg.等待秒 == 5.0
# 最大连接不接受非 int、bool、越界
with pytest.raises(配置错误, match="连接池.最大连接"):
连接池配置(最大连接=True) # type: ignore
with pytest.raises(配置错误, match="连接池.最大连接"):
连接池配置(最大连接=0)
with pytest.raises(配置错误, match="连接池.最大连接"):
连接池配置(最大连接=101)
with pytest.raises(配置错误, match="连接池.最大连接"):
连接池配置(最大连接="8") # type: ignore
# 等待秒不接受 bool、NaN、Inf、越界
with pytest.raises(配置错误, match="连接池.等待秒"):
连接池配置(等待秒=False) # type: ignore
with pytest.raises(配置错误, match="连接池.等待秒"):
连接池配置(等待秒=float("nan"))
with pytest.raises(配置错误, match="连接池.等待秒"):
连接池配置(等待秒=float("inf"))
with pytest.raises(配置错误, match="连接池.等待秒"):
连接池配置(等待秒=0.05)
with pytest.raises(配置错误, match="连接池.等待秒"):
连接池配置(等待秒=61.0)
@pytest.mark.case_id("TC-POOL-CONFIG-DECOUPLING")
def test_连接池配置不改变数据库引用相等性():
ref1 = 数据库引用("环境变量", "DB_URL")
ref2 = 数据库引用("环境变量", "DB_URL")
assert ref1 == ref2
# 配置解析产出独立连接池配置
cfg = 解析配置(
{
"数据库": {
"取值方式": "环境变量",
"位置": "DB_URL",
"连接池": {"最大连接": 10, "等待秒": 3.0},
},
"资源": {"发布身份": "identity-1"},
}
)
assert cfg.数据库 == ref1
assert cfg.连接池.最大连接 == 10
assert cfg.连接池.等待秒 == 3.0
@pytest.mark.case_id("TC-POOL-CONFIG-STARTUP-INTEGRATION")
def test_启动构建消费连接池配置(monkeypatch):
from muse.启动 import 构建
# 1. 显式配置生效
cfg = 解析配置(
{
"数据库": {
"取值方式": "环境变量",
"位置": "DB_URL",
"连接池": {"最大连接": 12, "等待秒": 4.5},
},
"资源": {"发布身份": "identity-1"},
"运行": {"用途": "production"},
}
)
装配 = 构建(cfg)
assert 装配.数据库.最大连接 == 12
assert 装配.数据库.等待秒 == 4.5
# 2. 未配置时回退默认(生产 4 / 2.0,评测 2 / 2.0)
cfg_prod = 解析配置(
{
"数据库": {"取值方式": "环境变量", "位置": "DB_URL"},
"资源": {"发布身份": "identity-1"},
"运行": {"用途": "production"},
}
)
装配_prod = 构建(cfg_prod)
assert 装配_prod.数据库.最大连接 == 4
assert 装配_prod.数据库.等待秒 == 2.0
cfg_eval = 解析配置(
{
"数据库": {"取值方式": "环境变量", "位置": "DB_URL"},
"资源": {"发布身份": "identity-1"},
"运行": {"用途": "evaluation"},
}
)
装配_eval = 构建(cfg_eval)
assert 装配_eval.数据库.最大连接 == 2
assert 装配_eval.数据库.等待秒 == 2.0
@pytest.mark.case_id("TC-SSE-BACKOFF-ADAPTIVE")
def test_SSE空轮询退避机制(monkeypatch):
from types import SimpleNamespace
from muse.任务运行.模型 import 事件类型, 事件续接, 运行事件
from muse.接入.http import 事件订阅
def 构造事件(seq):
from datetime import UTC, datetime
return 运行事件(
f"event-{seq}",
"task",
None,
None,
seq,
事件类型.步骤完成,
datetime(2026, 9, 17, tzinfo=UTC),
1,
{"step": f"step-{seq}"},
)
class 模拟事件服务:
def __init__(self, 页集):
self.页集 = 页集
def 续接事件(self, task_id, cursor, 数量=100):
return self.页集.get(cursor, 事件续接((), False, cursor, cursor, None))
def 运行订阅(protocol, 页集, 循环次数=4):
服务 = 模拟事件服务(页集)
monkeypatch.setattr(事件订阅, "取得任务服务", lambda _: 服务)
monkeypatch.setattr(事件订阅, "读取作者任务", lambda req, tid, wid=None: None)
等待记录 = []
async def 记录等待(秒):
等待记录.append(秒)
monkeypatch.setattr(事件订阅.asyncio, "sleep", 记录等待)
检查次数 = [0]
async def 已断开():
检查次数[0] += 1
return 检查次数[0] > 循环次数
async def 执行():
响应 = await 事件订阅.订阅事件(
"task-1",
SimpleNamespace(is_disconnected=已断开),
"author",
cursor=0,
last_event_id=None,
protocol=protocol,
)
return [块 async for 块 in 响应.body_iterator]
asyncio.run(执行())
return 等待记录
# 1. bounded-v2: 连续空轮询从 1.0s 退避到 2.0s
等待_v2 = 运行订阅("bounded-v2", {}, 循环次数=4)
assert 等待_v2 == [1.0, 2.0, 2.0, 2.0]
# 2. legacy: 始终保持 1.0s
等待_legacy = 运行订阅("legacy", {}, 循环次数=4)
assert 等待_legacy == [1.0, 1.0, 1.0, 1.0]
# 3. bounded-v2: 退避后有新事件到达,立即重置回 1.0s
class 动态服务:
def __init__(self):
self.count = 0
def 续接事件(self, task_id, cursor, 数量=100):
self.count += 1
if self.count == 3:
return 事件续接((构造事件(1),), False, 1, 1, None)
return 事件续接((), False, cursor, cursor, None)
服务动态 = 动态服务()
monkeypatch.setattr(事件订阅, "取得任务服务", lambda _: 服务动态)
monkeypatch.setattr(事件订阅, "读取作者任务", lambda req, tid, wid=None: None)
动态等待 = []
async def 记录动态等待(秒):
动态等待.append(秒)
monkeypatch.setattr(事件订阅.asyncio, "sleep", 记录动态等待)
检查 = [0]
async def 已断开_动态():
检查[0] += 1
return 检查[0] > 4
async def 执行动态():
响应 = await 事件订阅.订阅事件(
"task-1",
SimpleNamespace(is_disconnected=已断开_动态),
"author",
cursor=0,
last_event_id=None,
protocol="bounded-v2",
)
return [块 async for 块 in 响应.body_iterator]
asyncio.run(执行动态())
# 轮询1: 空 (等待1.0)
# 轮询2: 空 (等待2.0)
# 轮询3: 有事件 (重置为0,等待1.0)
# 轮询4: 空 (等待1.0)
assert 动态等待 == [1.0, 2.0, 1.0, 1.0]