From 938396ef9099ebc5c0056169286325bf1d2bae0a Mon Sep 17 00:00:00 2001 From: mateo-berri <277851410+mateo-berri@users.noreply.github.com> Date: Mon, 10 Aug 2026 21:37:07 -0700 Subject: [PATCH] fix(proxy): recognize crlf sse frame boundaries in passthrough reassembly --- litellm/proxy/common_request_processing.py | 2 +- .../streaming_handler.py | 8 +++++--- .../test_streaming_handler_interrupt.py | 18 ++++++++++++++++++ 3 files changed, 24 insertions(+), 4 deletions(-) diff --git a/litellm/proxy/common_request_processing.py b/litellm/proxy/common_request_processing.py index 8dfa08ff19b..4e422ee49d7 100644 --- a/litellm/proxy/common_request_processing.py +++ b/litellm/proxy/common_request_processing.py @@ -3066,7 +3066,7 @@ class ProxyBaseLLMRequestProcessing: elif isinstance(chunk, (bytes, bytearray)): try: s: Final = chunk.decode("utf-8") - if s.endswith("\n\n"): + if s.endswith(("\n\n", "\r\n\r\n")): maybe_mod = ProxyBaseLLMRequestProcessing._inject_cost_into_sse_frame_str(s, model_name) if maybe_mod is not None: return maybe_mod.encode("utf-8") diff --git a/litellm/proxy/pass_through_endpoints/streaming_handler.py b/litellm/proxy/pass_through_endpoints/streaming_handler.py index da4a8eceb89..8428c7dcbe2 100644 --- a/litellm/proxy/pass_through_endpoints/streaming_handler.py +++ b/litellm/proxy/pass_through_endpoints/streaming_handler.py @@ -141,10 +141,12 @@ class PassThroughStreamingHandler: @staticmethod def _split_complete_sse_frames(pending: bytes) -> tuple[bytes, bytes]: - frame_boundary: Final = pending.rfind(b"\n\n") - if frame_boundary == -1: + lf_boundary_end: Final = pending.rfind(b"\n\n") + 2 + crlf_boundary_end: Final = pending.rfind(b"\r\n\r\n") + 4 + boundary_end: Final = max(lf_boundary_end if lf_boundary_end >= 2 else 0, crlf_boundary_end if crlf_boundary_end >= 4 else 0) + if boundary_end == 0: return b"", pending - return pending[: frame_boundary + 2], pending[frame_boundary + 2 :] + return pending[:boundary_end], pending[boundary_end:] @staticmethod async def _route_streaming_logging_to_handler( diff --git a/tests/test_litellm/proxy/pass_through_endpoints/test_streaming_handler_interrupt.py b/tests/test_litellm/proxy/pass_through_endpoints/test_streaming_handler_interrupt.py index 7b649db43ed..d559faba1c2 100644 --- a/tests/test_litellm/proxy/pass_through_endpoints/test_streaming_handler_interrupt.py +++ b/tests/test_litellm/proxy/pass_through_endpoints/test_streaming_handler_interrupt.py @@ -443,6 +443,24 @@ async def test_chunk_processor_injects_cost_into_usage_frame_fragmented_across_c assert reassembled.endswith("data: [DONE]\n\n") +@pytest.mark.asyncio +async def test_chunk_processor_streams_crlf_delimited_frames_live_and_injects_cost(monkeypatch): + """Regression: CRLF-delimited SSE frames must flow as they complete instead of + buffering until EOF, and the usage frame must still get cost injected.""" + monkeypatch.setattr(litellm, "include_cost_in_streaming_usage", True) + chunks = [chunk.replace(b"\n\n", b"\r\n\r\n") for chunk in _openai_passthrough_stream_chunks()] + + received = await _collect_openai_passthrough_chunks(chunks, EndpointType.OPENAI) + + assert len(received) == len(chunks) + assert received[0] == chunks[0] + reassembled = b"".join(received).decode("utf-8") + usage_lines = [ln for ln in reassembled.replace("\r\n", "\n").split("\n") if '"total_tokens"' in ln] + assert len(usage_lines) == 1 + final_payload = json.loads(usage_lines[0].split("data:", 1)[1].strip()) + assert final_payload["usage"]["cost"] > 0 + + @pytest.mark.asyncio async def test_chunk_processor_flag_off_leaves_openai_passthrough_stream_byte_identical(monkeypatch): monkeypatch.setattr(litellm, "include_cost_in_streaming_usage", False)