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_<response_id> and every answer delta msg_<response_id>, so any
interleaving hits the stale mapping.
This commit is contained in:
zhanghanduo 2026-08-16 21:25:36 +08:00
parent 428380fb1e
commit a40b9983c5
2 changed files with 41 additions and 1 deletions

View file

@ -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

View file

@ -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