Coverage for src/lexigram/web/transport/reactive.py: 22%
37 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-25 04:37 +0800
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-25 04:37 +0800
1"""SSE responses driven by lexigram.reactive EventStream sources."""
3from __future__ import annotations
5import asyncio
6from collections.abc import AsyncGenerator, Callable
7from contextlib import suppress
8from typing import Any
10from lexigram.reactive import EventStream
11from lexigram.web.transport.sse import EventSourceResponse, ServerSentEvent
14def sse_from_stream(
15 stream: EventStream[Any],
16 serializer: Callable[[Any], str] | None = None,
17 event_name: str | None = None,
18 keepalive: float | None = 15.0,
19) -> EventSourceResponse:
20 """Expose a reactive stream as an EventSource-compatible SSE response.
22 Args:
23 stream: Any reactive stream.
24 serializer: Optional item → frame-string serializer.
25 Defaults to ``str(item)``.
26 event_name: Optional SSE ``event:`` field value.
27 keepalive: Emit ``: keepalive`` comments after this many seconds
28 of silence; ``None`` disables.
30 Returns:
31 A Starlette StreamingResponse with ``text/event-stream``.
33 Example:
34 ```python
35 @app.get("/dashboard/events")
36 async def dashboard_events() -> EventSourceResponse:
37 return sse_from_stream(
38 orders.pipe(ops.filter(lambda e: e.type == "Order")),
39 serializer=order_serializer,
40 )
41 ```
42 """
44 async def _items() -> AsyncGenerator[ServerSentEvent, None]:
45 queue: asyncio.Queue[ServerSentEvent] = asyncio.Queue()
46 done = asyncio.Event()
47 background_tasks: set[asyncio.Task[Any]] = set()
49 async def _drain() -> None:
50 try:
51 async for item in stream:
52 frame = serializer(item) if serializer else str(item)
53 await queue.put(ServerSentEvent(data=frame, event=event_name))
54 finally:
55 done.set()
57 drain_task = asyncio.create_task(_drain())
58 background_tasks.add(drain_task)
59 drain_task.add_done_callback(background_tasks.discard)
61 try:
62 while not done.is_set():
63 try:
64 if keepalive is None:
65 event = await queue.get()
66 else:
67 event = await asyncio.wait_for(queue.get(), timeout=keepalive)
68 except TimeoutError:
69 yield ServerSentEvent(comment="keepalive")
70 continue
71 yield event
72 while not queue.empty():
73 yield queue.get_nowait()
74 finally:
75 drain_task.cancel()
76 with suppress(asyncio.CancelledError):
77 await drain_task
79 return EventSourceResponse(_items())