Skip to content

Commit 5e7a2bd

Browse files
committed
fix: preserve oversized public responses SSE events
1 parent c840145 commit 5e7a2bd

3 files changed

Lines changed: 101 additions & 24 deletions

File tree

lib/codex_pooler/gateway/transports/streaming/stream_protocol.ex

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -37,7 +37,8 @@ defmodule CodexPooler.Gateway.Transports.Streaming.StreamProtocol do
3737
@type public_openai_responses_stream_state :: %{
3838
required(:buffer) => binary(),
3939
required(:created?) => boolean(),
40-
required(:text_delta?) => boolean()
40+
required(:text_delta?) => boolean(),
41+
required(:passthrough?) => boolean()
4142
}
4243
@type websocket_frame_headers :: %{optional(String.t()) => String.t()}
4344

lib/codex_pooler/gateway/transports/streaming/stream_protocol/public_responses.ex

Lines changed: 62 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -7,37 +7,71 @@ defmodule CodexPooler.Gateway.Transports.Streaming.StreamProtocol.PublicResponse
77
@type state :: %{
88
required(:buffer) => binary(),
99
required(:created?) => boolean(),
10-
required(:text_delta?) => boolean()
10+
required(:text_delta?) => boolean(),
11+
required(:passthrough?) => boolean()
1112
}
1213

1314
@spec new_state() :: state()
14-
def new_state, do: %{buffer: "", created?: false, text_delta?: false}
15+
def new_state, do: %{buffer: "", created?: false, text_delta?: false, passthrough?: false}
1516

1617
@spec normalize_data(binary(), state()) :: {binary(), state()}
18+
def normalize_data(data, %{passthrough?: true} = state) when is_binary(data) do
19+
normalize_passthrough_data(data, state)
20+
end
21+
1722
def normalize_data(data, state) when is_binary(data) do
1823
buffered_data = state.buffer <> data
19-
{blocks, buffer} = StreamProtocol.complete_sse_blocks(buffered_data, bounded?: true)
20-
21-
if oversized_incomplete_sse_prefix?(blocks, buffer, buffered_data) do
22-
BufferTelemetry.record_oversized_incomplete(
23-
"public_openai_responses_sse",
24-
byte_size(buffered_data),
25-
StreamProtocol.max_incomplete_sse_block_bytes()
26-
)
27-
28-
{buffered_data, %{state | buffer: ""}}
29-
else
30-
normalize_blocks(blocks, buffer, state)
24+
{blocks, buffer} = StreamProtocol.complete_sse_blocks(buffered_data, bounded?: false)
25+
26+
cond do
27+
blocks == [] and StreamProtocol.oversized_incomplete_sse_block?(buffer) ->
28+
record_oversized_incomplete(byte_size(buffered_data))
29+
{buffered_data, %{state | buffer: "", passthrough?: true}}
30+
31+
StreamProtocol.oversized_incomplete_sse_block?(buffer) ->
32+
record_oversized_incomplete(byte_size(buffered_data))
33+
{iodata, state} = normalize_complete_blocks(blocks, state)
34+
{[iodata, buffer] |> IO.iodata_to_binary(), %{state | buffer: "", passthrough?: true}}
35+
36+
true ->
37+
normalize_blocks(blocks, buffer, state)
3138
end
3239
end
3340

3441
def normalize_data(data, state), do: {data, state}
3542

43+
defp normalize_passthrough_data(data, state) do
44+
case sse_block_separator(data) do
45+
{index, separator_size} ->
46+
passthrough_size = index + separator_size
47+
<<passthrough::binary-size(passthrough_size), rest::binary>> = data
48+
49+
state = %{state | passthrough?: false, buffer: ""}
50+
{normalized_rest, state} = normalize_data(rest, state)
51+
52+
{[passthrough, normalized_rest] |> IO.iodata_to_binary(), state}
53+
54+
nil ->
55+
{data, state}
56+
end
57+
end
58+
59+
defp record_oversized_incomplete(bytes) do
60+
BufferTelemetry.record_oversized_incomplete(
61+
"public_openai_responses_sse",
62+
bytes,
63+
StreamProtocol.max_incomplete_sse_block_bytes()
64+
)
65+
end
66+
67+
defp normalize_complete_blocks(blocks, state) do
68+
Enum.map_reduce(blocks, state, fn block, stream_state ->
69+
normalize_block(block, stream_state)
70+
end)
71+
end
72+
3673
defp normalize_blocks(blocks, buffer, state) do
37-
{iodata, state} =
38-
Enum.map_reduce(blocks, %{state | buffer: buffer}, fn block, stream_state ->
39-
normalize_block(block, stream_state)
40-
end)
74+
{iodata, state} = normalize_complete_blocks(blocks, %{state | buffer: buffer})
4175

4276
state = if stream_terminal?(blocks), do: new_state(), else: state
4377

@@ -72,11 +106,6 @@ defmodule CodexPooler.Gateway.Transports.Streaming.StreamProtocol.PublicResponse
72106
end
73107
end
74108

75-
defp oversized_incomplete_sse_prefix?([], "", data),
76-
do: StreamProtocol.oversized_incomplete_sse_block?(data)
77-
78-
defp oversized_incomplete_sse_prefix?(_blocks, _buffer, _data), do: false
79-
80109
defp terminal_prefix(decoded, state) do
81110
{created_prefix, state} =
82111
if state.created? do
@@ -161,6 +190,16 @@ defmodule CodexPooler.Gateway.Transports.Streaming.StreamProtocol.PublicResponse
161190
defp codex_public_event?(type) when is_binary(type), do: String.starts_with?(type, "codex.")
162191
defp codex_public_event?(_type), do: false
163192

193+
defp sse_block_separator(data) do
194+
["\n\n", "\r\n\r\n"]
195+
|> Enum.map(fn separator -> {separator, :binary.match(data, separator)} end)
196+
|> Enum.flat_map(fn
197+
{separator, {index, _size}} -> [{index, byte_size(separator)}]
198+
{_separator, :nomatch} -> []
199+
end)
200+
|> Enum.min_by(fn {index, _size} -> index end, fn -> nil end)
201+
end
202+
164203
defp stream_block_event(block) do
165204
data = StreamProtocol.sse_field(block, "data")
166205

test/codex_pooler/gateway/transports/stream_protocol_test.exs

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,43 @@ defmodule CodexPooler.Gateway.Transports.Streaming.StreamProtocolTest do
1212
end
1313

1414
describe "normalize_public_openai_responses_sse_data/2" do
15+
test "preserves oversized split reasoning events until the SSE block is complete" do
16+
state = StreamProtocol.public_openai_responses_stream_state()
17+
18+
event =
19+
sse_event("response.output_item.added", %{
20+
"type" => "response.output_item.added",
21+
"output_index" => 0,
22+
"sequence_number" => 2,
23+
"item" => %{
24+
"id" => "rs_oversized_reasoning",
25+
"type" => "reasoning",
26+
"summary" => [],
27+
"encrypted_content" => String.duplicate("synthetic-obfuscated-content", 4_000)
28+
}
29+
})
30+
31+
split_at = StreamProtocol.max_incomplete_sse_block_bytes() + 1
32+
<<first::binary-size(split_at), second::binary>> = event
33+
34+
{first_out, state} =
35+
StreamProtocol.normalize_public_openai_responses_sse_data(first, state)
36+
37+
{second_out, _state} =
38+
StreamProtocol.normalize_public_openai_responses_sse_data(second, state)
39+
40+
combined = first_out <> second_out
41+
42+
assert combined == event
43+
assert [block] = StreamProtocol.complete_sse_blocks(combined, bounded?: false) |> elem(0)
44+
assert "response.output_item.added" == StreamProtocol.sse_field(block, "event")
45+
46+
assert %{"item" => %{"type" => "reasoning"}} =
47+
block
48+
|> StreamProtocol.sse_field("data")
49+
|> StreamProtocol.decode_sse_data()
50+
end
51+
1552
test "carries incomplete response stream state explicitly" do
1653
state = StreamProtocol.public_openai_responses_stream_state()
1754

0 commit comments

Comments
 (0)