Skip to content

Commit e74cbc6

Browse files
monody0007claude
andcommitted
fix(connection): honor $/cancel_request and reply -32800 to cancelled requests
ACP v1 defines the `$/cancel_request` notification (part of the stable v1 schema vendored in schema/, 1.23.0; see https://agentclientprotocol.com/protocol/v1/cancellation), and the SDK already ships `PROTOCOL_METHODS["cancel_request"]`, but `Connection` never handled it: - An incoming `$/cancel_request` fell through to the router, which logged an ERROR traceback (method not found) while the targeted handler kept running and the peer never received the required response. - A handler that ended with `CancelledError` (cancelled from inside the agent) produced no response at all, leaving the peer waiting forever. - Cancelling the task awaiting `send_request` only dropped the local future; the peer kept working on the request. - `close()` hung forever when a handler answered its own cancellation with a result, because the response was queued on an already-closed sender. `Connection` now tracks in-flight incoming requests by id, cancels the handler task on `$/cancel_request`, answers cancelled requests with `-32800 Request cancelled` (a handler may still return a partial result), sends `$/cancel_request` when a still-pending outgoing request is cancelled locally, and skips replies once the connection is closed. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
1 parent 9d07d78 commit e74cbc6

5 files changed

Lines changed: 379 additions & 12 deletions

File tree

‎docs/experimental-v2.md‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -150,6 +150,8 @@ Only the initial v2 request is reduced to the common v1 initialization fields
150150
when an agent selects v1.
151151

152152
Client-side fallback is application controlled and may require opening a new
153-
transport. Protocol-level request cancellation is not yet exposed by the
154-
experimental runtime; `session/cancel` remains available for cancelling active
153+
transport. Protocol-level request cancellation (`$/cancel_request`) is handled
154+
by the shared connection layer, as in v1: cancelling the task awaiting a request
155+
notifies the peer, and an incoming cancellation cancels the handler task and
156+
answers with `-32800`. `session/cancel` remains available for cancelling active
155157
session work.

‎src/acp/connection.py‎

Lines changed: 66 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
import inspect
66
import json
77
import logging
8+
import sys
89
from collections.abc import Awaitable, Callable
910
from dataclasses import dataclass
1011
from enum import Enum
@@ -14,6 +15,7 @@
1415

1516
from ._transport import NdjsonTransport, Transport
1617
from .exceptions import RequestError
18+
from .meta import PROTOCOL_METHODS
1719
from .task import MessageSender, TaskSupervisor
1820
from .telemetry import span_context
1921

@@ -37,6 +39,8 @@ class StreamEvent:
3739

3840
StreamObserver = Callable[[StreamEvent], Awaitable[None] | None]
3941

42+
_CANCEL_REQUEST_METHOD = PROTOCOL_METHODS["cancel_request"]
43+
4044

4145
class Connection:
4246
"""Minimal JSON-RPC 2.0 connection over newline-delimited JSON frames."""
@@ -54,6 +58,7 @@ def __init__(
5458
self._handler = handler
5559
self._next_request_id = 0
5660
self._pending: dict[int, asyncio.Future[Any]] = {}
61+
self._incoming: dict[Any, asyncio.Task[Any]] = {}
5762
self._tasks = TaskSupervisor(source="acp.Connection")
5863
self._tasks.add_error_handler(self._on_task_error)
5964
self._closed = False
@@ -117,16 +122,14 @@ async def send_request(self, method: str, params: JsonValue | None = None) -> An
117122
payload = {"jsonrpc": "2.0", "id": request_id, "method": method, "params": params}
118123
try:
119124
await self._transport.send(payload)
120-
except BaseException:
121-
self._pending.pop(request_id, None)
122-
future.cancel()
123-
raise
124-
self._notify_observers(StreamDirection.OUTGOING, payload)
125-
try:
125+
self._notify_observers(StreamDirection.OUTGOING, payload)
126126
return await future
127-
except asyncio.CancelledError:
128-
self._pending.pop(request_id, None)
127+
except BaseException as exc:
128+
still_pending = self._pending.pop(request_id, None) is not None
129129
future.cancel()
130+
if isinstance(exc, asyncio.CancelledError) and still_pending:
131+
# The request may already be on the wire; ask the peer to stop working on it.
132+
self._send_cancel_request(request_id)
130133
raise
131134

132135
async def send_notification(self, method: str, params: JsonValue | None = None) -> None:
@@ -152,13 +155,18 @@ async def _receive_loop(self) -> None:
152155
def _process_message(self, message: dict[str, Any]) -> None:
153156
method = message.get("method")
154157
has_id = "id" in message
158+
if method == _CANCEL_REQUEST_METHOD and not has_id:
159+
self._cancel_incoming(message.get("params"))
160+
return
155161
if method is not None: # this is a request or notification
156162
# {"jsonrpc": "2.0", "id": 1, "method": "foo", "params": {...}} # request
157163
# {"jsonrpc": "2.0", "method": "foo", "params: {...}} # notification
158-
self._tasks.create(
164+
task = self._tasks.create(
159165
self._run_request(message) if has_id else self._run_notification(message),
160166
name="acp.Connection.request" if has_id else "acp.Connection.notification",
161167
)
168+
if has_id:
169+
self._track_incoming(message["id"], task)
162170
return
163171
if has_id: # this is a response, {"id", "result" | "error"}
164172
self._handle_response(message)
@@ -184,8 +192,56 @@ def _notify_observers(self, direction: StreamDirection, message: dict[str, Any])
184192
def _on_observer_error(self, task: asyncio.Task[Any], exc: BaseException) -> None:
185193
logging.exception("Stream observer coroutine failed", exc_info=exc)
186194

195+
def _track_incoming(self, request_id: Any, task: asyncio.Task[Any]) -> None:
196+
try:
197+
self._incoming[request_id] = task
198+
except TypeError: # unhashable id, nothing can refer to it
199+
return
200+
201+
def _forget(done: asyncio.Task[Any]) -> None:
202+
if self._incoming.get(request_id) is done:
203+
del self._incoming[request_id]
204+
205+
task.add_done_callback(_forget)
206+
207+
def _cancel_incoming(self, params: Any) -> None:
208+
request_id = params.get("requestId") if isinstance(params, dict) else None
209+
try:
210+
task = self._incoming.get(request_id)
211+
except TypeError:
212+
return
213+
# Unknown or already finished requests are ignored, as the protocol allows.
214+
if task is not None:
215+
task.cancel()
216+
217+
def _send_cancel_request(self, request_id: int) -> None:
218+
if self._closed or self._disconnected:
219+
return
220+
self._tasks.create(
221+
self.send_notification(_CANCEL_REQUEST_METHOD, {"requestId": request_id}),
222+
name="acp.Connection.cancel_request",
223+
on_error=self._on_cancel_request_error,
224+
)
225+
226+
def _on_cancel_request_error(self, task: asyncio.Task[Any], exc: BaseException) -> None:
227+
logging.debug("Failed to send %s", _CANCEL_REQUEST_METHOD, exc_info=exc)
228+
187229
async def _run_request(self, message: dict[str, Any]) -> None:
188-
payload = await self._execute_request(message)
230+
try:
231+
payload = await self._execute_request(message)
232+
except asyncio.CancelledError:
233+
if self._closed:
234+
raise
235+
# Cancelled by ``$/cancel_request`` or from inside the handler; either way the
236+
# protocol still requires a response for the original request.
237+
task = asyncio.current_task()
238+
if sys.version_info >= (3, 11) and task is not None:
239+
task.uncancel()
240+
payload = {"jsonrpc": "2.0", "id": message["id"], "error": RequestError.request_cancelled().to_error_obj()}
241+
if self._closed:
242+
# The transport is gone (e.g. the handler returned a result while being
243+
# cancelled by ``close()``); sending would never complete.
244+
return
189245
await self._transport.send(payload)
190246
self._notify_observers(StreamDirection.OUTGOING, payload)
191247

‎src/acp/exceptions.py‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,5 +42,9 @@ def resource_not_found(cls, uri: str | None = None) -> RequestError:
4242
data = {"uri": uri} if uri is not None else None
4343
return cls(-32002, "Resource not found", data)
4444

45+
@classmethod
46+
def request_cancelled(cls, data: dict[str, Any] | None = None) -> RequestError:
47+
return cls(-32800, "Request cancelled", data)
48+
4549
def to_error_obj(self) -> dict[str, Any]:
4650
return {"code": self.code, "message": str(self), "data": self.data}

0 commit comments

Comments
 (0)