fix(a2a): suppress a repeated snapshot after a single delta too

Snapshot suppression was gated on more than one delta having been emitted, to
avoid collapsing a delta that genuinely repeats the accumulated text. That
left the ordinary shape of a short reply broken: one delta "OK" followed by
the terminal cumulative snapshot "OK" rendered as "OKOK".

A2A marks no event as delta-or-snapshot, so an event whose text equals
everything emitted so far is inherently ambiguous. Read it as a snapshot: a
server repeating the whole reply at the end of a stream is common, a delta
that exactly reproduces the accumulated text is not, and duplicating a reply
is far worse for a reader than dropping one repeated fragment. The tradeoff
is pinned by the `identical_deltas_collapse` case rather than left implicit.

Reported by Greptile on #39513.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Peter Boers 2026-09-03 11:33:29 +02:00
parent edb455ee78
commit 13a313987a
No known key found for this signature in database
2 changed files with 31 additions and 11 deletions

View file

@ -56,7 +56,6 @@ class A2AModelResponseIterator(BaseModelResponseIterator):
self.model = model
# Text already emitted downstream, used to collapse cumulative snapshots.
self._emitted_text: str = ""
self._delta_count: int = 0
def chunk_parser(self, chunk: dict) -> GenericStreamingChunk | ModelResponseStream:
"""
@ -121,6 +120,12 @@ class A2AModelResponseIterator(BaseModelResponseIterator):
Handles both streaming styles: servers that send deltas ("O", "K") and servers that
send growing snapshots ("O", "OK") collapse to the same output.
A2A marks no event as delta-or-snapshot, so an event whose text equals everything
emitted so far is necessarily ambiguous. It is read as a snapshot, because servers
repeating the whole reply at the end of a stream are common while a delta that
exactly reproduces the accumulated text is not, and duplicating a reply is far
worse for a reader than dropping one repeated fragment.
"""
if not text:
return ""
@ -129,18 +134,13 @@ class A2AModelResponseIterator(BaseModelResponseIterator):
text_key: Final = _ignoring_whitespace(text)
if emitted_key and text_key.startswith(emitted_key):
# A cumulative snapshot: emit only its tail, which is empty when the snapshot
# just repeats everything sent so far.
suffix: Final = text[_index_after(text, len(emitted_key)) :]
if suffix.strip():
self._emitted_text += suffix
return suffix
# Same content as everything emitted so far: a snapshot repeat, except while
# only a single delta has been emitted, where a genuinely repeated delta is
# still indistinguishable from one.
if self._delta_count > 1:
return ""
self._emitted_text += suffix
return suffix
self._emitted_text += text
self._delta_count += 1
return text
def _get_finish_reason(self, chunk: dict) -> str | None:

View file

@ -121,7 +121,11 @@ def test_kagent_stream_finishes_on_completed_state():
pytest.param(["O", "OK", "OKAY"], "OKAY", id="cumulative_snapshots"),
pytest.param(["O", "K", "OK"], "OK", id="deltas_then_final_snapshot"),
pytest.param(["O", "K", "OK", "OK"], "OK", id="deltas_then_repeated_snapshots"),
pytest.param(["a", "a", "a"], "aaa", id="genuinely_repeated_deltas"),
pytest.param(["OK", "OK"], "OK", id="one_delta_then_equal_snapshot"),
# Known limitation: A2A marks no event as delta-or-snapshot, so a delta that
# exactly reproduces the accumulated text is indistinguishable from a snapshot
# and collapses. Duplicating a whole reply is the worse failure of the two.
pytest.param(["a", "a", "a"], "a", id="identical_deltas_collapse"),
pytest.param(["Hello", "world", "Hello world"], "Helloworld", id="multipart_snapshot_respaced"),
pytest.param(
["Hello", "world", "Hello world again"],
@ -153,3 +157,19 @@ def test_multipart_snapshot_is_not_re_emitted():
iterator = _iterator()
rendered = "".join(iterator.chunk_parser(e)["text"] for e in MULTIPART_SNAPSHOT_STREAM)
assert rendered == "Helloworld"
# A one-token reply: a single delta followed by the terminal cumulative snapshot. This is
# the ordinary shape of a short A2A answer, so the snapshot must not be forwarded again.
SINGLE_DELTA_STREAM = [
_status_update(text="Reply with exactly: OK", role="user", state="submitted"),
_status_update(text="OK"),
_artifact_update("OK"),
_status_update(state="completed", final=True),
]
def test_single_delta_then_snapshot_is_not_duplicated():
"""Regression: snapshot suppression must not depend on how many deltas preceded it."""
iterator = _iterator()
assert "".join(iterator.chunk_parser(e)["text"] for e in SINGLE_DELTA_STREAM) == "OK"