Coverage for src/lexigram/web/transport/sse.py: 28%
40 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"""Low-level Server-Sent Events transport.
3For higher-level SSE with backpressure, retry, and handler patterns,
4use ``lexigram_web.sse`` instead.
5"""
7from __future__ import annotations
9from collections.abc import AsyncGenerator
10from typing import Any
12from starlette.responses import StreamingResponse
14from lexigram.serialization import dumps
17class ServerSentEvent:
18 """Server-sent event"""
20 def __init__(
21 self,
22 data: Any = None,
23 event: str | None = None,
24 event_id: str | None = None,
25 retry: int | None = None,
26 comment: str | None = None,
27 ) -> None:
28 self.data = data
29 self.event = event
30 self.event_id = event_id
31 self.retry = retry
32 self.comment = comment
34 def encode(self) -> str:
35 """Encode the event as SSE format"""
36 # Build header lines conditionally
37 lines = []
38 if self.comment:
39 lines.append(f": {self.comment}")
40 if self.event:
41 lines.append(f"event: {self.event}")
42 if self.event_id:
43 lines.append(f"id: {self.event_id}")
44 if self.retry:
45 lines.append(f"retry: {self.retry}")
47 # Handle data - split string by newlines or serialize object
48 if self.data is None:
49 data_lines = []
50 elif isinstance(self.data, str):
51 data_lines = self.data.split("\n")
52 else:
53 data_lines = [dumps(self.data).decode("utf-8")]
55 # Append data lines (using list comprehension)
56 lines.extend(f"data: {line}" for line in data_lines)
57 lines.append("") # Empty line to end the event
58 return "\n".join(lines) + "\n"
61class EventSourceResponse(StreamingResponse):
62 """Server-sent events response"""
64 def __init__(
65 self,
66 content: AsyncGenerator[ServerSentEvent, None],
67 status_code: int = 200,
68 headers: dict[str, str] | None = None,
69 ) -> None:
70 headers = headers or {}
71 headers.update(
72 {
73 "Content-Type": "text/event-stream",
74 "Cache-Control": "no-cache",
75 "Connection": "keep-alive",
76 },
77 )
79 async def event_generator() -> AsyncGenerator[str, None]:
80 async for event in content:
81 yield event.encode()
83 super().__init__(
84 content=event_generator(),
85 status_code=status_code,
86 headers=headers,
87 media_type="text/event-stream",
88 )
91def sse_response(
92 content: AsyncGenerator[ServerSentEvent, None],
93 status_code: int = 200,
94 headers: dict[str, str] | None = None,
95) -> EventSourceResponse:
96 """Create a server-sent events response"""
97 return EventSourceResponse(content, status_code, headers)