Skip to content

fix(a2a): mark a failed remote task as an error event - #7398

Open
ferponse wants to merge 2 commits into
google:mainfrom
ferponse:fix/a2a-failed-task-error-event
Open

ferponse wants to merge 2 commits into
google:mainfrom
ferponse:fix/a2a-failed-task-error-event

Conversation

@ferponse

@ferponse ferponse commented Oct 3, 2026

Copy link
Copy Markdown
Contributor

Please ensure you have read the contribution guide before creating a pull request.

Link to Issue or Description of Change

1. Link to an existing issue (if applicable):

Problem:

When a remote A2A task ends in TASK_STATE_FAILED, RemoteA2aAgent yields the failure as an ordinary answer. The event has the remote agent's text as content, and both error_code and error_message are None. A caller cannot tell "the remote agent answered" from "the remote agent's task failed", and the failure's text goes into the conversation history as if it were a reply.

The response converters (convert_a2a_task_to_event, convert_a2a_status_update_to_event, and the legacy ones) read the task's message but not its state. task mode already handles the failure separately, with _create_task_failure_events. The default mode does nothing.

Two details mean the fix cannot live in just one place:

  • The handler is chosen per response. _handle_a2a_response_v2 only runs when the task metadata carries ADK's integration-extension marker. A non-ADK A2A server never sets it, so its responses always go through _handle_a2a_response.
  • The failure arrives in two shapes. When streamed, it is a TaskStatusUpdateEvent with the new status. When not streamed, it is the Task itself.

This is the follow-up #6708 asks for. #6713 implemented it, but it no longer applies: since a6a4052 RemoteA2aAgent lives in src/google/adk/a2a/agent/_remote_a2a_agent.py, and agents/remote_a2a_agent.py, the file that PR edits, only re-exports it. This PR redoes the fix on the current layout. The approach of setting a stable error_code with the remote's text as error_message comes from #6713.

Solution:

The fix goes where _run_async_impl dispatches each response, right after either handler returns. That one place covers both handlers, both response shapes, and any custom a2a_*_converter set in A2aRemoteAgentConfig:

  • _failed_task_status(a2a_response) returns the status of a response that reports TASK_STATE_FAILED:
    • from the TaskStatusUpdateEvent when the response is a status update;
    • from the Task when it is the task itself;
    • None for a message, an artifact update, or any other state.
  • _mark_task_failed(event, status, ctx, agent_name) then:
    • sets error_code = A2A_TASK_FAILED_ERROR_CODE ("A2A_TASK_FAILED"), keeping one the converter already set;
    • sets error_message to the status message's text, falling back to the event's own text and then to "Remote A2A task failed";
    • creates the event when the converter returned none, for example for a bare FAILED status update with no message.
  • The mark is applied before the after_request interceptors run, so they see the final event.

What does not change:

  • The event's content is kept, so anything that reads the remote's text today still finds it. It is now also the error_message.
  • task mode is left as it is. It already yields its own error and finish_task events for a failed task, and the check is skipped there.
  • Other terminal states (COMPLETED, CANCELED, REJECTED, INPUT_REQUIRED, …) are not marked. Only FAILED means the remote task failed; a test pins this.
  • _compat.TS_FAILED is used for the comparison, and _compat.part_text/is_text_part to read the message, so the same code works with a2a-sdk 0.3.x and 1.x.

With error_code set, the failure event is a final response (Event.is_final_response()) and is stored with the session as an error.

Testing Plan

Unit Tests:

  • I have added or updated unit tests for my change.
  • All unit tests pass locally.

New tests in tests/unittests/a2a/agent/test_remote_a2a_agent.py. They go through the public Runner with the file's existing _run_remote_task_responses harness, so the real response handlers and converters run:

  • test_failed_remote_task_is_an_error_event (8 cases):
    • parametrized over the legacy and v2 handlers, streamed and non-streamed, and with or without a failure message;
    • checks error_code, error_message, is_final_response(), the task id metadata, and that the stored session event is the error;
    • checks that the working update before the failure, and the turns after it, are not marked.
  • test_only_a_failed_remote_task_is_an_error_event (8 cases): COMPLETED, CANCELED, REJECTED and INPUT_REQUIRED, through both handlers, yield no error.

Without the change, the 8 failure cases fail (assert None == 'A2A_TASK_FAILED') and the 8 other-state cases pass.

$ pytest tests/unittests/a2a
891 passed, 51 skipped

$ pytest tests/unittests/agents tests/unittests/tools/test_agent_tool.py
908 passed, 2 xfailed

$ pytest -n 8 tests/unittests
17444 passed, 87 skipped, 25 xfailed, 2 xpassed, 17 failed, 1 error

None of the failures is in a file this PR touches:

  • code_executors/test_gke_code_executor.py (5);
  • features/ (11);
  • tools/spanner/test_spanner_tool_settings.py (1);
  • tools/test_skill_toolset.py::test_integration_python_fifo_in_working_dir_does_not_block (1);
  • the error is the integrations/daytona ImportError on CreateSandboxFromImageParams.

Run on their own, those files give the same result with and without this change: the 5 GKE executor tests fail and the rest pass. The other 12 only fail under -n 8.

pre-commit run --files <changed files> passes. mypy on _remote_a2a_agent.py reports the same single error before and after the change.

Manual End-to-End (E2E) Tests:

The script below serves a plain a2a-sdk agent, with no ADK on the server side, whose task narrates while working and then fails. It consumes the agent with RemoteA2aAgent through a Runner, once streamed and once not. It needs no network or credentials.

Before (main at 63aed55, also released 2.11.0):

streaming=True
  state=TASK_STATE_SUBMITTED   error_code=None             error_message=None  text=''
  state=TASK_STATE_WORKING     error_code=None             error_message=None  text='Reading a.txt.'
  state=TASK_STATE_FAILED      error_code=None             error_message=None  text='claude exited with code 1'
streaming=False
  state=TASK_STATE_FAILED      error_code=None             error_message=None  text='claude exited with code 1'

After:

streaming=True
  state=TASK_STATE_SUBMITTED   error_code=None             error_message=None  text=''
  state=TASK_STATE_WORKING     error_code=None             error_message=None  text='Reading a.txt.'
  state=TASK_STATE_FAILED      error_code=A2A_TASK_FAILED  error_message='claude exited with code 1'  text='claude exited with code 1'
streaming=False
  state=TASK_STATE_FAILED      error_code=A2A_TASK_FAILED  error_message='claude exited with code 1'  text='claude exited with code 1'
Reproduction script
"""A remote A2A task that fails: what event does RemoteA2aAgent yield for it?

Serves a plain a2a-sdk agent (no ADK on the server side) whose task narrates
while working and then ends in TASK_STATE_FAILED, and consumes it with
RemoteA2aAgent through a Runner, streaming and not.
"""

import asyncio
import socket
import threading

from a2a.client.card_resolver import parse_agent_card
from a2a.client.client import ClientConfig
from a2a.client.client_factory import ClientFactory
from a2a.helpers.proto_helpers import new_task
from a2a.server.agent_execution import AgentExecutor
from a2a.server.agent_execution import RequestContext
from a2a.server.events import EventQueue
from a2a.server.tasks import InMemoryTaskStore
from a2a.server.tasks import TaskUpdater
from a2a.types import Part
from a2a.types import TaskState
from google.adk.a2a import _compat
from google.adk.agents.remote_a2a_agent import RemoteA2aAgent
from google.adk.runners import Runner
from google.adk.sessions.in_memory_session_service import InMemorySessionService
from google.genai import types
import httpx
from starlette.applications import Starlette
import uvicorn


class FailingExecutor(AgentExecutor):

  async def execute(self, context: RequestContext, event_queue: EventQueue):
    await event_queue.enqueue_event(
        new_task(
            context.task_id, context.context_id, TaskState.TASK_STATE_SUBMITTED
        )
    )
    updater = TaskUpdater(event_queue, context.task_id, context.context_id)
    await updater.start_work(
        message=updater.new_agent_message([Part(text="Reading a.txt.")])
    )
    await updater.failed(
        message=updater.new_agent_message(
            [Part(text="claude exited with code 1")]
        )
    )

  async def cancel(self, context, event_queue):
    pass


def serve(port: int) -> None:
  card = parse_agent_card({
      "name": "remote",
      "description": "d",
      "version": "1",
      "url": f"http://127.0.0.1:{port}/",
      "preferredTransport": "JSONRPC",
      "capabilities": {"streaming": True},
      "defaultInputModes": ["text"],
      "defaultOutputModes": ["text"],
      "skills": [],
  })
  app = Starlette()
  _compat.attach_a2a_routes_to_app(
      app,
      agent_card=card,
      agent_executor=FailingExecutor(),
      task_store=InMemoryTaskStore(),
  )
  uvicorn.run(app, host="127.0.0.1", port=port, log_level="error")


async def run(base: str, streaming: bool) -> None:
  agent = RemoteA2aAgent(
      name="remote",
      agent_card=f"{base}/.well-known/agent-card.json",
      a2a_client_factory=ClientFactory(
          config=ClientConfig(
              streaming=streaming,
              httpx_client=httpx.AsyncClient(timeout=30),
          )
      ),
  )
  runner = Runner(
      app_name="repro", agent=agent, session_service=InMemorySessionService()
  )
  session = await runner.session_service.create_session(
      app_name="repro", user_id="u"
  )
  print(f"streaming={streaming}")
  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")]),
  ):
    state = (
        ((event.custom_metadata or {}).get("a2a:response") or {}).get("status")
        or {}
    ).get("state")
    text = " ".join(
        p.text for p in (event.content.parts if event.content else []) if p.text
    )
    print(
        f"  state={state!s:22} error_code={event.error_code!s:16}"
        f" error_message={event.error_message!r}  text={text!r}"
    )


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)
  await run(base, streaming=True)
  await run(base, streaming=False)


asyncio.run(main())

Checklist

  • I have read the CONTRIBUTING.md document.
  • I have performed a self-review of my own code.
  • I have commented my code, particularly in hard-to-understand areas.
  • I have added tests that prove my fix is effective or that my feature works.
  • New and existing unit tests pass locally with my changes.
  • I have manually tested my changes end-to-end.
  • Any dependent changes have been merged and published in downstream modules.

Additional context

How we work around this today. In production, our orchestrator delegates to a coding agent running in a sandbox, through RemoteA2aAgent. When that agent's task fails (for example claude exited with code 1), the user saw the error text as an ordinary assistant message. We subclass RemoteA2aAgent and override both handlers, which is the workaround suggested in #4309:

class _SandboxRemoteA2aAgent(RemoteA2aAgent):
  async def _handle_a2a_response(self, a2a_response, ctx):
    event = await super()._handle_a2a_response(a2a_response, ctx)
    return self._mark_if_failed(event, a2a_response, ctx)

  async def _handle_a2a_response_v2(self, a2a_response, ctx):
    event = await super()._handle_a2a_response_v2(a2a_response, ctx)
    return self._mark_if_failed(event, a2a_response, ctx)

  def _mark_if_failed(self, event, a2a_response, ctx):
    if event is None or not _is_failed_task(a2a_response):
      return event
    event.error_code = "SANDBOX_TASK_FAILED"
    event.error_message = f"The code agent could not finish the task: {_text_of(event)}"
    return event


def _is_failed_task(a2a_response) -> bool:
  if not isinstance(a2a_response, tuple):
    return False
  task, update = a2a_response
  status = (
      update.status
      if isinstance(update, TaskStatusUpdateEvent)
      else getattr(task, "status", None)
  )
  return getattr(status, "state", None) == TaskState.TASK_STATE_FAILED

It works, but every consumer has to find out on its own that both handlers need overriding: overriding only _v2 silently does nothing against a non-ADK server. It also depends on two private methods. With this change, callers can read event.error_code instead.

When a remote A2A task ends in TASK_STATE_FAILED, RemoteA2aAgent yielded
the failure as an ordinary answer: the remote's text as content, and no
error_code or error_message. The response converters read the task's
message but not its state, and only task mode handled the failure.

The failure is now marked where _run_async_impl dispatches each
response, so it covers both handlers (a non-ADK server never sets the
extension marker that routes to _handle_a2a_response_v2), both shapes
(a TaskStatusUpdateEvent when streamed, the Task when not) and any
custom a2a_*_converter. The event gets error_code A2A_TASK_FAILED and the
remote's text as error_message, or one is created when the response
converted to none. Its content is kept. Task mode keeps its own failure
events, and other terminal states are not marked.

Redoes google#6713 on the current layout: RemoteA2aAgent moved to
a2a/agent/_remote_a2a_agent.py in a6a4052.

Closes google#6708
@astrogilda

Copy link
Copy Markdown

thanks @ferponse, ran 23d253f with an A2A 1.2.1 server that writes before FAILED, streamed and nonstreamed. The write stays visible and ADK reports A2A_TASK_FAILED.

One wrinkle: bare nonstreamed FAILED reused earlier "working" text as error_message; streaming used the generic fallback. runnable cases, with a separate installed reader. Happy to lift the after-write cases into your tests.

A non-streamed failure arrives as the task itself. With no status message,
the task converter falls back to the task's history, whose last agent
message is what the remote agent said while working, and that text became
the error message. A streamed failure with no reason already got the generic
message. Take the error message only from the failed status's own message,
so both shapes agree.
@ferponse

ferponse commented Oct 5, 2026

Copy link
Copy Markdown
Contributor Author

Thanks @astrogilda, good catch. Reproduced it: with no status message, a non-streamed failure arrives as the task itself, and the task converter falls back to the last agent message in the history, so the working text became the error_message. Fixed in 78aad40: the error message now comes only from the failed status's own message, with the generic fallback otherwise, so streamed and non-streamed agree. The content is unchanged, so the progress text is still on the event. Added test_a_failure_without_a_reason_is_not_worded_with_earlier_progress (legacy/v2 × streamed/non-streamed); it fails on 23d253f for the non-streamed legacy case and passes now.

On lifting your after-write cases into the tests: the new test covers the same shape (progress before a bare FAILED) with the file's own harness, so I'd rather keep the suite self-contained, but thanks for the offer and for the runnable cases.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Follow-up to #4309: fix for failed A2A tasks needs both _handle_a2a_response/_v2, _compat.TS_FAILED, and error_code — still reproducible in 2.6.2

2 participants