From a40b9983c55039485e65258fb6b75004a3c67816 Mon Sep 17 00:00:00 2001 From: zhanghanduo Date: Sun, 16 Aug 2026 21:25:36 +0800 Subject: [PATCH] fix(anthropic): reopen a content block when a resumed item id returns _get_or_start_block trusted the item_id -> block index map without checking whether that block was still open, so a provider that reuses one item id for a whole run and interleaves channels got a delta addressed to a stopped block. Replaying reasoning, text, reasoning, text produced content_block_delta index=0 after content_block_stop index=0, which is not a valid Anthropic stream. Treat the mapping as valid only while it points at the open block, and rebind the item to a fresh block otherwise. Items registered through response.output_item.added still reuse their block, since that block is the open one while its deltas arrive. Apodex Deep Research is what surfaced this: it labels every reasoning delta of a run rs_ and every answer delta msg_, so any interleaving hits the stale mapping. --- .../responses_adapters/streaming_iterator.py | 5 ++- ...t_responses_adapters_streaming_iterator.py | 37 +++++++++++++++++++ 2 files changed, 41 insertions(+), 1 deletion(-) diff --git a/litellm/llms/anthropic/experimental_pass_through/responses_adapters/streaming_iterator.py b/litellm/llms/anthropic/experimental_pass_through/responses_adapters/streaming_iterator.py index 8b9a3ed6764..3e6e6df4565 100644 --- a/litellm/llms/anthropic/experimental_pass_through/responses_adapters/streaming_iterator.py +++ b/litellm/llms/anthropic/experimental_pass_through/responses_adapters/streaming_iterator.py @@ -101,8 +101,11 @@ class AnthropicResponsesStreamWrapper: content_block: Mapping[str, object], ) -> int: mapped_index: Final = self._item_id_to_block_index.get(item_id) if item_id else None - if mapped_index is not None: + if mapped_index is not None and mapped_index == self._open_block_index: return mapped_index + # A resumed item whose block already closed needs a fresh one: Anthropic rejects + # a delta addressed to a stopped block. Providers that reuse one item id for a + # whole run, then interleave channels, land here. if item_id is None and self._open_block_index is not None and self._open_block_type == block_type: return self._open_block_index diff --git a/tests/test_litellm/llms/anthropic/experimental_pass_through/responses_adapters/test_responses_adapters_streaming_iterator.py b/tests/test_litellm/llms/anthropic/experimental_pass_through/responses_adapters/test_responses_adapters_streaming_iterator.py index f8c34c9bf7d..81b207bb259 100644 --- a/tests/test_litellm/llms/anthropic/experimental_pass_through/responses_adapters/test_responses_adapters_streaming_iterator.py +++ b/tests/test_litellm/llms/anthropic/experimental_pass_through/responses_adapters/test_responses_adapters_streaming_iterator.py @@ -181,6 +181,43 @@ class TestProcessEventReasoningDeltaWithoutOutputItemAdded: ] assert chunks[3]["content_block"] == {"type": "text", "text": ""} + def test_resumed_item_id_opens_a_fresh_block(self): + """A provider that reuses one item id per run, then interleaves channels, would + otherwise address a delta to a block that has already stopped.""" + response = SimpleNamespace(status="completed", output=[], usage=None) + chunks = _process_all( + [ + {"type": "response.reasoning_summary_text.delta", "item_id": "rs_1", "delta": "Think A"}, + {"type": "response.output_text.delta", "item_id": "msg_1", "delta": "Answer A"}, + {"type": "response.reasoning_summary_text.delta", "item_id": "rs_1", "delta": "Think B"}, + {"type": "response.output_text.delta", "item_id": "msg_1", "delta": "Answer B"}, + {"type": "response.completed", "response": response}, + ] + ) + + assert [(chunk["type"], chunk.get("index")) for chunk in chunks] == [ + ("content_block_start", 0), + ("content_block_delta", 0), + ("content_block_stop", 0), + ("content_block_start", 1), + ("content_block_delta", 1), + ("content_block_stop", 1), + ("content_block_start", 2), + ("content_block_delta", 2), + ("content_block_stop", 2), + ("content_block_start", 3), + ("content_block_delta", 3), + ("content_block_stop", 3), + ("message_delta", None), + ("message_stop", None), + ] + assert [chunk["content_block"]["type"] for chunk in chunks if chunk["type"] == "content_block_start"] == [ + "thinking", + "text", + "thinking", + "text", + ] + class TestResponseCompletedUsage: """The Anthropic ``message_delta`` usage must report cache reads/writes and