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

1"""SSE responses driven by lexigram.reactive EventStream sources.""" 

2 

3from __future__ import annotations 

4 

5import asyncio 

6from collections.abc import AsyncGenerator, Callable 

7from contextlib import suppress 

8from typing import Any 

9 

10from lexigram.reactive import EventStream 

11from lexigram.web.transport.sse import EventSourceResponse, ServerSentEvent 

12 

13 

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. 

21 

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. 

29 

30 Returns: 

31 A Starlette StreamingResponse with ``text/event-stream``. 

32 

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 """ 

43 

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() 

48 

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() 

56 

57 drain_task = asyncio.create_task(_drain()) 

58 background_tasks.add(drain_task) 

59 drain_task.add_done_callback(background_tasks.discard) 

60 

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 

78 

79 return EventSourceResponse(_items())