Coverage for src / lexigram / ai / relay / gateway / web / sse.py: 90%
31 statements
« prev ^ index » next coverage.py v7.13.5, created at 2026-08-08 23:08 +0800
« prev ^ index » next coverage.py v7.13.5, created at 2026-08-08 23:08 +0800
1"""SSE framing from normalized relay events into client wire protocols.
3``RelayWireEvent`` values are framed as Server-Sent Events following the
4client's inbound wire protocol: OpenAI Chat streams data-only frames with
5a ``[DONE]`` terminator, OpenAI Responses and Claude streams carry
6``event:`` names above the data line, and Gemini streams data-only frames
7without a terminator. JSON serialization goes through
8``lexigram.serialization`` only.
9"""
11from __future__ import annotations
13from lexigram.contracts.ai.relay import RelayFormat, RelayWireEvent
14from lexigram.serialization import dumps_str
16__all__ = ["SSEEncoder"]
19class SSEEncoder:
20 """Frame ``RelayWireEvent`` values as SSE for one client protocol.
22 Encoders are stateless; a single instance can frame an entire stream.
23 Each ``encode`` call returns the complete bytes for one event,
24 including any ``[DONE]`` terminator, so streams never need module or
25 instance-level buffering.
26 """
28 def __init__(self, source: RelayFormat) -> None:
29 """Bind the encoder to the client's wire protocol.
31 Args:
32 source: The client's relay format; frames follow its syntax.
33 """
34 self._source = source
36 def encode(self, event: RelayWireEvent) -> bytes:
37 """Encode one event as a complete SSE frame.
39 Args:
40 event: The normalized wire event to frame.
42 Returns:
43 The full SSE frame bytes for the event, including the
44 trailing blank line and any ``[DONE]`` terminator.
45 """
46 data = dumps_str(event.data).encode("utf-8") if event.data is not None else b""
47 if self._source == RelayFormat.OPENAI_CHAT:
48 if event.terminal and data:
49 return b"data: " + data + b"\n\ndata: [DONE]\n\n"
50 if event.terminal:
51 return b"data: [DONE]\n\n"
52 return b"data: " + data + b"\n\n"
53 if self._source in {RelayFormat.OPENAI_RESPONSES, RelayFormat.CLAUDE}:
54 name = self._event_name(event)
55 return b"event: " + name.encode("utf-8") + b"\ndata: " + data + b"\n\n"
56 return b"data: " + data + b"\n\n"
58 def encode_terminal(
59 self, source: RelayFormat, terminal_event: RelayWireEvent | None
60 ) -> bytes:
61 """Emit the closing frame when the client protocol requires one.
63 OpenAI Chat streams terminate with ``data: [DONE]``. When the
64 terminal event already passed through ``encode`` (which appends
65 ``[DONE]`` for terminal events) nothing more is emitted; a stream
66 that ends without a terminal event gets the ``[DONE]`` frame here
67 so the client always sees exactly one terminator. All other
68 formats terminate in-band and return nothing.
70 Args:
71 source: The client's relay format.
72 terminal_event: The framed terminal event, if any.
74 Returns:
75 The closing frame, or ``b""`` when none is needed.
76 """
77 if source != RelayFormat.OPENAI_CHAT:
78 return b""
79 if terminal_event is not None:
80 return b""
81 return b"data: [DONE]\n\n"
83 def _event_name(self, event: RelayWireEvent) -> str:
84 """Derive the SSE event name from the event or its data type."""
85 name = event.event
86 if name is None:
87 fallback = (
88 event.data.get("type", "message")
89 if event.data is not None
90 else "message"
91 )
92 name = fallback if isinstance(fallback, str) else str(fallback)
93 return name