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