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

1"""Low-level Server-Sent Events transport. 

2 

3For higher-level SSE with backpressure, retry, and handler patterns, 

4use ``lexigram_web.sse`` instead. 

5""" 

6 

7from __future__ import annotations 

8 

9from collections.abc import AsyncGenerator 

10from typing import Any 

11 

12from starlette.responses import StreamingResponse 

13 

14from lexigram.serialization import dumps 

15 

16 

17class ServerSentEvent: 

18 """Server-sent event""" 

19 

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 

33 

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

46 

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

54 

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" 

59 

60 

61class EventSourceResponse(StreamingResponse): 

62 """Server-sent events response""" 

63 

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 ) 

78 

79 async def event_generator() -> AsyncGenerator[str, None]: 

80 async for event in content: 

81 yield event.encode() 

82 

83 super().__init__( 

84 content=event_generator(), 

85 status_code=status_code, 

86 headers=headers, 

87 media_type="text/event-stream", 

88 ) 

89 

90 

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)