Skip to content

A2A (default executor, streaming): RemoteA2aAgent receives the final reply and the adk_request_confirmation pause twice #7247

Description

@msteiner-google

Describe the Bug

With a stock to_a2a() server and a default RemoteA2aAgent, the caller gets content the remote already delivered a second time:

  1. A text reply arrives twice: once as thought=True (from a WORKING status update) and once as the answer (from the final artifact).
  2. A tool declared with require_confirmation=True produces its adk_request_confirmation function call twice with the same id. Both events have long_running_tool_ids set: the first is thought=True under WORKING, the second comes under INPUT_REQUIRED.

Both copies are yielded by Runner.run_async and persisted in the caller's session. This happens on the default (legacy) A2aAgentExecutor path with streaming, which is what to_a2a() + RemoteA2aAgent(agent_card=<url>) uses on a2a-sdk 1.x (AgentCardBuilder defaults to streaming=True). It does not happen with A2aAgentExecutor(..., force_new_version=True) or with non-streaming. Part 1 was reported in #3207 (ADK 1.16.0), which was closed as stale.

Steps to Reproduce

  1. Run the script below. It is self-contained: in process, scripted model, no network or credentials.
  2. Compare the printed counts.

Expected Behavior

Each piece of remote content becomes one ADK event on the caller: the answer text once, and the confirmation request once per call id.

Observed Behavior

text reply:
  text 'The answer is 42.'  thought=True  long_running=False  task_state=TASK_STATE_WORKING
  text 'The answer is 42.'  thought=False  long_running=False  task_state=TASK_STATE_WORKING
  counts: {"text 'The answer is 42.'": 2}
tool with require_confirmation=True:
  call publish id=adk-4320efab-...  thought=True  long_running=False  task_state=TASK_STATE_WORKING
  call adk_request_confirmation id=adk-0890223c-...  thought=True  long_running=True  task_state=TASK_STATE_WORKING
  call adk_request_confirmation id=adk-0890223c-...  thought=False  long_running=True  task_state=TASK_STATE_INPUT_REQUIRED
  counts: {... 'call adk_request_confirmation id=adk-0890223c-...': 2}

Raw stream the client receives for the confirmation case:

Task(SUBMITTED)
StatusUpdate(WORKING)
StatusUpdate(WORKING, msg=[function_call publish])
StatusUpdate(WORKING, msg=[function_call adk_request_confirmation id=X, long_running])
StatusUpdate(WORKING, msg=[function_response publish])
StatusUpdate(INPUT_REQUIRED, msg=[function_call adk_request_confirmation id=X, long_running])

For the text case: StatusUpdate(WORKING, msg=[text]), then ArtifactUpdate([same text]), then StatusUpdate(COMPLETED).

Cause (2.9.2):

  • a2a/executor/a2a_agent_executor.py:227-245 streams every ADK event as a status update with its message (a2a/converters/event_converter.py:588-600).
  • task_result_aggregator.py:73 rewrites the long-running update's state to WORKING.
  • At the end, the same content is sent again: the last message's parts become the artifact (a2a_agent_executor.py:264-282), or the stored long-running message goes out in the terminal status (a2a_agent_executor.py:291-301).
  • RemoteA2aAgent._handle_a2a_response (agents/remote_a2a_agent.py:1370-1412) turns each of these into an event.
  • The v2 executor avoids both problems: it strips long-running calls from streamed events (a2a_agent_executor_impl.py:204) and sends content only as artifacts.

Related: RemoteA2aAgent(use_legacy=False) does not switch the server to the v2 path on a2a-sdk 1.x. new_integration_extension.py:43-51 sets the header in state['http_kwargs'], but 1.x transports read only ClientCallContext.service_parameters. This affects the workaround suggested in #6343.

Environment Details

  • ADK: google-adk 2.9.2 (same code on main @ eb1a002)
  • a2a-sdk: 1.1.5
  • Python: 3.13.12
  • OS: Linux

Minimal Reproduction Code

import asyncio, collections, httpx
from google.adk.a2a.utils.agent_to_a2a import to_a2a
from google.adk.agents import LlmAgent
from google.adk.agents.remote_a2a_agent import RemoteA2aAgent
from google.adk.models.base_llm import BaseLlm
from google.adk.models.llm_response import LlmResponse
from google.adk.runners import Runner
from google.adk.sessions import InMemorySessionService
from google.adk.tools.function_tool import FunctionTool
from google.genai import types

class Scripted(BaseLlm):
    async def generate_content_async(self, req, stream=False):
        last = req.contents[-1].parts[-1]
        if self.model == "text" or last.function_response:
            part = types.Part(text="The answer is 42.")
        else:
            part = types.Part(function_call=types.FunctionCall(name="publish", args={"value": "391"}))
        yield LlmResponse(content=types.Content(role="model", parts=[part]))

def publish(value: str) -> dict:
    """Publish a value."""
    return {"status": "published", "value": value}

async def run(model):
    server = to_a2a(LlmAgent(name="remote", model=Scripted(model=model), instruction="x",
                             tools=[FunctionTool(publish, require_confirmation=True)]))
    async with server.router.lifespan_context(server):
        remote = RemoteA2aAgent(name="remote",
                                agent_card="http://localhost:8000/.well-known/agent-card.json",
                                httpx_client=httpx.AsyncClient(transport=httpx.ASGITransport(app=server)))
        runner = Runner(app_name="caller", agent=remote,
                        session_service=InMemorySessionService(), auto_create_session=True)
        seen = collections.Counter()
        async for e in runner.run_async(user_id="u", session_id="s",
                new_message=types.Content(role="user", parts=[types.Part(text="hi")])):
            for p in (e.content.parts if e.content else None) or []:
                if p.function_call: seen[f"call {p.function_call.name} {p.function_call.id}"] += 1
                elif p.text: seen[f"text {p.text!r}"] += 1
        print(model, dict(seen))

async def main():
    await run("text")   # text 'The answer is 42.': 2
    await run("gated")  # call adk_request_confirmation <id>: 2

asyncio.run(main())

Suggested direction

Fix it where the content is emitted instead of filtering on the client by text equality:

  • In the legacy A2aAgentExecutor._handle_request, don't republish content that was already streamed. Either close with a COMPLETED status and no repeated parts, or stream agent content as artifact updates, the way v2 does.
  • For long-running calls, send the call once, in the terminal INPUT_REQUIRED/AUTH_REQUIRED status. Keep it out of intermediate updates, as LongRunningFunctions.process_event does in v2.
  • Alternatively, make the v2 executor the default. If that route is taken, make use_legacy=False send the extension through service_parameters on a2a-sdk 1.x.

Activity

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

Metadata

Metadata

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