From feeea11a356c2d4ead825bf549732abb3b457ab8 Mon Sep 17 00:00:00 2001 From: Shakti Prasad Mohapatra Date: Mon, 20 Jul 2026 23:39:41 +0530 Subject: [PATCH 1/5] fix: merge list entries by logical index instead of physical position (#3201) accumulate_delta assumed that an indexed list entry's 'index' field matches its physical position in the Python list. This breaks when the first streamed chunk contains multiple tool_calls entries with the same index (e.g. from speculative decoding with vLLM). The first chunk is stored directly (acc[key] = delta_value), creating two physical entries with index: 0. Later chunks merge into acc_value[0] by physical position, stranding the second duplicate and producing invalid final JSON. Fix: when accumulating list entries, search for an existing entry by its 'index' field and merge into it, rather than indexing by physical position. New entries are inserted at their logical index. --- src/openai/lib/streaming/_deltas.py | 24 +++++--- tests/lib/streaming/test_deltas.py | 88 +++++++++++++++++++++++++++++ 2 files changed, 104 insertions(+), 8 deletions(-) create mode 100644 tests/lib/streaming/test_deltas.py diff --git a/src/openai/lib/streaming/_deltas.py b/src/openai/lib/streaming/_deltas.py index a5e1317612..9c305f32b9 100644 --- a/src/openai/lib/streaming/_deltas.py +++ b/src/openai/lib/streaming/_deltas.py @@ -49,15 +49,23 @@ def accumulate_delta(acc: dict[object, object], delta: dict[object, object]) -> if not isinstance(index, int): raise TypeError(f"Unexpected, list delta entry `index` value is not an integer; {index}") - try: - acc_entry = acc_value[index] - except IndexError: - acc_value.insert(index, delta_entry) - else: - if not is_dict(acc_entry): - raise TypeError("not handled yet") + # Merge by logical index, not physical position. (#3201) + # When the first chunk contains multiple entries with the same + # index (e.g. from speculative decoding), the physical position + # does not match the logical index. Find the existing entry by + # its index field and merge into it. + found = False + for i, existing in enumerate(acc_value): + if is_dict(existing) and existing.get("index") == index: + acc_value[i] = accumulate_delta(existing, delta_entry) + found = True + break - acc_value[index] = accumulate_delta(acc_entry, delta_entry) + if not found: + # Ensure the list is large enough + while len(acc_value) <= index: + acc_value.append({}) + acc_value[index] = delta_entry acc[key] = acc_value diff --git a/tests/lib/streaming/test_deltas.py b/tests/lib/streaming/test_deltas.py new file mode 100644 index 0000000000..00f346c035 --- /dev/null +++ b/tests/lib/streaming/test_deltas.py @@ -0,0 +1,88 @@ +"""Tests for the streaming delta accumulator.""" + +from __future__ import annotations + +from openai.lib.streaming._deltas import accumulate_delta + + +class TestAccumulateDelta: + """Tests for accumulate_delta — regression for #3201.""" + + def test_duplicate_index_first_chunk_merges(self) -> None: + """First chunk with two entries at the same index should merge into one.""" + acc: dict[object, object] = {} + delta = { + "tool_calls": [ + { + "index": 0, + "id": "call_abc", + "function": {"name": "list_files"}, + "type": "function", + }, + { + "index": 0, + "function": {"arguments": ' {"'}, + }, + ] + } + result = accumulate_delta(acc, delta) + calls = result["tool_calls"] + assert isinstance(calls, list) + # Should be a single entry at index 0, not two + assert len(calls) == 1 + assert calls[0]["index"] == 0 + assert calls[0]["id"] == "call_abc" + assert calls[0]["function"]["name"] == "list_files" + assert calls[0]["function"]["arguments"] == ' {"' + + def test_duplicate_index_subsequent_chunk_merges(self) -> None: + """Subsequent chunk with same index should merge into existing entry.""" + acc: dict[object, object] = { + "tool_calls": [ + { + "index": 0, + "id": "call_abc", + "function": {"name": "list_files", "arguments": ' {"'}, + "type": "function", + } + ] + } + delta = { + "tool_calls": [ + { + "index": 0, + "function": {"arguments": 'path": "."}'}, + } + ] + } + result = accumulate_delta(acc, delta) + calls = result["tool_calls"] + assert len(calls) == 1 + assert calls[0]["function"]["arguments"] == ' {"path": "."}' + + def test_different_indexes_accumulate_separately(self) -> None: + """Entries with different indexes should accumulate separately.""" + acc: dict[object, object] = {} + delta1 = { + "tool_calls": [ + {"index": 0, "id": "call_a", "function": {"name": "tool_a"}, "type": "function"}, + ] + } + delta2 = { + "tool_calls": [ + {"index": 1, "id": "call_b", "function": {"name": "tool_b"}, "type": "function"}, + ] + } + result = accumulate_delta(acc, delta1) + result = accumulate_delta(result, delta2) + calls = result["tool_calls"] + assert len(calls) == 2 + assert calls[0]["index"] == 0 + assert calls[1]["index"] == 1 + + def test_string_accumulation_unchanged(self) -> None: + """Basic string accumulation should still work.""" + acc: dict[object, object] = {"content": "hello"} + delta = {"content": " world"} + result = accumulate_delta(acc, delta) + assert result["content"] == "hello world" \ No newline at end of file From f2b61eb836851b242f6e0e81222f6b1c785b4441 Mon Sep 17 00:00:00 2001 From: Shakti Prasad Mohapatra Date: Mon, 20 Jul 2026 23:54:57 +0530 Subject: [PATCH 2/5] fix: coalesce duplicate-index entries when first chunk is stored directly (#3201) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Codex P1 review: when tool_calls is first added to the snapshot it is copied directly (acc[key] = delta_value), so a first chunk with two entries at the same index creates two physical entries that later merges can't fix — the merge only hits the first matching entry and breaks, leaving the earlier duplicate stranded. Fix: coalesce duplicate-index entries in the first chunk before storing it, using accumulate_delta to merge entries with the same index field. Added test_duplicate_index_first_chunk_then_subsequent_merge to verify the full round-trip: first chunk coalesces, subsequent chunk merges into the single coalesced entry. --- src/openai/lib/streaming/_deltas.py | 36 +++++++++++++++++++++++++++++ tests/lib/streaming/test_deltas.py | 30 +++++++++++++++++++++++- 2 files changed, 65 insertions(+), 1 deletion(-) diff --git a/src/openai/lib/streaming/_deltas.py b/src/openai/lib/streaming/_deltas.py index 9c305f32b9..dc23b85faf 100644 --- a/src/openai/lib/streaming/_deltas.py +++ b/src/openai/lib/streaming/_deltas.py @@ -6,6 +6,12 @@ def accumulate_delta(acc: dict[object, object], delta: dict[object, object]) -> dict[object, object]: for key, delta_value in delta.items(): if key not in acc: + # When the first chunk contains a list with multiple entries at the + # same index (e.g. from speculative decoding), storing it directly + # would leave duplicate entries that later merges can't fix. (#3201) + # Coalesce duplicate-index entries before storing. + if is_list(delta_value) and len(delta_value) > 1: + delta_value = _coalesce_list_by_index(delta_value) acc[key] = delta_value continue @@ -70,3 +76,33 @@ def accumulate_delta(acc: dict[object, object], delta: dict[object, object]) -> acc[key] = acc_value return acc + + +def _coalesce_list_by_index(lst: list[object]) -> list[object]: + """Merge list entries that share the same ``index`` field into a single entry. + + When the first streamed chunk contains multiple entries with the same + ``index`` (e.g. from speculative decoding), storing the list directly would + leave duplicate entries. This function coalesces them by merging entries + with the same index using :func:`accumulate_delta`, so the snapshot starts + in a clean state. (#3201) + """ + result: list[object] = [] + for entry in lst: + if not is_dict(entry): + result.append(entry) + continue + index = entry.get("index") + if not isinstance(index, int): + result.append(entry) + continue + # Find an existing entry with the same index + found = False + for i, existing in enumerate(result): + if is_dict(existing) and existing.get("index") == index: + result[i] = accumulate_delta(existing, entry) + found = True + break + if not found: + result.append(entry) + return result diff --git a/tests/lib/streaming/test_deltas.py b/tests/lib/streaming/test_deltas.py index 00f346c035..1d613a374f 100644 --- a/tests/lib/streaming/test_deltas.py +++ b/tests/lib/streaming/test_deltas.py @@ -85,4 +85,32 @@ def test_string_accumulation_unchanged(self) -> None: acc: dict[object, object] = {"content": "hello"} delta = {"content": " world"} result = accumulate_delta(acc, delta) - assert result["content"] == "hello world" \ No newline at end of file + assert result["content"] == "hello world" + + def test_duplicate_index_first_chunk_then_subsequent_merge(self) -> None: + """Full round-trip: first chunk with duplicate indexes, then subsequent chunk merges correctly.""" + acc: dict[object, object] = {} + # First chunk: two entries at index 0 + delta1 = { + "tool_calls": [ + {"index": 0, "id": "call_abc", "function": {"name": "list_files"}, "type": "function"}, + {"index": 0, "function": {"arguments": ' {"'}}, + ] + } + result = accumulate_delta(acc, delta1) + calls = result["tool_calls"] + assert len(calls) == 1, f"Expected 1 entry after coalescing, got {len(calls)}" + assert calls[0]["function"]["arguments"] == ' {"' + + # Second chunk: more arguments for index 0 + delta2 = { + "tool_calls": [ + {"index": 0, "function": {"arguments": 'path": "."}'}}, + ] + } + result = accumulate_delta(result, delta2) + calls = result["tool_calls"] + assert len(calls) == 1 + assert calls[0]["function"]["arguments"] == ' {"path": "."}' + assert calls[0]["id"] == "call_abc" + assert calls[0]["function"]["name"] == "list_files" \ No newline at end of file From 1b63e4cd0f0bb605a58ea78fda9d266197c8461a Mon Sep 17 00:00:00 2001 From: Shakti Prasad Mohapatra Date: Thu, 23 Jul 2026 23:23:31 +0530 Subject: [PATCH 3/5] fix: merge delta into all matching index entries, not just the first When acc_value already contains duplicate-index entries (e.g. from a prior chunk that wasn't coalesced), the merge loop only merged into the first matching entry and broke, leaving the second duplicate stranded. Remove the break so all matching entries get the delta merged in. Addresses Codex P1 review feedback. --- src/openai/lib/streaming/_deltas.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/src/openai/lib/streaming/_deltas.py b/src/openai/lib/streaming/_deltas.py index dc23b85faf..e37d1101fa 100644 --- a/src/openai/lib/streaming/_deltas.py +++ b/src/openai/lib/streaming/_deltas.py @@ -60,12 +60,15 @@ def accumulate_delta(acc: dict[object, object], delta: dict[object, object]) -> # index (e.g. from speculative decoding), the physical position # does not match the logical index. Find the existing entry by # its index field and merge into it. + # + # If acc_value already contains duplicate-index entries + # (e.g. from a prior chunk that wasn't coalesced), merge into + # all of them so none are stranded. found = False for i, existing in enumerate(acc_value): if is_dict(existing) and existing.get("index") == index: acc_value[i] = accumulate_delta(existing, delta_entry) found = True - break if not found: # Ensure the list is large enough From 78360f905023a52b00d88c6a3121777292567a7c Mon Sep 17 00:00:00 2001 From: Shakti Prasad Mohapatra Date: Fri, 7 Aug 2026 12:30:08 +0530 Subject: [PATCH 4/5] fix: normalize first-chunk duplicates, fix sparse-index data loss, fix types MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Addresses review feedback from @jbeckwith-oai on #3521: 1. First-chunk duplicate indexes — _convert_initial_chunk_into_snapshot now applies _coalesce_list_by_index to tool_calls in the initial chunk, so duplicate-index entries from speculative decoding are merged before the snapshot is seeded. 2. Data-loss in not-found branch — when acc_value has entries at higher indexes (e.g. [{"index": 1}]) and a lower index arrives later, the old code would overwrite the existing entry. Now appends instead of assigning by position. Added test_sparse_out_of_order_indexes_no_data_loss. 3. Test file Pyright errors — all 20 errors fixed with proper type annotations and cast calls. Pyright and Ruff both pass clean. --- src/openai/lib/streaming/_deltas.py | 43 ++++++- src/openai/lib/streaming/chat/_completions.py | 22 +++- tests/lib/streaming/test_deltas.py | 119 ++++++++++++++++-- 3 files changed, 162 insertions(+), 22 deletions(-) diff --git a/src/openai/lib/streaming/_deltas.py b/src/openai/lib/streaming/_deltas.py index e37d1101fa..3201cc992f 100644 --- a/src/openai/lib/streaming/_deltas.py +++ b/src/openai/lib/streaming/_deltas.py @@ -71,10 +71,29 @@ def accumulate_delta(acc: dict[object, object], delta: dict[object, object]) -> found = True if not found: - # Ensure the list is large enough - while len(acc_value) <= index: - acc_value.append({}) - acc_value[index] = delta_entry + # Add the new entry. Don't assume the logical index is a + # safe physical slot — if acc_value already has entries at + # higher indexes (e.g. [{"index": 1, ...}] and index 0 + # arrives), acc_value[index] would overwrite the existing + # entry. Place the entry at the position matching the + # logical index so downstream code that does + # tool_calls[index] (treating logical index as physical + # position) reads the right entry. + if len(acc_value) <= index: + while len(acc_value) < index: + acc_value.append({}) + acc_value.append(delta_entry) + else: + # The list is large enough but no entry has this + # index. If the slot at `index` is an empty + # placeholder ({}), replace it in-place. Otherwise + # insert at the correct position to keep the list + # addressable by logical index. + existing = acc_value[index] + if isinstance(existing, dict) and not existing: + acc_value[index] = delta_entry + else: + acc_value.insert(index, delta_entry) acc[key] = acc_value @@ -89,6 +108,10 @@ def _coalesce_list_by_index(lst: list[object]) -> list[object]: leave duplicate entries. This function coalesces them by merging entries with the same index using :func:`accumulate_delta`, so the snapshot starts in a clean state. (#3201) + + The result is sorted by the ``index`` field so the list stays addressable + by logical index — downstream code does ``tool_calls[index]`` treating + logical index as physical position. """ result: list[object] = [] for entry in lst: @@ -107,5 +130,15 @@ def _coalesce_list_by_index(lst: list[object]) -> list[object]: found = True break if not found: - result.append(entry) + # Place at the position matching the logical index, padding + # with empty dicts if needed, so the list is addressable by + # logical index. + while len(result) <= index: + result.append({}) + # Replace the placeholder at `index` or shift if occupied + existing = result[index] + if isinstance(existing, dict) and not existing: + result[index] = entry + else: + result.insert(index, entry) return result diff --git a/src/openai/lib/streaming/chat/_completions.py b/src/openai/lib/streaming/chat/_completions.py index 5f072cafbd..6d391d5378 100644 --- a/src/openai/lib/streaming/chat/_completions.py +++ b/src/openai/lib/streaming/chat/_completions.py @@ -22,9 +22,9 @@ FunctionToolCallArgumentsDoneEvent, FunctionToolCallArgumentsDeltaEvent, ) -from .._deltas import accumulate_delta +from .._deltas import accumulate_delta, _coalesce_list_by_index from ...._types import Omit, IncEx, omit -from ...._utils import is_given, consume_sync_iterator, consume_async_iterator +from ...._utils import is_list, is_given, consume_sync_iterator, consume_async_iterator from ...._compat import model_dump from ...._models import build, construct_type from ..._parsing import ( @@ -409,13 +409,19 @@ def _accumulate_chunk(self, chunk: ChatCompletionChunk) -> ParsedChatCompletionS elif TYPE_CHECKING: # type: ignore[unreachable] assert_never(prev_tool) except IndexError: + # A new choice appeared that wasn't in the initial chunk. + # Coalesce tool_calls by index to handle duplicate-index entries + # from speculative decoding, same as _convert_initial_chunk_into_snapshot. + delta_dict = choice.delta.to_dict() + if is_list(delta_dict.get("tool_calls")): + delta_dict["tool_calls"] = _coalesce_list_by_index(cast("list[object]", delta_dict["tool_calls"])) choice_snapshot = cast( ParsedChoiceSnapshot, construct_type( type_=ParsedChoiceSnapshot, value={ **choice.model_dump(exclude_unset=True, exclude={"delta"}), - "message": choice.delta.to_dict(), + "message": delta_dict, }, ), ) @@ -742,9 +748,17 @@ def _convert_initial_chunk_into_snapshot(chunk: ChatCompletionChunk) -> ParsedCh choices = cast("list[object]", data["choices"]) for choice in chunk.choices: + message_dict = choice.delta.to_dict() + # Coalesce duplicate-index tool_calls in the initial chunk. (#3201) + # When the first chunk contains multiple tool_calls with the same index + # (e.g. from speculative decoding), storing them directly would leave + # duplicate entries that later merges can't fix. + tool_calls = message_dict.get("tool_calls") + if is_list(tool_calls) and len(tool_calls) > 1: + message_dict["tool_calls"] = _coalesce_list_by_index(tool_calls) choices[choice.index] = { **choice.model_dump(exclude_unset=True, exclude={"delta"}), - "message": choice.delta.to_dict(), + "message": message_dict, } return cast( diff --git a/tests/lib/streaming/test_deltas.py b/tests/lib/streaming/test_deltas.py index 1d613a374f..aa1594d0ce 100644 --- a/tests/lib/streaming/test_deltas.py +++ b/tests/lib/streaming/test_deltas.py @@ -2,6 +2,8 @@ from __future__ import annotations +from typing import Any, cast + from openai.lib.streaming._deltas import accumulate_delta @@ -11,7 +13,7 @@ class TestAccumulateDelta: def test_duplicate_index_first_chunk_merges(self) -> None: """First chunk with two entries at the same index should merge into one.""" acc: dict[object, object] = {} - delta = { + delta: dict[object, object] = { "tool_calls": [ { "index": 0, @@ -26,7 +28,7 @@ def test_duplicate_index_first_chunk_merges(self) -> None: ] } result = accumulate_delta(acc, delta) - calls = result["tool_calls"] + calls = cast(list[dict[str, Any]], result["tool_calls"]) assert isinstance(calls, list) # Should be a single entry at index 0, not two assert len(calls) == 1 @@ -47,7 +49,7 @@ def test_duplicate_index_subsequent_chunk_merges(self) -> None: } ] } - delta = { + delta: dict[object, object] = { "tool_calls": [ { "index": 0, @@ -56,26 +58,26 @@ def test_duplicate_index_subsequent_chunk_merges(self) -> None: ] } result = accumulate_delta(acc, delta) - calls = result["tool_calls"] + calls = cast(list[dict[str, Any]], result["tool_calls"]) assert len(calls) == 1 assert calls[0]["function"]["arguments"] == ' {"path": "."}' def test_different_indexes_accumulate_separately(self) -> None: """Entries with different indexes should accumulate separately.""" acc: dict[object, object] = {} - delta1 = { + delta1: dict[object, object] = { "tool_calls": [ {"index": 0, "id": "call_a", "function": {"name": "tool_a"}, "type": "function"}, ] } - delta2 = { + delta2: dict[object, object] = { "tool_calls": [ {"index": 1, "id": "call_b", "function": {"name": "tool_b"}, "type": "function"}, ] } result = accumulate_delta(acc, delta1) result = accumulate_delta(result, delta2) - calls = result["tool_calls"] + calls = cast(list[dict[str, Any]], result["tool_calls"]) assert len(calls) == 2 assert calls[0]["index"] == 0 assert calls[1]["index"] == 1 @@ -83,7 +85,7 @@ def test_different_indexes_accumulate_separately(self) -> None: def test_string_accumulation_unchanged(self) -> None: """Basic string accumulation should still work.""" acc: dict[object, object] = {"content": "hello"} - delta = {"content": " world"} + delta: dict[object, object] = {"content": " world"} result = accumulate_delta(acc, delta) assert result["content"] == "hello world" @@ -91,26 +93,117 @@ def test_duplicate_index_first_chunk_then_subsequent_merge(self) -> None: """Full round-trip: first chunk with duplicate indexes, then subsequent chunk merges correctly.""" acc: dict[object, object] = {} # First chunk: two entries at index 0 - delta1 = { + delta1: dict[object, object] = { "tool_calls": [ {"index": 0, "id": "call_abc", "function": {"name": "list_files"}, "type": "function"}, {"index": 0, "function": {"arguments": ' {"'}}, ] } result = accumulate_delta(acc, delta1) - calls = result["tool_calls"] + calls = cast(list[dict[str, Any]], result["tool_calls"]) assert len(calls) == 1, f"Expected 1 entry after coalescing, got {len(calls)}" assert calls[0]["function"]["arguments"] == ' {"' # Second chunk: more arguments for index 0 - delta2 = { + delta2: dict[object, object] = { "tool_calls": [ {"index": 0, "function": {"arguments": 'path": "."}'}}, ] } result = accumulate_delta(result, delta2) - calls = result["tool_calls"] + calls = cast(list[dict[str, Any]], result["tool_calls"]) assert len(calls) == 1 assert calls[0]["function"]["arguments"] == ' {"path": "."}' assert calls[0]["id"] == "call_abc" - assert calls[0]["function"]["name"] == "list_files" \ No newline at end of file + assert calls[0]["function"]["name"] == "list_files" + + def test_sparse_out_of_order_indexes_no_data_loss(self) -> None: + """Regression for the data-loss bug: if acc_value has [{"index": 1, ...}] + and index 0 arrives later, the index-1 entry must not be overwritten.""" + acc: dict[object, object] = { + "tool_calls": [ + {"index": 1, "id": "call_b", "function": {"name": "tool_b"}, "type": "function"}, + ] + } + delta: dict[object, object] = { + "tool_calls": [ + {"index": 0, "id": "call_a", "function": {"name": "tool_a"}, "type": "function"}, + ] + } + result = accumulate_delta(acc, delta) + calls = cast(list[dict[str, Any]], result["tool_calls"]) + # Both entries should survive + assert len(calls) == 2 + # The index-1 entry should not be overwritten + ids = [c["id"] for c in calls] + assert "call_a" in ids + assert "call_b" in ids + + def test_out_of_order_index_stays_addressable_by_logical_index(self) -> None: + """Regression for Codex P2: when index 1 arrives before index 0, the + list must stay addressable by logical index — downstream code does + ``tool_calls[tool_call_delta.index]`` treating logical index as + physical position. If the list is ``[{"index": 1}, {"index": 0}]`` + then ``tool_calls[0]`` returns the wrong entry.""" + acc: dict[object, object] = { + "tool_calls": [ + {"index": 1, "id": "call_b", "function": {"name": "tool_b"}, "type": "function"}, + ] + } + delta: dict[object, object] = { + "tool_calls": [ + {"index": 0, "id": "call_a", "function": {"name": "tool_a"}, "type": "function"}, + ] + } + result = accumulate_delta(acc, delta) + calls = cast(list[dict[str, Any]], result["tool_calls"]) + # The list must be addressable by logical index: calls[0] should have + # index 0, calls[1] should have index 1. + assert calls[0]["index"] == 0 + assert calls[0]["id"] == "call_a" + assert calls[1]["index"] == 1 + assert calls[1]["id"] == "call_b" + + def test_gap_placeholder_replaced_not_shifted(self) -> None: + """Regression for Codex P2: when indexes 0 then 2 arrive, slot 1 is + padded with {}. If index 1 arrives later, it must replace the + placeholder in-place, not insert before it (which would shift the + placeholder ahead of index 2, breaking tool_calls[2] lookups).""" + acc: dict[object, object] = { + "tool_calls": [ + {"index": 0, "id": "call_a", "function": {"name": "tool_a"}, "type": "function"}, + {}, + {"index": 2, "id": "call_c", "function": {"name": "tool_c"}, "type": "function"}, + ] + } + delta: dict[object, object] = { + "tool_calls": [ + {"index": 1, "id": "call_b", "function": {"name": "tool_b"}, "type": "function"}, + ] + } + result = accumulate_delta(acc, delta) + calls = cast(list[dict[str, Any]], result["tool_calls"]) + # The placeholder at index 1 should be replaced, not shifted + assert len(calls) == 3 + assert calls[0]["index"] == 0 + assert calls[0]["id"] == "call_a" + assert calls[1]["index"] == 1 + assert calls[1]["id"] == "call_b" + assert calls[2]["index"] == 2 + assert calls[2]["id"] == "call_c" + + def test_coalesce_list_by_index_sorts_by_logical_index(self) -> None: + """Regression for Codex P2: _coalesce_list_by_index must sort entries + by logical index so the list is addressable by tool_calls[index].""" + from openai.lib.streaming._deltas import _coalesce_list_by_index + + lst: list[object] = [ + {"index": 1, "id": "call_b", "function": {"name": "tool_b"}, "type": "function"}, + {"index": 0, "id": "call_a", "function": {"name": "tool_a"}, "type": "function"}, + ] + result = _coalesce_list_by_index(lst) + calls = cast(list[dict[str, Any]], result) + assert calls[0]["index"] == 0 + assert calls[0]["id"] == "call_a" + assert calls[1]["index"] == 1 + assert calls[1]["id"] == "call_b" From cbcc4273b71a8de20bb36759258ef2b62335bc4b Mon Sep 17 00:00:00 2001 From: Shakti Prasad Mohapatra Date: Fri, 7 Aug 2026 17:54:17 +0530 Subject: [PATCH 5/5] fix: detect dumped placeholders after model_dump round-trip After the snapshot is round-tripped through model_dump, a gap-filler {} placeholder becomes a dict of unset tool-call fields (e.g. {"id": None, "function": None, "type": None}). The previous check (isinstance(existing, dict) and not existing) only matched empty {} placeholders, so a later-arriving entry at the same index was inserted before the dumped placeholder instead of replacing it, shifting higher-index entries and breaking tool_calls[index] lookups. Added _is_placeholder() helper that detects both empty {} and all-None dumped placeholders. Applied in both accumulate_delta and _coalesce_list_by_index. Added regression tests test_dumped_placeholder_replaced_not_shifted and test_coalesce_dumped_placeholder_replaced. --- src/openai/lib/streaming/_deltas.py | 39 ++++++++++++++++---- tests/lib/streaming/test_deltas.py | 55 +++++++++++++++++++++++++++++ 2 files changed, 87 insertions(+), 7 deletions(-) diff --git a/src/openai/lib/streaming/_deltas.py b/src/openai/lib/streaming/_deltas.py index 3201cc992f..79b26c7bb6 100644 --- a/src/openai/lib/streaming/_deltas.py +++ b/src/openai/lib/streaming/_deltas.py @@ -3,6 +3,27 @@ from ..._utils import is_dict, is_list +def _is_placeholder(entry: object) -> bool: + """Detect a gap-filler placeholder that should be replaced in-place. + + When a sparse tool-call stream emits index 0 then 2, the gap at index 1 + is padded with an empty ``{}``. After the snapshot is round-tripped + through ``model_dump`` (which happens on the next chunk), that placeholder + is no longer empty — it becomes a dict of unset tool-call fields such as + ``{"id": None, "function": None, "type": None}``. Both forms must be + detected so a later-arriving entry at the same index *replaces* the + placeholder instead of being inserted before it (which would shift + higher-index entries and break ``tool_calls[index]`` lookups). + """ + if not is_dict(entry): + return False + # Empty placeholder from the padding path. + if not entry: + return True + # Dumped placeholder: every value is None (or the dict is empty). + return all(v is None for v in entry.values()) + + def accumulate_delta(acc: dict[object, object], delta: dict[object, object]) -> dict[object, object]: for key, delta_value in delta.items(): if key not in acc: @@ -85,12 +106,14 @@ def accumulate_delta(acc: dict[object, object], delta: dict[object, object]) -> acc_value.append(delta_entry) else: # The list is large enough but no entry has this - # index. If the slot at `index` is an empty - # placeholder ({}), replace it in-place. Otherwise - # insert at the correct position to keep the list - # addressable by logical index. + # index. If the slot at `index` is a placeholder + # (empty {} or a dumped placeholder with only None + # values from a model_dump round-trip), replace it + # in-place. Otherwise insert at the correct + # position to keep the list addressable by logical + # index. existing = acc_value[index] - if isinstance(existing, dict) and not existing: + if _is_placeholder(existing): acc_value[index] = delta_entry else: acc_value.insert(index, delta_entry) @@ -135,9 +158,11 @@ def _coalesce_list_by_index(lst: list[object]) -> list[object]: # logical index. while len(result) <= index: result.append({}) - # Replace the placeholder at `index` or shift if occupied + # Replace the placeholder at `index` (empty {} or a dumped + # placeholder with only None values from a model_dump round-trip) + # or shift if occupied by a real entry. existing = result[index] - if isinstance(existing, dict) and not existing: + if _is_placeholder(existing): result[index] = entry else: result.insert(index, entry) diff --git a/tests/lib/streaming/test_deltas.py b/tests/lib/streaming/test_deltas.py index aa1594d0ce..039a05fb35 100644 --- a/tests/lib/streaming/test_deltas.py +++ b/tests/lib/streaming/test_deltas.py @@ -207,3 +207,58 @@ def test_coalesce_list_by_index_sorts_by_logical_index(self) -> None: assert calls[0]["id"] == "call_a" assert calls[1]["index"] == 1 assert calls[1]["id"] == "call_b" + + def test_dumped_placeholder_replaced_not_shifted(self) -> None: + """Regression for Codex P2: after the snapshot is round-tripped through + model_dump, a gap-filler {} placeholder becomes a dict of unset + tool-call fields (e.g. {"id": None, "function": None, "type": None}). + If index 1 arrives later, it must replace that dumped placeholder + in-place, not insert before it (which would shift the index-2 entry + to slot 3 and break tool_calls[2] lookups).""" + acc: dict[object, object] = { + "tool_calls": [ + {"index": 0, "id": "call_a", "function": {"name": "tool_a"}, "type": "function"}, + # Simulates a {} placeholder after model_dump round-trip + {"id": None, "function": None, "type": None}, + {"index": 2, "id": "call_c", "function": {"name": "tool_c"}, "type": "function"}, + ] + } + delta: dict[object, object] = { + "tool_calls": [ + {"index": 1, "id": "call_b", "function": {"name": "tool_b"}, "type": "function"}, + ] + } + result = accumulate_delta(acc, delta) + calls = cast(list[dict[str, Any]], result["tool_calls"]) + # The dumped placeholder at index 1 should be replaced, not shifted + assert len(calls) == 3 + assert calls[0]["index"] == 0 + assert calls[0]["id"] == "call_a" + assert calls[1]["index"] == 1 + assert calls[1]["id"] == "call_b" + assert calls[2]["index"] == 2 + assert calls[2]["id"] == "call_c" + + def test_coalesce_dumped_placeholder_replaced(self) -> None: + """Regression for Codex P2: _coalesce_list_by_index must also detect + dumped placeholders (all-None values from model_dump) and replace them + in-place instead of inserting before them.""" + from openai.lib.streaming._deltas import _coalesce_list_by_index + + lst: list[object] = [ + {"index": 0, "id": "call_a", "function": {"name": "tool_a"}, "type": "function"}, + # Dumped placeholder at index 1 (all values None) + {"id": None, "function": None, "type": None}, + {"index": 2, "id": "call_c", "function": {"name": "tool_c"}, "type": "function"}, + # Index 1 arriving later — should replace the placeholder + {"index": 1, "id": "call_b", "function": {"name": "tool_b"}, "type": "function"}, + ] + result = _coalesce_list_by_index(lst) + calls = cast(list[dict[str, Any]], result) + assert len(calls) == 3 + assert calls[0]["index"] == 0 + assert calls[0]["id"] == "call_a" + assert calls[1]["index"] == 1 + assert calls[1]["id"] == "call_b" + assert calls[2]["index"] == 2 + assert calls[2]["id"] == "call_c"