From 3601840f2e6e59b352facccd0109d17e9871dc40 Mon Sep 17 00:00:00 2001 From: HughhhhCoder Date: Tue, 8 Sep 2026 14:31:33 +0800 Subject: [PATCH] fix(streaming): preserve first SSE event after BOM --- src/openai/_streaming.py | 6 ++++++ tests/test_sse_framing.py | 8 ++++++++ 2 files changed, 14 insertions(+) diff --git a/src/openai/_streaming.py b/src/openai/_streaming.py index 78e2d20aa7..3ed5f911f5 100644 --- a/src/openai/_streaming.py +++ b/src/openai/_streaming.py @@ -327,6 +327,7 @@ class SSEDecoder: _event: str | None _retry: int | None _last_event_id: str | None + _first_line: bool def __init__(self) -> None: self._reset() @@ -336,6 +337,7 @@ def _reset(self) -> None: self._data = [] self._last_event_id = None self._retry = None + self._first_line = True def _decode_lines(self, lines: Iterator[bytes]) -> Iterator[ServerSentEvent]: for raw_line in lines: @@ -368,6 +370,10 @@ async def aiter_bytes(self, iterator: AsyncIterator[bytes]) -> AsyncIterator[Ser def decode(self, line: str) -> ServerSentEvent | None: # See: https://html.spec.whatwg.org/multipage/server-sent-events.html#event-stream-interpretation # noqa: E501 + if self._first_line: + self._first_line = False + line = line.removeprefix("\ufeff") + if not line: if not self._event and not self._data and not self._last_event_id and self._retry is None: return None diff --git a/tests/test_sse_framing.py b/tests/test_sse_framing.py index 74f814dfd8..419d9b5b4e 100644 --- a/tests/test_sse_framing.py +++ b/tests/test_sse_framing.py @@ -45,6 +45,14 @@ async def test_every_chunk_boundary(sync: bool, ending: bytes) -> None: assert [event.data for event in events] == ["ok"] +@pytest.mark.parametrize("sync", [True, False]) +@pytest.mark.parametrize("fragment_size", [1, 2, 100000]) +async def test_leading_utf8_bom(sync: bool, fragment_size: int) -> None: + frame = b"\xef\xbb\xbfdata:first\n\ndata:second\n\n" + events = await _decode(SSEDecoder(), _fragment(frame, fragment_size), sync) + assert [event.data for event in events] == ["first", "second"] + + @pytest.mark.parametrize("sync", [True, False]) async def test_many_complete_events(sync: bool) -> None: events = await _decode(SSEDecoder(), [b"data:ok\n\n" * 2000], sync)