"""合成持久事件页驱动真实SSE迭代器;不建立网络、数据库或真实浏览器。""" import asyncio import json from datetime import UTC, datetime from types import SimpleNamespace import pytest from fastapi import FastAPI from fastapi.testclient import TestClient from muse.任务运行.模型 import 事件类型, 事件续接, 运行事件 from muse.接入.http import 事件订阅 from muse.接入.http.作者会话 import 要求作者 from muse.接入.http.错误响应 import 安装错误响应 def 事件(序号): return 运行事件( f"event-{序号}", "task", None, None, 序号, 事件类型.步骤完成, datetime(2026, 9, 17, tzinfo=UTC), 1, {"step": f"step-{序号}"}, ) class 合成事件服务: def __init__(self, 页集): self.页集 = 页集 self.调用 = [] def 续接事件(self, task_id, cursor, *, 数量=100): self.调用.append((task_id, cursor, 数量)) assert cursor in self.页集, f"意外续接游标:{cursor}" return self.页集[cursor] def 装配替身(monkeypatch, 服务): 授权 = [] monkeypatch.setattr(事件订阅, "取得任务服务", lambda _: 服务) monkeypatch.setattr( 事件订阅, "读取作者任务", lambda request, task_id, work_id=None: 授权.append((task_id, work_id)), ) return 授权 def 读流( monkeypatch, 服务, *, protocol="bounded-v2", cursor=0, last_event_id=None, 断开检查次数=10 ): 授权 = 装配替身(monkeypatch, 服务) 检查 = [] 等待 = [] async def 已断开(): 检查.append(True) return len(检查) > 断开检查次数 async def 不等待(秒): 等待.append(秒) monkeypatch.setattr(事件订阅.asyncio, "sleep", 不等待) async def 收集(): 响应 = await 事件订阅.订阅事件( "task", SimpleNamespace(is_disconnected=已断开), "author", cursor=cursor, work_id="work", last_event_id=last_event_id, protocol=protocol, ) assert 响应.media_type == "text/event-stream" assert 响应.headers["cache-control"] == "no-cache" return [块 async for 块 in 响应.body_iterator] return asyncio.run(收集()), 授权, 等待, 检查 def 拆块(块): return dict(行.split(": ", 1) for 行 in 块.strip().splitlines() if ": " in 行) @pytest.mark.case_id("TC-O08-SSE-001") def test_v2逐页发完最终序号才收尾__o08101(monkeypatch): 服务 = 合成事件服务( { 0: 事件续接((事件(1), 事件(2)), False, 1, 3, 3), 2: 事件续接((事件(3),), False, 1, 3, 3), } ) 块, 授权, 等待, 检查 = 读流(monkeypatch, 服务) 数据 = [拆块(项) for 项 in 块] assert [项["event"] for 项 in 数据] == ["task-event"] * 3 + ["stream-complete"] assert [项["id"] for 项 in 数据[:3]] == ["1", "2", "3"] assert [json.loads(项["data"])["sequence"] for 项 in 数据[:3]] == [1, 2, 3] assert json.loads(数据[-1]["data"]) == {"execution_final_sequence": 3} assert 服务.调用 == [("task", 0, 100), ("task", 2, 100)] assert 授权 == [("task", "work")] assert 等待 == [1] and len(检查) == 2 @pytest.mark.case_id("TC-O08-SSE-002") @pytest.mark.parametrize("cursor,last_event_id", [(3, None), (1, 3), (3, 1)]) def test_v2最终游标重连直接结束__o08102(monkeypatch, cursor, last_event_id): 服务 = 合成事件服务({3: 事件续接((), False, 1, 3, 3)}) 块, _, 等待, _ = 读流(monkeypatch, 服务, cursor=cursor, last_event_id=last_event_id) assert len(块) == 1 assert 拆块(块[0])["event"] == "stream-complete" assert 服务.调用 == [("task", 3, 100)] assert 等待 == [] @pytest.mark.case_id("TC-O08-SSE-003") def test_legacy终态继续心跳保持原协议__o08103(monkeypatch): 服务 = 合成事件服务( { 0: 事件续接((事件(1),), False, 1, 1, 1), 1: 事件续接((), False, 1, 1, 1), } ) 块, _, 等待, 检查 = 读流(monkeypatch, 服务, protocol="legacy", 断开检查次数=2) assert 拆块(块[0])["event"] == "task-event" assert 块[1:] == [": keep-alive\n\n"] assert "stream-complete" not in "".join(块) assert 服务.调用 == [("task", 0, 100), ("task", 1, 100)] assert 等待 == [1, 1] and len(检查) == 3 @pytest.mark.case_id("TC-O08-SSE-004") @pytest.mark.parametrize("protocol", ["legacy", "bounded-v2"]) def test_reset要求快照后立即结束__o08104(monkeypatch, protocol): 服务 = 合成事件服务({0: 事件续接((), True, 5, 8, 8)}) 块, _, 等待, 检查 = 读流(monkeypatch, 服务, protocol=protocol) assert len(块) == 1 数据 = 拆块(块[0]) assert 数据["event"] == "reset" assert json.loads(数据["data"]) == { "items": [], "reset_required": True, "first_sequence": 5, "last_sequence": 8, "execution_final_sequence": 8, } assert len(服务.调用) == len(检查) == 1 assert 等待 == [] @pytest.mark.case_id("TC-O08-SSE-005") def test_未终态保留心跳且断开不再查询__o08105(monkeypatch): 服务 = 合成事件服务({0: 事件续接((), False, 1, 0)}) 块, _, 等待, 检查 = 读流(monkeypatch, 服务, 断开检查次数=1) assert 块 == [": keep-alive\n\n"] assert len(服务.调用) == 1 and len(检查) == 2 assert 等待 == [1] 先断 = 合成事件服务({}) assert 读流(monkeypatch, 先断, 断开检查次数=0)[0] == [] assert 先断.调用 == [] def 合成应用(monkeypatch, 服务): 装配替身(monkeypatch, 服务) app = FastAPI() app.include_router(事件订阅.路由) app.dependency_overrides[要求作者] = lambda: "author" 安装错误响应(app) return app @pytest.mark.case_id("TC-O08-SSE-006") def test_HTTP事件页边界与v2终态合同__o08106(monkeypatch): 服务 = 合成事件服务({0: 事件续接((事件(1),), False, 1, 1, 1)}) with TestClient(合成应用(monkeypatch, 服务)) as 客户端: 页 = 客户端.get("/api/v1/tasks/task/events", params={"limit": 1}) assert 页.status_code == 200 assert 页.json()["execution_final_sequence"] == 1 assert 服务.调用[-1] == ("task", 0, 1) assert 客户端.get("/api/v1/tasks/task/events", params={"limit": 100}).status_code == 200 assert 服务.调用[-1] == ("task", 0, 100) 响应 = 客户端.get("/api/v1/tasks/task/events/stream", params={"protocol": "bounded-v2"}) assert 响应.status_code == 200 assert 响应.headers["content-type"].startswith("text/event-stream") 帧 = [拆块(块) for 块 in 响应.text.strip().split("\n\n")] assert [项["event"] for 项 in 帧] == ["task-event", "stream-complete"] assert json.loads(帧[-1]["data"])["execution_final_sequence"] == 1 schema = 合成应用(monkeypatch, 服务).openapi() 参数 = schema["paths"]["/api/v1/tasks/{task_id}/events/stream"]["get"]["parameters"] assert next(p["schema"] for p in 参数 if p["name"] == "protocol")["default"] == "legacy" @pytest.mark.case_id("NC-O08-SSE-007") @pytest.mark.parametrize( "路径,参数,头", [ ("events", {"limit": 0}, {}), ("events", {"limit": 101}, {}), ("events", {"cursor": -1}, {}), ("events/stream", {"protocol": "v3"}, {}), ("events/stream", {"cursor": -1}, {}), ("events/stream", {}, {"Last-Event-ID": "-1"}), ], ) def test_HTTP无效范围不触发服务__o08107(monkeypatch, 路径, 参数, 头): 服务 = 合成事件服务({}) with TestClient(合成应用(monkeypatch, 服务)) as 客户端: 响应 = 客户端.get(f"/api/v1/tasks/task/{路径}", params=参数, headers=头) assert 响应.status_code == 422 assert 服务.调用 == []