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