Repository navigation
fix(a2a): give up on an agent whose sandbox died, and cancel a prompt that timed out #82
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
87fa4b4
3d4e4f0
5ec660c
9da082f
c2d9063
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -18,6 +18,19 @@ | |
| TaskState.completed, TaskState.failed, TaskState.canceled, TaskState.rejected, | ||
| }) | ||
|
|
||
| # How long an agent on a sandbox may go without answering a poll before ``poll_a2a_task`` gives up on it, over at | ||
| # least _UNREACHABLE_MIN_POLLS polls in a row: a sandbox that died answers none, a tunnel blip drops one or two. | ||
| UNREACHABLE_AFTER_SECONDS = 60 | ||
| _UNREACHABLE_MIN_POLLS = 3 | ||
| # What a proxy in front of the agent answers when it can't reach it; any other status is the agent answering. | ||
| _GATEWAY_STATUSES = frozenset({502, 503, 504}) | ||
| _CANCEL_TIMEOUT_SECONDS = 10 | ||
|
|
||
|
|
||
| class AgentUnreachableError(TimeoutError): | ||
| """The agent stopped answering, as one whose sandbox died does. A ``TimeoutError``, so whatever handles an | ||
| agent running out of time handles this the same way, only sooner.""" | ||
|
|
||
|
|
||
| _ERROR_CODE = re.compile(r"[a-z][a-z0-9_]{0,63}") | ||
|
|
||
|
|
@@ -134,10 +147,18 @@ async def poll_a2a_task( | |
| task_id: str, | ||
| timeout_seconds: int, | ||
| poll_interval_seconds: int = 10, | ||
| *, | ||
| sandbox_id: Optional[str] = None, | ||
| ) -> dict: | ||
| """POST /a2a tasks/get until status.state is 'completed' or 'failed'. Returns the result dict.""" | ||
| """POST /a2a tasks/get until the task reaches a terminal state, and return it. | ||
|
|
||
| Out of time, it raises ``TimeoutError``. ``sandbox_id`` names the sandbox the agent runs on: then it raises | ||
| ``AgentUnreachableError``, naming it, once the agent has gone UNREACHABLE_AFTER_SECONDS without answering, | ||
| rather than waiting out the timeout. Either way it first asks the agent to cancel the task, so an agent still | ||
| working on it stops.""" | ||
| deadline = time.monotonic() + timeout_seconds | ||
| consecutive_failures = 0 | ||
| answered_at, unanswered, last_failure = time.monotonic(), 0, "" | ||
| while time.monotonic() < deadline: | ||
| backoff = min(poll_interval_seconds + consecutive_failures * 5, 60) | ||
| await asyncio.sleep(backoff) | ||
|
|
@@ -149,30 +170,60 @@ async def poll_a2a_task( | |
| }, timeout=30) | ||
| resp.raise_for_status() | ||
| data = resp.json() | ||
| if "error" in data: | ||
| consecutive_failures += 1 | ||
| logger.warning(f"A2A poll returned error (consec={consecutive_failures}, will retry): {data['error']}") | ||
| continue | ||
| result = data["result"] | ||
| except httpx.HTTPStatusError as e: | ||
| if 400 <= e.response.status_code < 500: | ||
| raise | ||
| consecutive_failures += 1 | ||
| logger.warning(f"A2A poll got {e.response.status_code} (consec={consecutive_failures}, will retry)") | ||
| continue | ||
| if e.response.status_code in _GATEWAY_STATUSES: | ||
| unanswered, last_failure = unanswered + 1, f"HTTP {e.response.status_code}" | ||
| else: | ||
| answered_at, unanswered = time.monotonic(), 0 | ||
| except httpx.HTTPError as e: | ||
| consecutive_failures += 1 | ||
| logger.warning(f"A2A poll failed (consec={consecutive_failures}, will retry): {type(e).__name__}: {e}") | ||
| unanswered, last_failure = unanswered + 1, f"{type(e).__name__}: {e}" if str(e) else type(e).__name__ | ||
| else: | ||
| answered_at, unanswered = time.monotonic(), 0 | ||
| if "error" in data: | ||
| consecutive_failures += 1 | ||
| logger.warning(f"A2A poll returned error (consec={consecutive_failures}, will retry): {data['error']}") | ||
| continue | ||
| consecutive_failures = 0 | ||
| if data["result"]["status"]["state"] in _TERMINAL_TASK_STATES: | ||
| return data["result"] | ||
| continue | ||
| silent = time.monotonic() - answered_at | ||
| if sandbox_id is not None and unanswered >= _UNREACHABLE_MIN_POLLS and silent >= UNREACHABLE_AFTER_SECONDS: | ||
| await cancel_a2a_task(a2a_url, task_id) | ||
| raise AgentUnreachableError( | ||
| f"The agent on sandbox {sandbox_id} stopped answering: {unanswered} polls in a row over " | ||
| f"{silent:.0f}s got no reply (last: {last_failure})") | ||
|
|
||
| consecutive_failures = 0 | ||
| state = result["status"]["state"] | ||
| if state in _TERMINAL_TASK_STATES: | ||
| return result | ||
|
|
||
| await cancel_a2a_task(a2a_url, task_id) | ||
| raise TimeoutError(f"A2A task {task_id} did not complete within {timeout_seconds}s") | ||
|
|
||
|
|
||
| async def cancel_a2a_task(a2a_url: str, task_id: str) -> None: | ||
| """POST /a2a tasks/cancel, so the agent stops working on ``task_id``. Best effort: an agent that can't be | ||
| reached, or won't cancel, is only logged.""" | ||
| try: | ||
| async with httpx.AsyncClient() as client: | ||
| resp = await client.post(f"{a2a_url}/a2a", json={ | ||
| "jsonrpc": "2.0", "id": "cancel", "method": "tasks/cancel", | ||
| "params": {"id": task_id}, | ||
| }, timeout=_CANCEL_TIMEOUT_SECONDS) | ||
| resp.raise_for_status() | ||
| error = resp.json().get("error") | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
If the cancel endpoint returns HTTP 200 with valid JSON such as Prompt To Fix With AIThis is a comment left during a code review.
Path: src/agent_env/a2a_agent/protocol.py
Line: 215
Comment:
**Bad cancel reply hides timeout**
If the cancel endpoint returns HTTP 200 with valid JSON such as `null` or a list, calling `.get("error")` raises `AttributeError`. This hides the original `TimeoutError`, so a user-sim timeout fails the run instead of ending the conversation. Treat an unexpected cancel reply as a failed best-effort cancel.
---
For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Leaving this as is. An A2A endpoint answers with a JSON-RPC object; an error status or a non-JSON body (say, a proxy's HTML page) raises There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. That’s fair. The existing handling covers the failure modes produced by an A2A endpoint or an intervening proxy, while a successful 200 response with non-object JSON is outside the protocol contract and not produced by the callers here. Given the repository’s preference against speculative defensive code, I withdraw this finding; leaving it as-is is reasonable. |
||
| except (httpx.HTTPError, ValueError) as e: | ||
| logger.warning(f"Couldn't cancel A2A task {task_id}: {type(e).__name__}: {e}") | ||
| return | ||
| if error: | ||
| logger.warning(f"The agent didn't cancel A2A task {task_id}: {error}") | ||
| else: | ||
| logger.info(f"Canceled A2A task {task_id}") | ||
|
|
||
|
|
||
| @dataclass(frozen=True) | ||
| class TerminalResponse: | ||
| """What's pulled from an A2A terminal status message: the text reply, the typed | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.