218 lines
7.3 KiB
Python
218 lines
7.3 KiB
Python
"""连接池配置解耦、校验及 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]
|