Skip to content

RemoteA2aAgent: streamed A2A chunks and the final event that replaces them get different event ids, so the answer can be shown twice #7395

Description

@ferponse

🔴 Required Information

Describe the Bug:

When a remote ADK agent streams its answer over A2A (A2aAgentExecutor + SSE), RemoteA2aAgent yields:

  • one event per chunk;
  • then a final event that carries the whole text again.

Every one of these events has a different Event.id. So a consumer cannot tell that the final event replaces the chunks rather than adding to them. Unless it deduplicates on its own, it shows the answer twice.

Local streaming does not have this problem. In BaseLlmFlow, the partial chunks of one response and the complete response that closes them share one event id:

# Partial chunks of one streaming response share the base id; mint a
# fresh id only after a complete event so distinct responses differ.
if not event.partial:
  model_response_event.id = Event.new_id()

Consumers rely on that id. adk-web's chat, for example, replaces the partial bubble with the final event whose id matches.

The A2A server already marks the final event correctly. In from_adk_event.py, all chunks are updates of one artifact:

  • the first chunk is sent with append=False;
  • the following chunks with append=True;
  • the final one with last_chunk=True, append=False, which in A2A means "replace".

On the client, convert_a2a_artifact_update_to_event and the legacy handler in _handle_a2a_response only translate last_chunk into partial. Each Event is created without an id, so model_post_init mints a fresh one. Both append and the artifact id are dropped, and nothing left on the event links the final event to its chunks.

Steps to Reproduce:

  1. Serve an LlmAgent with to_a2a(...), using an A2aAgentExecutor whose request converter sets RunConfig(streaming_mode=StreamingMode.SSE). The model streams partial chunks; a fake LLM is enough (script below).
  2. Consume it with RemoteA2aAgent(..., use_legacy=False) and a streaming ClientFactory, through a Runner.
  3. Print event.partial, event.id and the text of every event.

Expected Behavior:

As with local streaming, all chunks of one answer and the final event that replaces them share one Event.id. Different answers (different artifacts) keep different ids.

Observed Behavior (adk-python main at 63aed55, also released 2.11.0):

partial=True   id=620eee9f-eae3…  text='Hello, '
partial=True   id=cac82b85-4908…  text='this is '
partial=True   id=7c2f138e-8a3e…  text='a streamed answer.'
partial=False  id=f849c4b4-e9a7…  text='Hello, this is a streamed answer.'
text events: 4  distinct ids: 4

A consumer that appends text per event shows Hello, this is a streamed answer.Hello, this is a streamed answer.

Environment Details:

  • ADK Library Version (pip show google-adk): 2.11.0, and main at 63aed55
  • Desktop OS: macOS
  • Python Version (python -V): 3.12

Model Information:

  • Are you using LiteLLM: No in the reproduction (fake BaseLlm). In production, yes.
  • Which model is being used: N/A in the reproduction. In production, Gemini and Claude behind the remote agents.

🟡 Optional Information

Regression:
No. ADK 2.7.1 behaves the same, on both the v2 and the legacy client paths.

Additional Context:

Without a shared id, consumers each work around this differently, and every option has drawbacks:

  • Skip the final text whenever partials were seen. This loses any text the final adds; the litellm_streaming sample and AG-UI's Python middleware do this.
  • Replace whatever partial bubble is last. adk-web's fallback does this; it breaks if another event lands in between.
  • Diff the final text against the streamed text. This is what we do in production.

An id-based replace would be exact. adk-samples' long-horizon-harness notes that clients concatenating every artifact part (it names Gemini Enterprise) render the reply twice.

Proposed fix:

  • Give every event converted from an artifact update an id derived from the artifact id, scoped to the invocation id, so that an artifact id a peer reuses in another invocation does not collide with it in the session.
  • Do this on both client paths: the v2 converter and the legacy handler.

It is a few lines and does not change partial/persistence semantics: partial events are still not persisted, and the stored final event keeps the same id. A PR will follow.

Minimal Reproduction Code:

import asyncio, socket, threading
from typing import AsyncGenerator

import httpx, uvicorn
from a2a.client.client import ClientConfig
from a2a.client.client_factory import ClientFactory
from google.adk.a2a.converters.request_converter import convert_a2a_request_to_agent_run_request
from google.adk.a2a.executor.a2a_agent_executor import A2aAgentExecutor
from google.adk.a2a.executor.config import A2aAgentExecutorConfig
from google.adk.a2a.utils.agent_to_a2a import to_a2a
from google.adk.agents.llm_agent import LlmAgent
from google.adk.agents.remote_a2a_agent import RemoteA2aAgent
from google.adk.agents.run_config import RunConfig, StreamingMode
from google.adk.models.base_llm import BaseLlm
from google.adk.models.llm_request import LlmRequest
from google.adk.models.llm_response import LlmResponse
from google.adk.runners import Runner
from google.adk.sessions.in_memory_session_service import InMemorySessionService
from google.genai import types

CHUNKS = ["Hello, ", "this is ", "a streamed answer."]


class StreamingFakeLlm(BaseLlm):
  model: str = "fake-streaming"

  async def generate_content_async(self, llm_request: LlmRequest, stream: bool = False) -> AsyncGenerator[LlmResponse, None]:
    if stream:
      for text in CHUNKS:
        yield LlmResponse(content=types.Content(role="model", parts=[types.Part(text=text)]), partial=True)
    yield LlmResponse(content=types.Content(role="model", parts=[types.Part(text="".join(CHUNKS))]), partial=False, turn_complete=True)


def sse_request(request, part_converter=None):
  run_request = convert_a2a_request_to_agent_run_request(request, part_converter) if part_converter else convert_a2a_request_to_agent_run_request(request)
  run_request.run_config = RunConfig(streaming_mode=StreamingMode.SSE)
  return run_request


def serve(port: int) -> None:
  agent = LlmAgent(name="streamer", model=StreamingFakeLlm(), instruction="Answer.")
  app = to_a2a(agent, host="127.0.0.1", port=port, agent_executor_factory=lambda runner: A2aAgentExecutor(
      runner=runner, config=A2aAgentExecutorConfig(request_converter=sse_request), force_new_version=True))
  uvicorn.run(app, host="127.0.0.1", port=port, log_level="error")


async def main() -> None:
  with socket.socket() as s:
    s.bind(("127.0.0.1", 0)); port = s.getsockname()[1]
  threading.Thread(target=serve, args=(port,), daemon=True).start()
  base = f"http://127.0.0.1:{port}"
  async with httpx.AsyncClient() as probe:
    for _ in range(100):
      try:
        if (await probe.get(f"{base}/.well-known/agent-card.json")).status_code == 200:
          break
      except httpx.HTTPError:
        pass
      await asyncio.sleep(0.1)
  client = RemoteA2aAgent(name="streamer", agent_card=f"{base}/.well-known/agent-card.json", use_legacy=False,
      a2a_client_factory=ClientFactory(config=ClientConfig(streaming=True, httpx_client=httpx.AsyncClient(timeout=30))))
  runner = Runner(app_name="e2e", agent=client, session_service=InMemorySessionService())
  session = await runner.session_service.create_session(app_name="e2e", user_id="u")
  ids = set()
  async for event in runner.run_async(user_id="u", session_id=session.id, new_message=types.Content(role="user", parts=[types.Part(text="Hi")])):
    text = "".join(p.text or "" for p in (event.content.parts if event.content else []) if not p.thought)
    if text:
      ids.add(event.id)
      print(f"partial={str(event.partial):5}  id={event.id[:13]}…  text={text!r}")
  print(f"distinct ids: {len(ids)}")


asyncio.run(main())

How often has this issue occurred?:

  • Always (100%)

Activity

  1. ferponse commented on Oct 3, 2026

    @ferponse
    ContributorAuthor

    PR with the fix: #7396

    What we do today to work around this. Our orchestrator consumes several remote ADK agents over A2A. Without an id to match on, it keeps the text already sent for the current answer. For each non-partial event, it sends only what that event adds:

    streamed_text = ""  # text already sent for the current answer
    
    
    def on_partial(text: str) -> None:
      global streamed_text
      send(text)
      streamed_text += text
    
    
    def on_non_partial(text: str) -> None:
      """A non-partial A2A event: the whole answer again, or genuinely new text."""
      global streamed_text
      new_text = (
          text[len(streamed_text):] if text.startswith(streamed_text) else text
      )
      if new_text:
        send(new_text)
        streamed_text += new_text
    
    
    def on_step_end() -> None:
      """A tool result closes the step: the next answer starts from scratch."""
      global streamed_text
      streamed_text = ""

    It works, but it compares text to guess whether an event is a repetition, and that guess has blind spots:

    • New text that starts like the old one is treated as a continuation. For example, if a remote agent narrates Done. and then sends a separate message Done. Next step…, only Next step… goes out, and the two messages are merged into one.
    • Every consumer has to reinvent it. AG-UI's ADK middleware drops partial=False text while a stream is open (MCP Tool Failure Crashes Entire ADK Multi-Agent Workflow #742 / Vertex AI (reasoningEngine) API functions 404 #400 there), so any text the final adds is lost. adk-web replaces "the last partial row", which breaks if another event arrives in between.

    With a shared id, consumers can do what they already do for local streaming: replace by id.

  2. added a commit that references this issue on Oct 3, 2026
    4416954
  3. added theissue type on Oct 5, 2026
  4. added
    a2a[Component] This issue is related a2a support inside ADK.
    on Oct 5, 2026
  5. self-assigned this
    on Oct 5, 2026
  6. sanketpatil06 commented on Oct 7, 2026

    @sanketpatil06

    Hi @ferponse,

    Thanks for the detailed report and repro. Confirmed on main. Both client paths create the Event without an id and drop artifact_id, so the chunks and the final event each get a different id. Deriving the id from (invocation_id, artifact_id) is the right direction.

    Thanks for opening #7396. We'll continue the review there.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

a2a[Component] This issue is related a2a support inside ADK.

Type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions