"""WebSocket 日志流测试""" import asyncio import pytest from backend.services.log_streamer import LogStreamer @pytest.fixture def streamer(): return LogStreamer() async def test_emit_and_subscribe(streamer): """测试发送和接收日志""" streamer.create_queue("t1") await streamer.emit("t1", "hello", level="info") received = [] async for msg in streamer.subscribe("t1"): received.append(msg) break assert len(received) == 1 assert received[0]["message"] == "hello" assert received[0]["level"] == "info" async def test_emit_step(streamer): """测试步骤标记""" streamer.create_queue("t2") await streamer.emit_step("t2", "构建中") received = [] async for msg in streamer.subscribe("t2"): received.append(msg) break assert received[0]["message"] == "[构建中]" assert received[0]["level"] == "step" async def test_emit_error(streamer): """测试错误消息""" streamer.create_queue("t3") await streamer.emit_error("t3", "出错了") received = [] async for msg in streamer.subscribe("t3"): received.append(msg) break assert received[0]["level"] == "error" assert received[0]["message"] == "出错了" async def test_emit_warning(streamer): """测试警告消息""" streamer.create_queue("t4") await streamer.emit_warning("t4", "警告") received = [] async for msg in streamer.subscribe("t4"): received.append(msg) break assert received[0]["level"] == "warn" async def test_emit_to_nonexistent_queue(streamer): """测试向不存在的队列发送消息(不应报错)""" await streamer.emit("nonexistent", "msg") async def test_complete_signal(streamer): """测试完成信号终止订阅""" streamer.create_queue("t5") await streamer.emit("t5", "msg1") streamer.complete("t5") received = [] async for msg in streamer.subscribe("t5"): received.append(msg) assert len(received) == 1 assert received[0]["message"] == "msg1" async def test_multiple_messages(streamer): """测试多条消息顺序""" streamer.create_queue("t6") for i in range(5): await streamer.emit("t6", f"msg{i}") streamer.complete("t6") received = [] async for msg in streamer.subscribe("t6"): received.append(msg) assert len(received) == 5 for i, msg in enumerate(received): assert msg["message"] == f"msg{i}" async def test_queue_full_drops_oldest(streamer): """测试队列满时丢弃最旧消息""" q = streamer.create_queue("t7") for i in range(q.maxsize): await streamer.emit("t7", f"old{i}") await streamer.emit("t7", "new") messages = [] while not q.empty(): messages.append(await q.get()) assert messages[-1]["message"] == "new" async def test_cleanup(streamer): """测试清理队列""" streamer.create_queue("t8") await streamer.emit("t8", "msg") streamer.cleanup("t8") assert "t8" not in streamer._queues # 清理后再发送不应报错 await streamer.emit("t8", "msg2") async def test_subscribe_without_queue(streamer): """测试订阅不存在的队列(应立即结束)""" collected = [] async for msg in streamer.subscribe("nonexistent"): collected.append(msg) assert collected == []