"""Серверные события (ТЗ 3.14): шина публикации и SSE-генератор.
TestClient буферизует весь ответ (portal.call ждёт завершения приложения),
поэтому бесконечный SSE-стрим через client.stream нечитаем в принципе —
генератор `_event_source` тестируем напрямую, HTTP-публикации — перехватом
publish.
"""
import asyncio
import threading
from collections.abc import AsyncIterator, Iterator
from typing import Any
import pytest
from app.api import events as events_api
from app.realtime import Event, bus, publish
@pytest.fixture(autouse=True)
def _fresh_bus() -> Iterator[None]:
"""Шина между тестами пуста: asyncio.run юнит-тестов оставляет мёртвый loop."""
bus._subscribers.clear()
bus._loop = None
yield
bus._subscribers.clear()
bus._loop = None
class _FakeRequest:
"""Минимальный двойник Request для _event_source (is_disconnected → False)."""
async def is_disconnected(self) -> bool:
return False
async def _anext_timeout(gen: AsyncIterator[str], timeout: float = 2.0) -> str:
return await asyncio.wait_for(gen.__anext__(), timeout=timeout)
# --- юниты шины ---
def test_dispatch_same_loop() -> None:
"""Подписка и публикация в текущем loop: событие попадает в очередь."""
async def scenario() -> Any:
queue = bus.subscribe("u1")
try:
publish("u1", "task.changed", {"id": 1})
return await asyncio.wait_for(queue.get(), timeout=2)
finally:
bus.unsubscribe("u1", queue)
event = asyncio.run(scenario())
assert event.kind == "task.changed"
assert event.data == {"id": 1}
def test_publish_from_thread() -> None:
"""publish из другого потока (как MCP/detail_task) доезжает через call_soon_threadsafe."""
async def scenario() -> Any:
queue = bus.subscribe("u1")
try:
thread = threading.Thread(
target=publish, args=(None, "garden.changed"), kwargs={"data": {"reason": "shop"}}
)
thread.start()
return await asyncio.wait_for(queue.get(), timeout=2)
finally:
bus.unsubscribe("u1", queue)
event = asyncio.run(scenario())
assert event.kind == "garden.changed"
def test_queue_overflow_skips() -> None:
"""Переполнение очереди не бросает: событие пропускается (coarse-модель)."""
async def scenario() -> None:
queue = bus.subscribe("u1")
try:
for i in range(100):
publish("u1", "task.changed", {"i": i})
# очередь забита до maxsize — следующий publish не бросает
publish("u1", "task.changed", {"overflow": True})
finally:
bus.unsubscribe("u1", queue)
asyncio.run(scenario())
def test_no_subscribers_is_noop() -> None:
"""Publish без подписчиков — тихий no-op (клиент догонит refetch'ем)."""
publish(None, "task.changed", {"id": 1})
# --- SSE-генератор ---
def test_event_source_ready_then_publish() -> None:
"""Первый кадр ready; событие шины доезжает в стрим; отсоединение завершает."""
async def scenario() -> list[str]:
gen = events_api._event_source(_FakeRequest(), "u1") # type: ignore[arg-type]
frames: list[str] = []
frames.append(await _anext_timeout(gen)) # subscribe внутри → loop забинден
publish("u1", "project.changed", {"id": 5})
frames.append(await _anext_timeout(gen))
await gen.aclose()
return frames
frames = asyncio.run(scenario())
assert frames[0] == "data: " + events_api._sse_json("ready") + "\n\n"
assert '"project.changed"' in frames[1]
assert '"id": 5' in frames[1]
# после aclose подписчик отписан
assert bus._subscribers == {}
def test_event_source_unsubscribe_on_disconnect() -> None:
"""is_disconnected → генератор завершается, подписка снята (finally)."""
class _Disconnecting:
def __init__(self) -> None:
self.calls = 0
async def is_disconnected(self) -> bool:
self.calls += 1
return True # ready уходит без проверки, затем сразу disconnect
async def scenario() -> None:
gen = events_api._event_source(_Disconnecting(), "u1") # type: ignore[arg-type]
await _anext_timeout(gen)
with pytest.raises(StopAsyncIteration):
await _anext_timeout(gen)
asyncio.run(scenario())
assert bus._subscribers == {}
def test_event_source_ping_on_idle() -> None:
"""Тишина в шине дольше PING_SECONDS → кадр ping (keep-alive для прокси)."""
async def scenario() -> str:
gen = events_api._event_source(_FakeRequest(), "u1") # type: ignore[arg-type]
await _anext_timeout(gen) # ready
# подменить таймаут на минимум, чтобы не ждать 20 секунд
saved = events_api.PING_SECONDS
events_api.PING_SECONDS = 0.01
try:
return await _anext_timeout(gen)
finally:
events_api.PING_SECONDS = saved
await gen.aclose()
frame = asyncio.run(scenario())
assert '"ping"' in frame
# --- HTTP-публикации (перехват publish: TestClient стрим не читает) ---
@pytest.fixture
def published(monkeypatch: pytest.MonkeyPatch) -> list[Event]:
"""Перехват publish в модуле api.tasks: события собираются в список."""
seen: list[Event] = []
import app.api.tasks as tasks_api
def _spy(user_id: str | None, kind: str, data: dict[str, Any] | None = None) -> None:
seen.append(Event(kind=kind, data=data))
publish(user_id, kind, data)
monkeypatch.setattr(tasks_api, "publish", _spy)
return seen
def test_create_task_publishes(client: Any, published: list[Event]) -> None:
"""Создание задачи → task.changed + xp.changed (без конфетти)."""
created = client.post("/api/tasks", json={"title": "Тест SSE"})
assert created.status_code == 200
kinds = [e.kind for e in published]
assert "task.changed" in kinds
assert "xp.changed" in kinds
xp = next(e for e in published if e.kind == "xp.changed")
assert xp.data == {"celebrate": False}
def test_close_task_publishes_xp(client: Any, published: list[Event]) -> None:
"""Закрытие задачи → task.changed + xp.changed {celebrate: True} + garden.changed."""
task_id = client.post("/api/tasks", json={"title": "Закрыть"}).json()["id"]
published.clear()
approved = client.post(f"/api/tasks/{task_id}/approve")
assert approved.status_code == 200
published.clear()
closed = client.patch(f"/api/tasks/{task_id}", json={"status": "done"})
assert closed.status_code == 200
kinds = [e.kind for e in published]
assert "task.changed" in kinds
assert "xp.changed" in kinds
assert "garden.changed" in kinds
xp = next(e for e in published if e.kind == "xp.changed")
assert xp.data is not None and xp.data["celebrate"] is True
def test_delete_task_publishes_deleted(client: Any, published: list[Event]) -> None:
"""Удаление задачи → task.deleted {id}."""
task_id = client.post("/api/tasks", json={"title": "Удалить"}).json()["id"]
published.clear()
deleted = client.delete(f"/api/tasks/{task_id}")
assert deleted.status_code == 200
assert [e for e in published if e.kind == "task.deleted" and e.data == {"id": task_id}]