Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions src/google/adk/a2a/executor/a2a_agent_executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@
from .utils import execute_after_agent_interceptors
from .utils import execute_after_event_interceptors
from .utils import execute_before_agent_interceptors
from .utils import failure_summary

logger = logging.getLogger('google_adk.' + __name__)

Expand Down Expand Up @@ -152,7 +153,13 @@ async def execute(
try:
await self._handle_request(context, event_queue)
except Exception as e:
logger.error('Error handling A2A request: %s', e, exc_info=True)
error_id, peer_text = failure_summary(e)
logger.error(
'Error handling A2A request [error_id=%s]: %s',
error_id,
e,
exc_info=True,
)
# Publish failure event
try:
await event_queue.enqueue_event(
Expand All @@ -164,7 +171,7 @@ async def execute(
message=Message(
message_id=platform_uuid.new_uuid(),
role=_compat.ROLE_AGENT,
parts=[_compat.make_text_part(str(e))],
parts=[_compat.make_text_part(peer_text)],
),
),
final=True,
Expand Down
11 changes: 9 additions & 2 deletions src/google/adk/a2a/executor/a2a_agent_executor_impl.py
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@
from .utils import execute_after_agent_interceptors
from .utils import execute_after_event_interceptors
from .utils import execute_before_agent_interceptors
from .utils import failure_summary

logger = logging.getLogger('google_adk.' + __name__)

Expand Down Expand Up @@ -153,7 +154,13 @@ async def execute(
run_request,
)
except Exception as e:
logger.error('Error handling A2A request: %s', e, exc_info=True)
error_id, peer_text = failure_summary(e)
logger.error(
'Error handling A2A request [error_id=%s]: %s',
error_id,
e,
exc_info=True,
)
# Publish failure event
try:
await event_queue.enqueue_event(
Expand All @@ -165,7 +172,7 @@ async def execute(
message=Message(
message_id=str(uuid.uuid4()),
role=_compat.ROLE_AGENT,
parts=[_compat.make_text_part(str(e))],
parts=[_compat.make_text_part(peer_text)],
),
),
final=True,
Expand Down
15 changes: 15 additions & 0 deletions src/google/adk/a2a/executor/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,9 @@
# limitations under the License.
from __future__ import annotations

import os
from typing import Optional
import uuid

from a2a.server.agent_execution.context import RequestContext
from a2a.server.events import Event as A2AEvent
Expand All @@ -28,6 +30,19 @@
from .executor_context import ExecutorContext


def failure_summary(error: Exception) -> tuple[str, str]:
"""Returns an error id and the peer-facing text for an execution failure.

The exception text can carry paths, hostnames and credentials, so the peer
gets the id only. `ADK_A2A_DEBUG_ERRORS=1` appends the exception text.
"""
error_id = uuid.uuid4().hex[:8]
text = f'Agent execution failed. (error_id: {error_id})'
if os.environ.get('ADK_A2A_DEBUG_ERRORS') == '1':
text = f'{text}: {type(error).__name__}: {error}'
return error_id, text


async def _enqueue_canceled_task_event(
context: RequestContext,
event_queue: EventQueue,
Expand Down
20 changes: 20 additions & 0 deletions tests/unittests/a2a/executor/test_a2a_agent_executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
# limitations under the License.

import asyncio
import re
from unittest.mock import AsyncMock
from unittest.mock import Mock
from unittest.mock import patch
Expand Down Expand Up @@ -800,6 +801,25 @@ async def test_execute_with_exception_handling(self):
assert failure_event.status.state == _compat.TS_FAILED
_assert_final(failure_event)

@pytest.mark.asyncio
async def test_execute_failure_message_omits_exception_text(self):
"""The peer gets a fixed summary and an id, never the exception text."""
self.mock_context.task_id = "test-task-id"
self.mock_context.current_task = None
self.mock_request_converter.side_effect = FileNotFoundError(
2, "No such file or directory", "/srv/secrets/sa-key.json"
)

await self.executor.execute(self.mock_context, self.mock_event_queue)

failure_event = self.mock_event_queue.enqueue_event.call_args_list[-1][0][0]
part = failure_event.status.message.parts[0]
text = part.text if _compat.IS_A2A_V1 else part.root.text
assert "/srv/secrets/sa-key.json" not in text
assert re.fullmatch(
r"Agent execution failed\. \(error_id: [0-9a-f]{8}\)", text
)

@pytest.mark.asyncio
async def test_handle_request_with_aggregator_message(self):
"""Test that the final task status event includes message from aggregator."""
Expand Down
45 changes: 41 additions & 4 deletions tests/unittests/a2a/executor/test_a2a_agent_executor_impl.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
from __future__ import annotations

import asyncio
import re
from unittest.mock import AsyncMock
from unittest.mock import Mock
from unittest.mock import patch
Expand Down Expand Up @@ -500,10 +501,46 @@ async def test_execute_with_exception_handling(self):
assert failure_event.status.state == _compat.TS_FAILED
_assert_final(failure_event)
_failure_part = failure_event.status.message.parts[0]
if _compat.IS_A2A_V1:
assert "Test error" in _failure_part.text
else:
assert "Test error" in _failure_part.root.text
_text = _failure_part.text if _compat.IS_A2A_V1 else _failure_part.root.text
assert "Agent execution failed." in _text
assert "Test error" not in _text

@pytest.mark.asyncio
async def test_execute_failure_message_omits_exception_text(self):
"""The peer gets a fixed summary and an id, never the exception text."""
self.mock_context.task_id = "test-task-id"
self.mock_context.current_task = None
self.mock_request_converter.side_effect = FileNotFoundError(
2, "No such file or directory", "/srv/secrets/sa-key.json"
)

await self.executor.execute(self.mock_context, self.mock_event_queue)

failure_event = self.mock_event_queue.enqueue_event.call_args_list[-1][0][0]
part = failure_event.status.message.parts[0]
text = part.text if _compat.IS_A2A_V1 else part.root.text
assert "/srv/secrets/sa-key.json" not in text
assert re.fullmatch(
r"Agent execution failed\. \(error_id: [0-9a-f]{8}\)", text
)

@pytest.mark.asyncio
async def test_execute_failure_message_includes_detail_when_opted_in(
self, monkeypatch
):
"""`ADK_A2A_DEBUG_ERRORS=1` puts the exception text back."""
monkeypatch.setenv("ADK_A2A_DEBUG_ERRORS", "1")
self.mock_context.task_id = "test-task-id"
self.mock_context.current_task = None
self.mock_request_converter.side_effect = Exception("Test error")

await self.executor.execute(self.mock_context, self.mock_event_queue)

failure_event = self.mock_event_queue.enqueue_event.call_args_list[-1][0][0]
part = failure_event.status.message.parts[0]
text = part.text if _compat.IS_A2A_V1 else part.root.text
assert "Agent execution failed." in text
assert "Test error" in text

@pytest.mark.asyncio
async def test_handle_request_with_non_working_state(self):
Expand Down
Loading