From a4b19f8f3f7a67c90c86004317a9b0c6b7b56c8e Mon Sep 17 00:00:00 2001 From: Ben Brandt Date: Thu, 8 Oct 2026 13:13:12 +0200 Subject: [PATCH 1/2] feat(acp): add retained outgoing request cancellation handles --- md/request-cancellation.md | 23 +- md/sending-requests.md | 22 ++ src/agent-client-protocol/CHANGELOG.md | 7 + .../src/concepts/cancellation.rs | 34 ++- src/agent-client-protocol/src/jsonrpc.rs | 172 ++++++++--- .../src/jsonrpc/prepared_request_tests.rs | 127 +++++++- src/agent-client-protocol/src/lib.rs | 2 +- .../tests/jsonrpc_request_cancellation.rs | 103 ++++++- .../tests/prepared_requests.rs | 286 +++++++++++++++++- 9 files changed, 728 insertions(+), 48 deletions(-) diff --git a/md/request-cancellation.md b/md/request-cancellation.md index 5c3a56b4..f6ea5d5a 100644 --- a/md/request-cancellation.md +++ b/md/request-cancellation.md @@ -4,7 +4,8 @@ This chapter documents the `$/cancel_request` protocol-level notification and how the SDK implements it. For API usage (cancelling a `SentRequest`, observing cancellation from a -`Responder`), see the `concepts::cancellation` chapter in the +`Responder`, or retaining a `RequestCancellationHandle` alongside a callback), +see the `concepts::cancellation` chapter in the [agent-client-protocol rustdoc](https://docs.rs/agent-client-protocol). ## The `$/cancel_request` Notification @@ -56,6 +57,26 @@ but which should continue running on the peer. The peer is still expected to answer the JSON-RPC request eventually; use a notification instead when no response is expected at all. +## Retaining Cancellation Control + +Before consuming a `PreparedRequest` or `SentRequest`, call +`cancellation_handle()` to retain a cloneable `RequestCancellationHandle`. +Its `cancel()` method requests cancellation independently of response consumption: +an ordered callback or response future still receives the eventual result. + +The handle shares once-only cancellation and settlement state with +`SentRequest::cancel`, forwarded cancellation, and automatic request-drop +cancellation. It uses the same peer and proxy wrapping. Dropping the handle +neither sends cancellation nor disables automatic request-drop cancellation. + +Calling it before a prepared request is published does nothing and is not +remembered for later publication. A call racing publication may also do nothing; +call after the publishing method returns to target a pending request. Once the +SDK routes a response, fails the request, or the request is detached, retained +handles also do nothing—even if the application has not consumed the result yet. +Cancellation is cooperative; it does not guarantee that the peer stops work or +responds successfully. + ## Proxy Chains Cancellation propagates **hop by hop** rather than end to end. Request IDs are diff --git a/md/sending-requests.md b/md/sending-requests.md index 1d0d2fe3..c53d1973 100644 --- a/md/sending-requests.md +++ b/md/sending-requests.md @@ -66,6 +66,28 @@ without cancelling. It does not select ordered consumption. Outgoing order follows publication, not preparation. A notification sent between preparation and consumption enters the queue before the request. +## Cancellation after selecting a consumer + +Obtain `cancellation_handle()` before consuming a prepared or sent request to +retain explicit cancellation control alongside its callback or response future: + +```rust +let prepared = connection.prepare_request(request); +let cancellation = prepared.cancellation_handle(); +prepared.on_receiving_result(async move |result| { + application_queue.enqueue(result)?; + Ok(()) +})?; +cancellation.cancel()?; +``` + +Cancelling does not discard the eventual result. The cloneable handle shares +SDK settlement and once-only cancellation state, rather than sending an +unconditional notification for a saved request ID. Dropping the handle does +nothing. Calls before publication, after a response or local failure, or after +detaching the request do nothing. See [Request Cancellation](./request-cancellation.md) +for the full contract, including publication races and peer wrapping. + ## Drop and errors Dropping an unconsumed `PreparedRequest` sends nothing. Dropping the response diff --git a/src/agent-client-protocol/CHANGELOG.md b/src/agent-client-protocol/CHANGELOG.md index 1e4869a9..7eebbb19 100644 --- a/src/agent-client-protocol/CHANGELOG.md +++ b/src/agent-client-protocol/CHANGELOG.md @@ -2,6 +2,13 @@ ## [Unreleased] +### Added + +- Add `RequestCancellationHandle` and `cancellation_handle()` on prepared and + sent requests. Retain explicit cancellation control while an ordered callback + or future consumes the response, sharing SDK publication/settlement state and + once-only cancellation without a separate cancel-on-drop policy. + ## [3.1.0](https://github.com/agentclientprotocol/rust-sdk/compare/v3.0.0...v3.1.0) - 2026-10-07 ### Added diff --git a/src/agent-client-protocol/src/concepts/cancellation.rs b/src/agent-client-protocol/src/concepts/cancellation.rs index 2d426821..459d9335 100644 --- a/src/agent-client-protocol/src/concepts/cancellation.rs +++ b/src/agent-client-protocol/src/concepts/cancellation.rs @@ -35,6 +35,33 @@ //! original request, so this also works for requests sent through //! [`ConnectionTo::send_request_to`]. //! +//! When another task needs to request cancellation after the request is consumed, +//! retain a [`RequestCancellationHandle`] first: +//! +//! ``` +//! # use agent_client_protocol::{ConnectionTo, Error, UntypedRole}; +//! # use agent_client_protocol_test::{MyRequest, MyResponse}; +//! # fn apply_result(_result: Result) {} +//! # async fn example(cx: ConnectionTo) -> Result<(), Error> { +//! let request = cx.prepare_request(MyRequest {}); +//! let cancellation = request.cancellation_handle(); +//! request.on_receiving_result(async |result| { +//! apply_result(result); +//! Ok(()) +//! })?; +//! cancellation.cancel()?; +//! # Ok(()) +//! # } +//! ``` +//! +//! The response is still delivered to the selected callback or future. Cloned +//! handles share the same once-only cancellation state as [`SentRequest::cancel`] +//! and automatic request-drop cancellation, and remember the original peer and +//! proxy wrapping. Dropping a cancellation handle does nothing. Cancelling a +//! prepared request before publication is also a no-op; it is not remembered for +//! later publication. Responses and local failures disarm cancellation before +//! application callbacks run. Detaching a request disarms retained handles too. +//! //! Dropping a [`SentRequest`] before the SDK receives a response also sends //! `$/cancel_request`. This covers abandoned request handles and futures. For a //! request whose eventual response should be ignored, but which should continue @@ -164,8 +191,9 @@ //! If you are implementing custom routing and already know the JSON-RPC request //! ID on the peer connection you are targeting, use //! [`ConnectionTo::send_cancel_request_to`]. Most code should use -//! [`SentRequest::cancel`] instead, because the request handle already knows the -//! correct peer, request ID, and proxy wrapping. +//! [`SentRequest::cancel`] or [`RequestCancellationHandle::cancel`] instead, +//! because they share SDK settlement state and already know the correct peer, +//! request ID, and proxy wrapping. //! //! [`block_task`]: crate::SentRequest::block_task //! [`on_receiving_result`]: crate::SentRequest::on_receiving_result @@ -180,6 +208,8 @@ //! [`ConnectionTo::spawn`]: crate::ConnectionTo::spawn //! [`SentRequest`]: crate::SentRequest //! [`SentRequest::cancel`]: crate::SentRequest::cancel +//! [`RequestCancellationHandle`]: crate::RequestCancellationHandle +//! [`RequestCancellationHandle::cancel`]: crate::RequestCancellationHandle::cancel //! [`SentRequest::detach`]: crate::SentRequest::detach //! [`forward_cancellation_from`]: crate::SentRequest::forward_cancellation_from //! [`ConnectionTo::send_cancel_request_to`]: crate::ConnectionTo::send_cancel_request_to diff --git a/src/agent-client-protocol/src/jsonrpc.rs b/src/agent-client-protocol/src/jsonrpc.rs index 364800f4..6cc5c102 100644 --- a/src/agent-client-protocol/src/jsonrpc.rs +++ b/src/agent-client-protocol/src/jsonrpc.rs @@ -16,7 +16,7 @@ use std::panic::Location; use std::pin::pin; use std::sync::{ Arc, Mutex, Weak, - atomic::{AtomicBool, Ordering}, + atomic::{AtomicBool, AtomicU8, Ordering}, }; use uuid::Uuid; @@ -4507,7 +4507,6 @@ impl ConnectionTo { let remote_style = self.counterpart.remote_style(peer); let cancellation = SentRequestCancellation::new(self.message_tx.clone(), remote_style, id.clone()); - cancellation.disarm(); let pending_reply = PendingReply { method: method.clone(), role_id, @@ -5926,7 +5925,7 @@ impl RequestPublication { }; let id = id.clone(); let method = method.clone(); - self.pending_reply.cancellation_disarm.arm(); + let cancellation_disarm = self.pending_reply.cancellation_disarm.clone(); // Register before enqueueing so incoming EOF can fail every observable // request before close callbacks begin. The outgoing actor checks that // the registration still exists before sending the request. @@ -5943,6 +5942,9 @@ impl RequestPublication { } return Err(error); } + // An escaped cancellation handle must not enqueue cancellation before + // the request. A fast response or EOF may already have disarmed it. + cancellation_disarm.arm(); Ok(()) } } @@ -5969,6 +5971,17 @@ impl PreparedRequest { self.sent.method() } + /// Retain explicit cancellation control without publishing this request. + /// + /// The handle can outlive consumption of this request by an ordered callback + /// or response future. Calling it before publication is a no-op, not a + /// cancellation to apply when the request is later published. Dropping the + /// handle does not cancel. See [`RequestCancellationHandle`]. + #[must_use] + pub fn cancellation_handle(&self) -> RequestCancellationHandle { + self.sent.cancellation_handle() + } + /// Map a successful response without publishing the request. /// /// The mapper has the same contract as [`SentRequest::map`]. @@ -6144,32 +6157,127 @@ impl PreparedRequest { #[derive(Clone, Debug)] pub(crate) struct SentRequestCancellationDisarm { - armed: Arc, + state: Arc, +} + +#[repr(u8)] +enum OutgoingCancellationState { + Unpublished, + Armed, + Disarmed, } impl SentRequestCancellationDisarm { fn new() -> Self { Self { - armed: Arc::new(AtomicBool::new(true)), + state: Arc::new(AtomicU8::new(OutgoingCancellationState::Unpublished as u8)), } } fn disarm(&self) { - self.armed.store(false, Ordering::Release); + self.state + .store(OutgoingCancellationState::Disarmed as u8, Ordering::Release); } - fn arm(&self) { - self.armed.store(true, Ordering::Release); + fn arm(&self) -> bool { + self.state + .compare_exchange( + OutgoingCancellationState::Unpublished as u8, + OutgoingCancellationState::Armed as u8, + Ordering::AcqRel, + Ordering::Acquire, + ) + .is_ok() + } + + fn take_armed(&self) -> bool { + self.state + .compare_exchange( + OutgoingCancellationState::Armed as u8, + OutgoingCancellationState::Disarmed as u8, + Ordering::AcqRel, + Ordering::Acquire, + ) + .is_ok() + } + + fn is_armed(&self) -> bool { + self.state.load(Ordering::Acquire) == OutgoingCancellationState::Armed as u8 } } -struct SentRequestCancellation { +/// Explicit cancellation control for one outgoing request. +/// +/// Obtain this handle from [`PreparedRequest::cancellation_handle`] or +/// [`SentRequest::cancellation_handle`] before consuming the request. It retains +/// neither the response consumer nor application callback, so cancellation can +/// be requested independently while the eventual response is still delivered. +/// +/// Clones share the request's cancellation state with [`SentRequest::cancel`], +/// forwarded cancellation, and request-drop automatic cancellation. At most one +/// cancellation notification is attempted, with the original peer and proxy +/// wrapping. Dropping this handle neither cancels the request nor disables its +/// automatic cancellation. +/// +/// Cancellation before publication is a no-op and is not remembered for later +/// publication. A call racing publication may also be a no-op; call after the +/// publishing method returns to target a pending request. Once the SDK routes a +/// response or fails the request, or the request is detached, cancellation is +/// also a no-op, even if its callback has not yet run. +/// +/// This controls outgoing requests, unlike [`RequestCancellation`], which +/// observes a peer's cancellation of an incoming request. +#[derive(Clone)] +pub struct RequestCancellationHandle { message_tx: OutgoingMessageTx, remote_style: crate::role::RemoteStyle, request_id: RequestId, disarm: SentRequestCancellationDisarm, } +impl RequestCancellationHandle { + /// Ask the peer to cancel this request without discarding its response. + /// + /// Cancellation is cooperative: the peer may respond normally or with a + /// cancellation error. Repeated calls, including calls through other clones + /// or the original request, return `Ok(())` without sending another + /// notification. Calls before publication or after settlement are no-ops. + /// + /// # Errors + /// + /// Only the call that attempts to send reports serialization or enqueue + /// failure. A failed send is not retried by later cancellation calls. + pub fn cancel(&self) -> Result<(), crate::Error> { + if !self.disarm.take_armed() { + return Ok(()); + } + + // Build the notification lazily: most requests are never cancelled, + // so this avoids serializing a notification per outgoing request. + let untyped = self.remote_style.transform_outgoing_message( + crate::schema::v1::CancelRequestNotification::new(self.request_id.clone()), + )?; + + send_raw_message(&self.message_tx, OutgoingMessage::Notification { untyped }) + } +} + +impl Debug for RequestCancellationHandle { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("RequestCancellationHandle") + .field("request_id", &self.request_id) + .field("remote_style", &self.remote_style) + .field("armed", &self.disarm.is_armed()) + .finish_non_exhaustive() + } +} + +#[derive(Debug)] +struct SentRequestCancellation { + handle: RequestCancellationHandle, +} + impl SentRequestCancellation { fn new( message_tx: OutgoingMessageTx, @@ -6177,33 +6285,25 @@ impl SentRequestCancellation { request_id: RequestId, ) -> Self { Self { - message_tx, - remote_style, - request_id, - disarm: SentRequestCancellationDisarm::new(), + handle: RequestCancellationHandle { + message_tx, + remote_style, + request_id, + disarm: SentRequestCancellationDisarm::new(), + }, } } fn disarm(&self) { - self.disarm.disarm(); + self.handle.disarm.disarm(); } fn disarm_handle(&self) -> SentRequestCancellationDisarm { - self.disarm.clone() + self.handle.disarm.clone() } fn send(&self) -> Result<(), crate::Error> { - if !self.disarm.armed.swap(false, Ordering::AcqRel) { - return Ok(()); - } - - // Build the notification lazily: most requests are never cancelled, - // so this avoids serializing a notification per outgoing request. - let untyped = self.remote_style.transform_outgoing_message( - crate::schema::v1::CancelRequestNotification::new(self.request_id.clone()), - )?; - - send_raw_message(&self.message_tx, OutgoingMessage::Notification { untyped }) + self.handle.cancel() } } @@ -6215,16 +6315,6 @@ impl Drop for SentRequestCancellation { } } -impl Debug for SentRequestCancellation { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.debug_struct("SentRequestCancellation") - .field("request_id", &self.request_id) - .field("remote_style", &self.remote_style) - .field("armed", &self.disarm.armed.load(Ordering::Acquire)) - .finish_non_exhaustive() - } -} - /// Await the response payload for an outgoing request, watching `sources` for /// cancellation of the upstream requests it was registered with. /// @@ -6345,6 +6435,16 @@ impl SentRequest { self.cancellation.send() } + /// Retain explicit cancellation control after consuming this request. + /// + /// The handle shares the once-only cancellation state used by + /// [`cancel`](Self::cancel), response routing, and automatic request-drop + /// cancellation, but has no cancel-on-drop behavior of its own. + #[must_use] + pub fn cancellation_handle(&self) -> RequestCancellationHandle { + self.cancellation.handle.clone() + } + /// Forward cancellation of another request to this one. /// /// When the request that `source` belongs to is cancelled by its peer, diff --git a/src/agent-client-protocol/src/jsonrpc/prepared_request_tests.rs b/src/agent-client-protocol/src/jsonrpc/prepared_request_tests.rs index 2c2157f1..b0eefdcd 100644 --- a/src/agent-client-protocol/src/jsonrpc/prepared_request_tests.rs +++ b/src/agent-client-protocol/src/jsonrpc/prepared_request_tests.rs @@ -14,9 +14,13 @@ fn preparation_configuration_and_drop_send_nothing() { .map(|response| Ok(response.to_string())); assert_eq!(request.method(), "prepared"); assert!(!fixture.pending_replies.contains(request.id())); + let handle = request.cancellation_handle(); + handle.cancel().unwrap(); + drop(handle.clone()); assert!(fixture.message_rx.next().now_or_never().is_none()); assert!(fixture.task_rx.next().now_or_never().is_none()); drop(request); + handle.cancel().unwrap(); assert!(fixture.message_rx.next().now_or_never().is_none()); } @@ -74,8 +78,10 @@ fn detach_publishes_without_ordering_or_cancellation() { let mut fixture = Fixture::new(); let prepared = fixture.connection.prepare_request(request()); let id = prepared.id().clone(); + let handle = prepared.cancellation_handle(); prepared.detach().unwrap(); fixture.next_message(); + handle.cancel().unwrap(); assert!( fixture .route(id, Err(crate::Error::invalid_params())) @@ -94,6 +100,7 @@ fn callback_ordering_is_installed_before_publication_and_task_polling() { let mut fixture = Fixture::new(); let prepared = fixture.connection.prepare_request(request()); let id = prepared.id().clone(); + let handle = prepared.cancellation_handle(); let expected = result.clone(); prepared .on_receiving_result(async move |actual| { @@ -104,6 +111,8 @@ fn callback_ordering_is_installed_before_publication_and_task_polling() { fixture.next_message(); let mut acknowledgment = fixture.route(id, result).unwrap(); assert_eq!(acknowledgment.try_recv().unwrap(), None); + handle.cancel().unwrap(); + assert!(fixture.message_rx.next().now_or_never().is_none()); let task = fixture.next_task(); futures::executor::block_on(task.run_for_test()).unwrap(); futures::executor::block_on(acknowledgment).unwrap(); @@ -130,8 +139,10 @@ fn task_registration_failure_does_not_publish_or_cancel() { let mut fixture = Fixture::new(); let prepared = fixture.connection.prepare_request(request()); let id = prepared.id().clone(); + let handle = prepared.cancellation_handle(); fixture.task_rx.close(); assert!(prepared.on_receiving_result(async |_| Ok(())).is_err()); + handle.cancel().unwrap(); assert!(!fixture.pending_replies.contains(&id)); assert!(fixture.message_rx.next().now_or_never().is_none()); } @@ -152,6 +163,7 @@ fn publication_failures_reach_each_consumer_without_cancellation() { fixture.connection.prepare_request(request()) }; let id = prepared.id().clone(); + let handle = prepared.cancellation_handle(); match failure { Failure::Serialization => {} Failure::OutgoingClosed => fixture.message_rx.close(), @@ -193,6 +205,7 @@ fn publication_failures_reach_each_consumer_without_cancellation() { } _ => unreachable!(), } + handle.cancel().unwrap(); assert!(!fixture.pending_replies.contains(&id)); assert!(!matches!( fixture.message_rx.next().now_or_never(), @@ -205,9 +218,9 @@ fn publication_failures_reach_each_consumer_without_cancellation() { #[test] fn eof_after_publication_runs_the_callback_without_a_response_barrier() { let mut fixture = Fixture::new(); - fixture - .connection - .prepare_request(request()) + let prepared = fixture.connection.prepare_request(request()); + let handle = prepared.cancellation_handle(); + prepared .on_receiving_result(async |result| { assert!(is_incoming_transport_closed(&result.unwrap_err())); Ok(()) @@ -216,6 +229,7 @@ fn eof_after_publication_runs_the_callback_without_a_response_barrier() { fixture.next_message(); fixture.connection.incoming_closed.begin_close(); assert_eq!(fixture.pending_replies.close_incoming(), 1); + handle.cancel().unwrap(); futures::executor::block_on(fixture.next_task().run_for_test()).unwrap(); assert!(fixture.message_rx.next().now_or_never().is_none()); } @@ -270,6 +284,7 @@ fn consumer_waits_for_publication_before_forwarding_cancellation() { .prepare_request(request()) .forward_cancellation_from(cancellation); let id = prepared.id().clone(); + let handle = prepared.cancellation_handle(); let published_tx = prepared.register_consumer(async |_| Ok(())).unwrap(); let mut task = Box::pin(fixture.next_task().run_for_test()); assert!(task.as_mut().now_or_never().is_none()); @@ -290,6 +305,8 @@ fn consumer_waits_for_publication_before_forwarding_cancellation() { untyped.params()["requestId"], serde_json::to_value(&id).unwrap() ); + handle.cancel().unwrap(); + assert!(fixture.message_rx.next().now_or_never().is_none()); let acknowledgment = fixture .route(id, Err(crate::Error::request_cancelled())) .unwrap(); @@ -350,6 +367,110 @@ fn response_before_consumer_handoff_keeps_cancellation_disarmed() { assert!(fixture.message_rx.next().now_or_never().is_none()); } +#[test] +fn publication_activation_preserves_cancellation_and_settlement() { + #[derive(Clone, Copy)] + enum Outcome { + Pending, + Response, + Eof, + } + + for outcome in [Outcome::Pending, Outcome::Response, Outcome::Eof] { + let mut fixture = Fixture::new(); + let prepared = fixture.connection.prepare_request(request()); + let id = prepared.id().clone(); + let handle = prepared.cancellation_handle(); + let PreparedRequest { sent, publication } = prepared; + let RequestPublication { + message, + pending_reply, + message_tx, + pending_replies, + incoming_closed, + } = publication; + let cancellation_disarm = pending_reply.cancellation_disarm.clone(); + pending_replies + .subscribe(id.clone(), pending_reply, &incoming_closed) + .unwrap(); + handle.cancel().unwrap(); + assert!(fixture.message_rx.next().now_or_never().is_none()); + message_tx.unbounded_send(message.unwrap()).unwrap(); + assert!(matches!( + fixture.next_message(), + OutgoingMessage::Request { id: sent_id, .. } if sent_id == id + )); + + handle.cancel().unwrap(); + assert!(fixture.message_rx.next().now_or_never().is_none()); + // Force each outcome in the enqueue-to-activation interval of publication. + match outcome { + Outcome::Pending => { + assert!(cancellation_disarm.arm()); + handle.cancel().unwrap(); + handle.clone().cancel().unwrap(); + let OutgoingMessage::Notification { untyped } = fixture.next_message() else { + panic!("activation must preserve later cancellation"); + }; + assert_eq!(untyped.method(), "$/cancel_request"); + assert_eq!( + untyped.params()["requestId"], + serde_json::to_value(&id).unwrap() + ); + assert!(fixture.route(id, Ok(json!({"value": 42}))).is_none()); + } + Outcome::Response => { + assert!(fixture.route(id, Ok(json!({"value": 42}))).is_none()); + assert!(!cancellation_disarm.arm()); + } + Outcome::Eof => { + incoming_closed.begin_close(); + assert_eq!(fixture.pending_replies.close_incoming(), 1); + assert!(!cancellation_disarm.arm()); + } + } + handle.cancel().unwrap(); + let response = futures::executor::block_on(sent.block_task()); + if matches!(outcome, Outcome::Eof) { + assert!(is_incoming_transport_closed(&response.unwrap_err())); + } else { + assert_eq!(response.unwrap(), json!({"value": 42})); + } + assert!(fixture.message_rx.next().now_or_never().is_none()); + } +} + +#[test] +fn only_the_first_cancellation_attempt_reports_send_failure() { + let mut fixture = Fixture::new(); + let prepared = fixture.connection.prepare_request(request()); + let handle = prepared.cancellation_handle(); + let response = prepared.block_task(); + fixture.next_message(); + fixture.message_rx.close(); + assert!(handle.cancel().is_err()); + handle.clone().cancel().unwrap(); + drop(response); + assert!(!matches!( + fixture.message_rx.next().now_or_never(), + Some(Some(_)) + )); +} + +#[test] +fn cancellation_handle_traits_do_not_depend_on_the_response_type() { + fn assert_traits() {} + assert_traits::(); + let fixture = Fixture::new(); + let prepared = fixture + .connection + .prepare_request(request()) + .map(|_| Ok(std::rc::Rc::new(()))); + let handle: crate::RequestCancellationHandle = prepared.cancellation_handle(); + drop(prepared); + handle.cancel().unwrap(); +} + #[cfg(feature = "unstable_protocol_v2")] #[test] fn v2_context_exposes_both_preparation_methods() { diff --git a/src/agent-client-protocol/src/lib.rs b/src/agent-client-protocol/src/lib.rs index 9c1b0d56..590a16ad 100644 --- a/src/agent-client-protocol/src/lib.rs +++ b/src/agent-client-protocol/src/lib.rs @@ -167,7 +167,7 @@ pub use jsonrpc::{ is_incoming_transport_closed, run::{ChainRun, NullRun, RunWithConnectionTo}, }; -pub use jsonrpc::{RequestCancellation, is_cancel_request_notification}; +pub use jsonrpc::{RequestCancellation, RequestCancellationHandle, is_cancel_request_notification}; #[cfg(feature = "unstable_protocol_v2")] pub use jsonrpc::{V2Builder, V2ConnectionContext, V2ConnectionTo}; diff --git a/src/agent-client-protocol/tests/jsonrpc_request_cancellation.rs b/src/agent-client-protocol/tests/jsonrpc_request_cancellation.rs index 3be9a7c2..6f8dac92 100644 --- a/src/agent-client-protocol/tests/jsonrpc_request_cancellation.rs +++ b/src/agent-client-protocol/tests/jsonrpc_request_cancellation.rs @@ -422,6 +422,91 @@ async fn wrapped_cancel_request_cancels_wrapped_request() { .await; } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn prepared_cancellation_handle_sends_successor_wrapped_cancel_and_delivers_callback() { + use agent_client_protocol::{RawJsonRpcMessage, TransportFrame, UntypedMessage}; + use futures::channel::oneshot; + use serde_json::json; + + let (transport, mut peer) = Channel::duplex(); + let (published_tx, published_rx) = oneshot::channel(); + let client = WrappedCounterpart + .builder() + .connect_with(transport, async move |connection| { + let prepared = connection.prepare_request_to( + WrappedSuccessor, + SimpleRequest { + message: "wrapped cancel".into(), + }, + ); + let expected_id = prepared.id().clone(); + let cancellation = prepared.cancellation_handle(); + let (callback_tx, callback_rx) = oneshot::channel(); + prepared.on_receiving_result(async move |response| { + callback_tx + .send(response) + .map_err(|_| agent_client_protocol::Error::internal_error()) + })?; + assert_eq!(published_rx.await.unwrap(), expected_id); + cancellation.cancel()?; + cancellation.clone().cancel()?; + connection + .send_notification(UntypedMessage::new("after-cancellation", json!({})).unwrap())?; + assert_eq!( + callback_rx.await.unwrap().unwrap().result, + "completed despite cancellation" + ); + connection.incoming_closed().await; + Ok(()) + }); + let peer = async move { + let Some(TransportFrame::Single(RawJsonRpcMessage::Request(request))) = + peer.rx.next().await + else { + panic!("expected the successor-wrapped request"); + }; + assert_eq!(request.method.as_ref(), "_proxy/successor"); + assert_eq!( + serde_json::to_value(request.params).unwrap(), + json!({"method": "simple_method", "params": {"message": "wrapped cancel"}}) + ); + published_tx.send(request.id.clone()).unwrap(); + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(cancellation))) = + peer.rx.next().await + else { + panic!("expected the successor-wrapped cancellation"); + }; + assert_eq!(cancellation.method.as_ref(), "_proxy/successor"); + assert_eq!( + serde_json::to_value(cancellation.params).unwrap(), + json!({"method": "$/cancel_request", "params": {"requestId": request.id}}) + ); + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(boundary))) = + peer.rx.next().await + else { + panic!("expected the cancellation barrier"); + }; + assert_eq!(boundary.method.as_ref(), "after-cancellation"); + peer.tx + .unbounded_send(TransportFrame::Single(RawJsonRpcMessage::response( + request.id, + Ok(json!({"result": "completed despite cancellation"})), + ))) + .unwrap(); + drop(peer.tx); + assert!( + peer.rx.next().await.is_none(), + "callback completion or cloned handles emitted another cancellation" + ); + }; + tokio::time::timeout(std::time::Duration::from_secs(10), async { + let (client, ()) = futures::join!(client, peer); + client.unwrap(); + }) + .await + .expect("wrapped cancellation prevented callback delivery"); +} + #[tokio::test(flavor = "current_thread")] async fn cancelling_request_sent_to_successor_peer_sends_wrapped_cancel() { use tokio::task::LocalSet; @@ -754,10 +839,22 @@ async fn detached_sent_request_does_not_send_cancellation() { UntypedRole .builder() .connect_with(client_transport, async |cx| { - cx.send_request(SimpleRequest { + let request = cx.send_request(SimpleRequest { message: "detached".into(), - }) - .detach(); + }); + let cancellation = request.cancellation_handle(); + request.detach(); + cancellation.cancel()?; + cancellation.clone().cancel()?; + drop(cancellation); + + let prepared = cx.prepare_request(SimpleRequest { + message: "prepared detached".into(), + }); + let cancellation = prepared.cancellation_handle(); + prepared.detach()?; + cancellation.cancel()?; + drop(cancellation); // Barrier round trip: a cancellation sent by dropping the // detached handle would reach the server before this diff --git a/src/agent-client-protocol/tests/prepared_requests.rs b/src/agent-client-protocol/tests/prepared_requests.rs index d1740d0d..6731005b 100644 --- a/src/agent-client-protocol/tests/prepared_requests.rs +++ b/src/agent-client-protocol/tests/prepared_requests.rs @@ -6,8 +6,8 @@ use std::time::Duration; use agent_client_protocol::{ - Channel, ConnectionTo, Error, RawJsonRpcMessage, TransportBatch, TransportFrame, - UntypedMessage, role::UntypedRole, + Channel, ConnectionTo, Error, RawJsonRpcMessage, RequestCancellationHandle, TransportBatch, + TransportFrame, UntypedMessage, role::UntypedRole, }; use futures::{ FutureExt as _, StreamExt as _, @@ -26,11 +26,13 @@ async fn external_prepared_callback_holds_following_batch_entry_and_eof() { mut peer, } = external_connection(); let (boundary_tx, boundary_rx) = oneshot::channel(); + let (routed_tx, routed_rx) = oneshot::channel(); let expected = result.clone(); let caller = async move { let connection = connection_rx.await.unwrap(); let prepared = connection.prepare_request(UntypedMessage::new("prepared", json!({})).unwrap()); + let cancellation = prepared.cancellation_handle(); connection .send_notification(UntypedMessage::new("before-publication", json!({})).unwrap()) .unwrap(); @@ -50,6 +52,12 @@ async fn external_prepared_callback_holds_following_batch_entry_and_eof() { assert_eq!(started_rx.await.unwrap(), expected); assert!(notifications.next().now_or_never().is_none()); assert!(!connection.is_incoming_closed()); + cancellation.cancel().unwrap(); + cancellation.clone().cancel().unwrap(); + connection + .send_notification(UntypedMessage::new("after-routed-response", json!({})).unwrap()) + .unwrap(); + routed_rx.await.unwrap(); // This release comes from the external caller, not inbound traffic. release_tx.send(()).unwrap(); assert_eq!(notifications.next().await.unwrap().method(), "following"); @@ -80,6 +88,17 @@ async fn external_prepared_callback_holds_following_batch_entry_and_eof() { )) .unwrap(); drop(peer.tx); + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(notification))) = + peer.rx.next().await + else { + panic!("expected the routed-response barrier"); + }; + assert_eq!( + notification.method.as_ref(), + "after-routed-response", + "cancellation must be disarmed while the ordered callback is held" + ); + routed_tx.send(()).unwrap(); assert!( peer.rx.next().await.is_none(), "completed request emitted another message" @@ -94,6 +113,269 @@ async fn external_prepared_callback_holds_following_batch_entry_and_eof() { } } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn prepared_cancellation_handle_survives_publication_and_preserves_callback_result() { + fn assert_handle_traits() {} + assert_handle_traits::(); + + for result in [Ok(json!({"value": 42})), Err(Error::invalid_params())] { + let ExternalConnection { + driver, + connection_rx, + stop_tx, + mut notifications, + mut peer, + } = external_connection(); + let (boundary_tx, boundary_rx) = oneshot::channel(); + let (published_tx, published_rx) = oneshot::channel(); + let expected = result.clone(); + let caller = async move { + let connection = connection_rx.await.unwrap(); + let prepared = + connection.prepare_request(UntypedMessage::new("prepared", json!({})).unwrap()); + let request_id = prepared.id().clone(); + let cancellation: RequestCancellationHandle = prepared.cancellation_handle(); + let cloned_cancellation = cancellation.clone(); + cancellation.cancel().unwrap(); + cloned_cancellation.cancel().unwrap(); + connection + .send_notification(UntypedMessage::new("before-publication", json!({})).unwrap()) + .unwrap(); + boundary_rx.await.unwrap(); + + let (callback_tx, callback_rx) = oneshot::channel(); + prepared + .on_receiving_result(async move |response| { + callback_tx + .send(response) + .map_err(|_| Error::internal_error()) + }) + .unwrap(); + assert_eq!(published_rx.await.unwrap(), request_id); + // A prepublication cancel must neither emit traffic nor consume + // the once-only cancellation available after publication. + cloned_cancellation.cancel().unwrap(); + cancellation.cancel().unwrap(); + connection + .send_notification(UntypedMessage::new("after-cancellation", json!({})).unwrap()) + .unwrap(); + assert_eq!(callback_rx.await.unwrap(), expected); + assert_eq!(notifications.next().await.unwrap().method(), "following"); + connection.incoming_closed().await; + stop_tx.send(()).unwrap(); + }; + let peer = async move { + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(boundary))) = + peer.rx.next().await + else { + panic!("prepared cancellation must not publish the request"); + }; + assert_eq!(boundary.method.as_ref(), "before-publication"); + boundary_tx.send(()).unwrap(); + let Some(TransportFrame::Single(RawJsonRpcMessage::Request(request))) = + peer.rx.next().await + else { + panic!("expected the prepared request after callback registration"); + }; + assert_eq!(request.method.as_ref(), "prepared"); + published_tx.send(request.id.clone()).unwrap(); + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(cancellation))) = + peer.rx.next().await + else { + panic!("expected cancellation after publication"); + }; + assert_eq!(cancellation.method.as_ref(), "$/cancel_request"); + assert_eq!( + serde_json::to_value(cancellation.params).unwrap(), + json!({"requestId": request.id}) + ); + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(boundary))) = + peer.rx.next().await + else { + panic!("expected the cancellation barrier"); + }; + assert_eq!(boundary.method.as_ref(), "after-cancellation"); + peer.tx + .unbounded_send(TransportFrame::Batch( + TransportBatch::from_messages([ + RawJsonRpcMessage::response(request.id, result), + RawJsonRpcMessage::notification("following".into(), json!({})).unwrap(), + ]) + .unwrap(), + )) + .unwrap(); + drop(peer.tx); + assert!( + peer.rx.next().await.is_none(), + "cloned handles or callback completion emitted another cancellation" + ); + }; + tokio::time::timeout(Duration::from_secs(10), async { + let ((), (), driver) = futures::join!(caller, peer, driver); + driver.unwrap().unwrap(); + }) + .await + .expect("retained cancellation handle lost the ordered callback result"); + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn cancellation_handles_share_once_only_state_with_sent_request_and_drop() { + let ExternalConnection { + driver, + connection_rx, + stop_tx, + notifications: _, + mut peer, + } = external_connection(); + let caller = async move { + let connection = connection_rx.await.unwrap(); + let request = connection.send_request(UntypedMessage::new("pending", json!({})).unwrap()); + let cancellation: RequestCancellationHandle = request.cancellation_handle(); + let cloned_cancellation = cancellation.clone(); + cancellation.cancel().unwrap(); + request.cancel().unwrap(); + cloned_cancellation.cancel().unwrap(); + drop(request); + cancellation.cancel().unwrap(); + drop(cloned_cancellation); + drop(cancellation); + connection + .send_notification(UntypedMessage::new("after-drop", json!({})).unwrap()) + .unwrap(); + connection.incoming_closed().await; + stop_tx.send(()).unwrap(); + }; + let peer = async move { + let Some(TransportFrame::Single(RawJsonRpcMessage::Request(request))) = + peer.rx.next().await + else { + panic!("expected the pending request"); + }; + assert_eq!(request.method.as_ref(), "pending"); + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(cancellation))) = + peer.rx.next().await + else { + panic!("expected one cancellation"); + }; + assert_eq!(cancellation.method.as_ref(), "$/cancel_request"); + assert_eq!( + serde_json::to_value(cancellation.params).unwrap(), + json!({"requestId": request.id}) + ); + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(boundary))) = + peer.rx.next().await + else { + panic!("expected the drop barrier"); + }; + assert_eq!(boundary.method.as_ref(), "after-drop"); + drop(peer.tx); + assert!( + peer.rx.next().await.is_none(), + "explicit cancellation or drop emitted a duplicate" + ); + }; + tokio::time::timeout(Duration::from_secs(10), async { + let ((), (), driver) = futures::join!(caller, peer, driver); + driver.unwrap().unwrap(); + }) + .await + .expect("shared cancellation did not reach the wire"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn dropping_cancellation_handles_leaves_request_drop_armed() { + let ExternalConnection { + driver, + connection_rx, + stop_tx, + notifications: _, + mut peer, + } = external_connection(); + let (boundary_tx, boundary_rx) = oneshot::channel(); + let (dropped_tx, dropped_rx) = oneshot::channel(); + let caller = async move { + let connection = connection_rx.await.unwrap(); + let application_owner = std::sync::Arc::new(()); + let application_owner_weak = std::sync::Arc::downgrade(&application_owner); + let request = connection + .send_request(UntypedMessage::new("pending", json!({})).unwrap()) + .map(move |response| { + drop(application_owner); + Ok(response) + }); + let cancellation = request.cancellation_handle(); + let cloned_cancellation = cancellation.clone(); + drop(cancellation); + drop(cloned_cancellation); + connection + .send_notification(UntypedMessage::new("handle-dropped", json!({})).unwrap()) + .unwrap(); + boundary_rx.await.unwrap(); + let retained_cancellation = request.cancellation_handle(); + drop(request); + assert!( + application_owner_weak.upgrade().is_none(), + "a cancellation handle must not retain the response consumer's captured owner" + ); + dropped_rx.await.unwrap(); + retained_cancellation.cancel().unwrap(); + connection + .send_notification(UntypedMessage::new("request-dropped", json!({})).unwrap()) + .unwrap(); + connection.incoming_closed().await; + drop(retained_cancellation); + stop_tx.send(()).unwrap(); + }; + let peer = async move { + let Some(TransportFrame::Single(RawJsonRpcMessage::Request(request))) = + peer.rx.next().await + else { + panic!("expected the pending request"); + }; + assert_eq!(request.method.as_ref(), "pending"); + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(boundary))) = + peer.rx.next().await + else { + panic!("expected the handle-drop barrier"); + }; + assert_eq!( + boundary.method.as_ref(), + "handle-dropped", + "dropping a cancellation handle must not send cancellation" + ); + boundary_tx.send(()).unwrap(); + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(cancellation))) = + peer.rx.next().await + else { + panic!("request drop must cancel with a cancellation handle still retained"); + }; + assert_eq!(cancellation.method.as_ref(), "$/cancel_request"); + assert_eq!( + serde_json::to_value(cancellation.params).unwrap(), + json!({"requestId": request.id}) + ); + dropped_tx.send(()).unwrap(); + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(boundary))) = + peer.rx.next().await + else { + panic!("expected the request-drop barrier"); + }; + assert_eq!(boundary.method.as_ref(), "request-dropped"); + drop(peer.tx); + assert!( + peer.rx.next().await.is_none(), + "dropping the retained cancellation handle emitted another message" + ); + }; + tokio::time::timeout(Duration::from_secs(10), async { + let ((), (), driver) = futures::join!(caller, peer, driver); + driver.unwrap().unwrap(); + }) + .await + .expect("cancellation handle changed automatic request-drop cancellation"); +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn external_blocking_request_publishes_before_polling_without_holding_dispatch() { let ExternalConnection { From 061b1f0bda706a9851b668e4598c8cc0294fef60 Mon Sep 17 00:00:00 2001 From: Ben Brandt Date: Thu, 8 Oct 2026 13:53:16 +0200 Subject: [PATCH 2/2] fix(acp): preserve explicit cancellation after detach --- md/request-cancellation.md | 30 +++- md/sending-requests.md | 15 +- src/agent-client-protocol/CHANGELOG.md | 3 +- .../src/concepts/cancellation.rs | 57 +++++-- src/agent-client-protocol/src/jsonrpc.rs | 48 ++++-- .../src/jsonrpc/prepared_request_tests.rs | 67 +++++++- .../tests/jsonrpc_request_cancellation.rs | 13 +- .../tests/prepared_requests.rs | 161 +++++++++++++++++- 8 files changed, 348 insertions(+), 46 deletions(-) diff --git a/md/request-cancellation.md b/md/request-cancellation.md index f6ea5d5a..5be6077d 100644 --- a/md/request-cancellation.md +++ b/md/request-cancellation.md @@ -68,14 +68,34 @@ The handle shares once-only cancellation and settlement state with `SentRequest::cancel`, forwarded cancellation, and automatic request-drop cancellation. It uses the same peer and proxy wrapping. Dropping the handle neither sends cancellation nor disables automatic request-drop cancellation. +To cancel when a caller is abandoned, its application-owned guard or abandonment +handler must explicitly call `cancel()` and handle any immediate error. Calling it before a prepared request is published does nothing and is not remembered for later publication. A call racing publication may also do nothing; -call after the publishing method returns to target a pending request. Once the -SDK routes a response, fails the request, or the request is detached, retained -handles also do nothing—even if the application has not consumed the result yet. -Cancellation is cooperative; it does not guarantee that the peer stops work or -responds successfully. +hand control to independent cancellers after the publishing method returns. +Once the SDK routes a response or fails the request, new cancellation calls do +nothing—even if the application has not consumed the result yet. An attempt begun +before settlement may still enqueue afterward. + +Detachment is separate from settlement: `detach()` discards the response and +suppresses automatic cancellation on drop, but does not revoke retained explicit +handles. They can still cancel the detached request while it remains pending. + +`Ok(())` means the call encountered no immediate error, not that a notification +was sent or the peer stopped work. Cancellation is cooperative; it does not +guarantee transmission, peer cooperation, or a successful response. A handle does +not abort local callback work or keep the connection driver/response consumer alive. + +## Request Cancellation Is Not Session Cancellation + +`$/cancel_request` targets one pending JSON-RPC request, not the lifetime of the +session work it may start. In ACP v2, `session/prompt` succeeds once the user +message is inserted and returns its `messageId`. After insertion the Agent must +return success, not `-32800`, even if cancellation races the response. Stopping +active session work uses `session/cancel`. See the +[v2 cancellation specification](https://agentclientprotocol.com/protocol/v2/cancellation) +and [prompt lifecycle](https://agentclientprotocol.com/protocol/v2/prompt-lifecycle#2-prompt-accepted). ## Proxy Chains diff --git a/md/sending-requests.md b/md/sending-requests.md index c53d1973..10614074 100644 --- a/md/sending-requests.md +++ b/md/sending-requests.md @@ -78,15 +78,22 @@ prepared.on_receiving_result(async move |result| { application_queue.enqueue(result)?; Ok(()) })?; -cancellation.cancel()?; +caller.on_abandoned(move || cancellation.cancel()); ``` +Hand cancellation control to the caller's abandonment handler only after +publication returns, and have that handler deal with any immediate error. +Merely dropping the public handle does not cancel. The callback owns eventual +completion and any late-resource cleanup independently of the caller. + Cancelling does not discard the eventual result. The cloneable handle shares SDK settlement and once-only cancellation state, rather than sending an unconditional notification for a saved request ID. Dropping the handle does -nothing. Calls before publication, after a response or local failure, or after -detaching the request do nothing. See [Request Cancellation](./request-cancellation.md) -for the full contract, including publication races and peer wrapping. +nothing. Calls before publication or after a response/local failure do nothing. +Detachment suppresses automatic cancellation but preserves retained explicit +control until settlement. `Ok(())` is not acknowledgment that cancellation was +sent or took effect. See [Request Cancellation](./request-cancellation.md) for +the full contract, including publication races and peer wrapping. ## Drop and errors diff --git a/src/agent-client-protocol/CHANGELOG.md b/src/agent-client-protocol/CHANGELOG.md index 7eebbb19..96319e0b 100644 --- a/src/agent-client-protocol/CHANGELOG.md +++ b/src/agent-client-protocol/CHANGELOG.md @@ -7,7 +7,8 @@ - Add `RequestCancellationHandle` and `cancellation_handle()` on prepared and sent requests. Retain explicit cancellation control while an ordered callback or future consumes the response, sharing SDK publication/settlement state and - once-only cancellation without a separate cancel-on-drop policy. + once-only cancellation. Handle drop is inert; detaching a request suppresses + automatic cancellation without revoking retained explicit control. ## [3.1.0](https://github.com/agentclientprotocol/rust-sdk/compare/v3.0.0...v3.1.0) - 2026-10-07 diff --git a/src/agent-client-protocol/src/concepts/cancellation.rs b/src/agent-client-protocol/src/concepts/cancellation.rs index 459d9335..f06f92ca 100644 --- a/src/agent-client-protocol/src/concepts/cancellation.rs +++ b/src/agent-client-protocol/src/concepts/cancellation.rs @@ -36,31 +36,66 @@ //! [`ConnectionTo::send_request_to`]. //! //! When another task needs to request cancellation after the request is consumed, -//! retain a [`RequestCancellationHandle`] first: +//! retain a [`RequestCancellationHandle`] first. Its drop is inert; an +//! application-owned guard can explicitly cancel when the caller is abandoned: //! //! ``` -//! # use agent_client_protocol::{ConnectionTo, Error, UntypedRole}; +//! # use agent_client_protocol::{ConnectionTo, Error, RequestCancellationHandle, UntypedRole}; //! # use agent_client_protocol_test::{MyRequest, MyResponse}; +//! # use futures::channel::oneshot; //! # fn apply_result(_result: Result) {} +//! # fn clean_up_abandoned_result(_result: Result) {} +//! struct CancelOnAbandonment(RequestCancellationHandle); +//! +//! impl Drop for CancelOnAbandonment { +//! fn drop(&mut self) { +//! if let Err(error) = self.0.cancel() { +//! eprintln!("Failed to request cancellation: {}", error.code); +//! } +//! } +//! } +//! //! # async fn example(cx: ConnectionTo) -> Result<(), Error> { //! let request = cx.prepare_request(MyRequest {}); //! let cancellation = request.cancellation_handle(); -//! request.on_receiving_result(async |result| { -//! apply_result(result); +//! let (result_sender, result_received) = oneshot::channel(); +//! request.on_receiving_result(async move |result| { +//! if let Err(result) = result_sender.send(result) { +//! clean_up_abandoned_result(result); +//! } //! Ok(()) //! })?; -//! cancellation.cancel()?; +//! let _caller_guard = CancelOnAbandonment(cancellation); +//! let result = result_received.await.map_err(Error::into_internal_error)?; +//! apply_result(result); //! # Ok(()) //! # } //! ``` //! -//! The response is still delivered to the selected callback or future. Cloned -//! handles share the same once-only cancellation state as [`SentRequest::cancel`] -//! and automatic request-drop cancellation, and remember the original peer and -//! proxy wrapping. Dropping a cancellation handle does nothing. Cancelling a -//! prepared request before publication is also a no-op; it is not remembered for +//! Install the guard in the caller's scope, not the response callback, and only +//! after the publishing method returns: cancellation before publication is not +//! remembered. If the caller is abandoned, the guard requests cancellation; the +//! selected callback remains responsible for the eventual result or cleanup. +//! +//! Explicit cancellation does not discard the selected callback's or future's +//! response. Cloned handles share the same once-only cancellation state as +//! [`SentRequest::cancel`] and automatic request-drop cancellation, and remember +//! the original peer and proxy wrapping. Dropping a cancellation handle does +//! nothing. Cancelling a prepared request before publication is also a no-op; it is not remembered for //! later publication. Responses and local failures disarm cancellation before -//! application callbacks run. Detaching a request disarms retained handles too. +//! application callbacks run. Detaching suppresses only automatic cancellation: +//! retained handles can explicitly cancel a detached request while it is pending. +//! +//! `cancel()` returning `Ok(())` is not acknowledgment of transmission or peer +//! cooperation. It may have been a no-op, and an attempt begun before settlement +//! can still enqueue afterward. Cancellation controls peer request work, not +//! local callback work, and retaining a handle does not keep the connection +//! driver or response consumer alive. +//! +//! In ACP v2, a successful `session/prompt` response means the user message was +//! inserted. The request is then complete, even if session work continues: +//! stopping that work requires `session/cancel`, not request cancellation. +//! See the [v2 prompt lifecycle](https://agentclientprotocol.com/protocol/v2/prompt-lifecycle#2-prompt-accepted). //! //! Dropping a [`SentRequest`] before the SDK receives a response also sends //! `$/cancel_request`. This covers abandoned request handles and futures. For a diff --git a/src/agent-client-protocol/src/jsonrpc.rs b/src/agent-client-protocol/src/jsonrpc.rs index 6cc5c102..17cc1290 100644 --- a/src/agent-client-protocol/src/jsonrpc.rs +++ b/src/agent-client-protocol/src/jsonrpc.rs @@ -6028,7 +6028,8 @@ impl PreparedRequest { /// Returns immediate preparation or enqueue failures. Later local /// transformation errors and peer response errors are discarded along with /// successful responses. Transport failures still propagate through the - /// connection future. + /// connection future. Retained [`RequestCancellationHandle`] values can + /// still explicitly cancel the request while it remains pending. pub fn detach(self) -> Result<(), crate::Error> { let result = self.publication.publish(); self.sent.detach(); @@ -6211,7 +6212,8 @@ impl SentRequestCancellationDisarm { /// Obtain this handle from [`PreparedRequest::cancellation_handle`] or /// [`SentRequest::cancellation_handle`] before consuming the request. It retains /// neither the response consumer nor application callback, so cancellation can -/// be requested independently while the eventual response is still delivered. +/// be requested without discarding the eventual response. It does not keep the +/// connection driver or response consumer alive, or abort local callback work. /// /// Clones share the request's cancellation state with [`SentRequest::cancel`], /// forwarded cancellation, and request-drop automatic cancellation. At most one @@ -6222,8 +6224,9 @@ impl SentRequestCancellationDisarm { /// Cancellation before publication is a no-op and is not remembered for later /// publication. A call racing publication may also be a no-op; call after the /// publishing method returns to target a pending request. Once the SDK routes a -/// response or fails the request, or the request is detached, cancellation is -/// also a no-op, even if its callback has not yet run. +/// response or fails the request, subsequent cancellation calls are no-ops, +/// even if its callback has not yet run. Detaching the request suppresses only +/// automatic cancellation; retained handles can still explicitly cancel it. /// /// This controls outgoing requests, unlike [`RequestCancellation`], which /// observes a peer's cancellation of an incoming request. @@ -6242,6 +6245,11 @@ impl RequestCancellationHandle { /// cancellation error. Repeated calls, including calls through other clones /// or the original request, return `Ok(())` without sending another /// notification. Calls before publication or after settlement are no-ops. + /// An attempt begun before settlement may still enqueue afterward. + /// + /// `Ok(())` means this call encountered no immediate error, not that a + /// notification was sent or the peer stopped work. This method does not + /// wait for transmission or acknowledgment. /// /// # Errors /// @@ -6276,6 +6284,7 @@ impl Debug for RequestCancellationHandle { #[derive(Debug)] struct SentRequestCancellation { handle: RequestCancellationHandle, + cancel_on_drop: bool, } impl SentRequestCancellation { @@ -6291,6 +6300,7 @@ impl SentRequestCancellation { request_id, disarm: SentRequestCancellationDisarm::new(), }, + cancel_on_drop: true, } } @@ -6309,6 +6319,9 @@ impl SentRequestCancellation { impl Drop for SentRequestCancellation { fn drop(&mut self) { + if !self.cancel_on_drop { + return; + } if let Err(error) = self.send() { tracing::debug!(?error, "failed to auto-cancel dropped request"); } @@ -6402,7 +6415,7 @@ impl SentRequest { impl SentRequest { /// Detach this request handle without waiting for its response. /// - /// The response will be discarded when it arrives. This also disarms the + /// The response will be discarded when it arrives. This also disables the /// drop-time automatic cancellation described in /// [Drop Behavior](Self#drop-behavior), so use it for requests whose /// eventual response should be ignored, but which should keep running on @@ -6411,9 +6424,11 @@ impl SentRequest { /// all. /// /// To ask the peer to stop the request, call `cancel` instead, or drop the - /// handle while automatic cancellation is armed. - pub fn detach(self) { - self.cancellation.disarm(); + /// handle while automatic cancellation is enabled. A retained + /// [`RequestCancellationHandle`] can still explicitly cancel the detached + /// request until the SDK receives its response or fails it. + pub fn detach(mut self) { + self.cancellation.cancel_on_drop = false; } /// Send a `$/cancel_request` notification for this outgoing request. @@ -6422,12 +6437,15 @@ impl SentRequest { /// original request, so it is the preferred way to cancel a [`SentRequest`] /// when the request handle is still available. /// - /// At most one `$/cancel_request` is ever sent per request: the first - /// `cancel` call sends it (and also prevents the drop-time automatic - /// cancellation described in [Drop Behavior](Self#drop-behavior)), while - /// later calls return `Ok(())` without sending anything. Likewise, once - /// the SDK has routed the response to this handle, `cancel` becomes a - /// no-op: there is nothing left to cancel. + /// At most one cancellation attempt is made per request, shared with + /// retained handles, forwarded cancellation, and automatic cancellation + /// described in [Drop Behavior](Self#drop-behavior). Later calls return + /// `Ok(())` without another attempt, including when the first attempt failed. + /// Once the SDK has routed the response, a new call is a no-op; an attempt + /// begun before settlement may still enqueue afterward. + /// + /// `Ok(())` means this call encountered no immediate error, not that a + /// notification was sent or the peer stopped work. /// /// Errors are only reported by the call that attempts to send the /// notification. @@ -6435,7 +6453,7 @@ impl SentRequest { self.cancellation.send() } - /// Retain explicit cancellation control after consuming this request. + /// Obtain a handle that remains usable after this request is consumed. /// /// The handle shares the once-only cancellation state used by /// [`cancel`](Self::cancel), response routing, and automatic request-drop diff --git a/src/agent-client-protocol/src/jsonrpc/prepared_request_tests.rs b/src/agent-client-protocol/src/jsonrpc/prepared_request_tests.rs index b0eefdcd..223d4b18 100644 --- a/src/agent-client-protocol/src/jsonrpc/prepared_request_tests.rs +++ b/src/agent-client-protocol/src/jsonrpc/prepared_request_tests.rs @@ -74,19 +74,82 @@ fn blocking_maps_the_response_without_holding_dispatch() { } #[test] -fn detach_publishes_without_ordering_or_cancellation() { +fn detach_publishes_without_ordering_or_automatic_cancellation() { let mut fixture = Fixture::new(); let prepared = fixture.connection.prepare_request(request()); let id = prepared.id().clone(); let handle = prepared.cancellation_handle(); prepared.detach().unwrap(); fixture.next_message(); - handle.cancel().unwrap(); + assert!(fixture.message_rx.next().now_or_never().is_none()); assert!( fixture .route(id, Err(crate::Error::invalid_params())) .is_none() ); + handle.cancel().unwrap(); + assert!(fixture.message_rx.next().now_or_never().is_none()); + assert!(fixture.task_rx.next().now_or_never().is_none()); +} + +#[test] +fn detached_handles_preserve_once_only_control_until_response_routing() { + for eager in [false, true] { + for result in [ + Ok(json!({"value": 42})), + Err(crate::Error::request_cancelled()), + ] { + let mut fixture = Fixture::new(); + let (id, handle) = if eager { + let sent = fixture.connection.send_request(request()); + let id = sent.id().clone(); + let handle = sent.cancellation_handle(); + sent.detach(); + (id, handle) + } else { + let prepared = fixture.connection.prepare_request(request()); + let id = prepared.id().clone(); + let handle = prepared.cancellation_handle(); + prepared.detach().unwrap(); + (id, handle) + }; + assert!(matches!( + fixture.next_message(), + OutgoingMessage::Request { id: sent_id, .. } if sent_id == id + )); + assert!(fixture.message_rx.next().now_or_never().is_none()); + handle.cancel().unwrap(); + handle.clone().cancel().unwrap(); + let OutgoingMessage::Notification { untyped } = fixture.next_message() else { + panic!("a detached pending request must retain explicit cancellation"); + }; + assert_eq!(untyped.method(), "$/cancel_request"); + assert_eq!( + untyped.params()["requestId"], + serde_json::to_value(&id).unwrap() + ); + assert!(fixture.message_rx.next().now_or_never().is_none()); + assert!(fixture.route(id, result).is_none()); + handle.cancel().unwrap(); + drop(handle); + assert!(fixture.message_rx.next().now_or_never().is_none()); + assert!(fixture.task_rx.next().now_or_never().is_none()); + } + } +} + +#[test] +fn detached_handles_disarm_on_eof_without_a_response_consumer() { + let mut fixture = Fixture::new(); + let prepared = fixture.connection.prepare_request(request()); + let handle = prepared.cancellation_handle(); + prepared.detach().unwrap(); + fixture.next_message(); + fixture.connection.incoming_closed.begin_close(); + assert_eq!(fixture.pending_replies.close_incoming(), 1); + handle.cancel().unwrap(); + handle.clone().cancel().unwrap(); + drop(handle); assert!(fixture.message_rx.next().now_or_never().is_none()); assert!(fixture.task_rx.next().now_or_never().is_none()); } diff --git a/src/agent-client-protocol/tests/jsonrpc_request_cancellation.rs b/src/agent-client-protocol/tests/jsonrpc_request_cancellation.rs index 6f8dac92..142e32b8 100644 --- a/src/agent-client-protocol/tests/jsonrpc_request_cancellation.rs +++ b/src/agent-client-protocol/tests/jsonrpc_request_cancellation.rs @@ -843,22 +843,23 @@ async fn detached_sent_request_does_not_send_cancellation() { message: "detached".into(), }); let cancellation = request.cancellation_handle(); + let cloned_cancellation = cancellation.clone(); request.detach(); - cancellation.cancel()?; - cancellation.clone().cancel()?; + drop(cloned_cancellation); drop(cancellation); let prepared = cx.prepare_request(SimpleRequest { message: "prepared detached".into(), }); let cancellation = prepared.cancellation_handle(); + let cloned_cancellation = cancellation.clone(); prepared.detach()?; - cancellation.cancel()?; + drop(cloned_cancellation); drop(cancellation); - // Barrier round trip: a cancellation sent by dropping the - // detached handle would reach the server before this - // request. + // Barrier round trip: any automatic cancellation from + // detaching or dropping retained handles would reach the + // server before this request. let barrier = cx .send_request(SimpleRequest { message: "barrier".into(), diff --git a/src/agent-client-protocol/tests/prepared_requests.rs b/src/agent-client-protocol/tests/prepared_requests.rs index 6731005b..10fbb33a 100644 --- a/src/agent-client-protocol/tests/prepared_requests.rs +++ b/src/agent-client-protocol/tests/prepared_requests.rs @@ -17,7 +17,11 @@ use serde_json::json; #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn external_prepared_callback_holds_following_batch_entry_and_eof() { - for result in [Ok(json!({"value": 42})), Err(Error::invalid_params())] { + for result in [ + Ok(json!({"value": 42})), + Err(Error::invalid_params()), + Err(Error::request_cancelled()), + ] { let ExternalConnection { driver, connection_rx, @@ -118,7 +122,11 @@ async fn prepared_cancellation_handle_survives_publication_and_preserves_callbac fn assert_handle_traits() {} assert_handle_traits::(); - for result in [Ok(json!({"value": 42})), Err(Error::invalid_params())] { + for result in [ + Ok(json!({"value": 42})), + Err(Error::invalid_params()), + Err(Error::request_cancelled()), + ] { let ExternalConnection { driver, connection_rx, @@ -219,6 +227,155 @@ async fn prepared_cancellation_handle_survives_publication_and_preserves_callbac } } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn detached_requests_retain_explicit_cancellation_control() { + for prepare in [false, true] { + for result in [Ok(json!({"value": 42})), Err(Error::request_cancelled())] { + let ExternalConnection { + driver, + connection_rx, + stop_tx, + mut notifications, + mut peer, + } = external_connection(); + let (detached_tx, detached_rx) = oneshot::channel(); + let caller = async move { + let connection = connection_rx.await.unwrap(); + let application_owner = std::sync::Arc::new(()); + let application_owner_weak = std::sync::Arc::downgrade(&application_owner); + let (request_id, cancellation) = if prepare { + let request = connection + .prepare_request(UntypedMessage::new("detached", json!({})).unwrap()) + .map(move |response| { + drop(application_owner); + Ok(response) + }); + let request_id = request.id().clone(); + let cancellation = request.cancellation_handle(); + request.detach().unwrap(); + (request_id, cancellation) + } else { + let request = connection + .send_request(UntypedMessage::new("detached", json!({})).unwrap()) + .map(move |response| { + drop(application_owner); + Ok(response) + }); + let request_id = request.id().clone(); + let cancellation = request.cancellation_handle(); + request.detach(); + (request_id, cancellation) + }; + assert!( + application_owner_weak.upgrade().is_none(), + "a retained cancellation handle must not retain the detached response consumer" + ); + connection + .send_notification(UntypedMessage::new("after-detach", json!({})).unwrap()) + .unwrap(); + assert_eq!(detached_rx.await.unwrap(), request_id); + let cloned_cancellation = cancellation.clone(); + cancellation.cancel().unwrap(); + cloned_cancellation.cancel().unwrap(); + connection + .send_notification( + UntypedMessage::new("after-cancellation", json!({})).unwrap(), + ) + .unwrap(); + + // The following notification is dispatched after routing the + // response, even though detach discarded its consumer. + assert_eq!(notifications.next().await.unwrap().method(), "following"); + cancellation.cancel().unwrap(); + cloned_cancellation.cancel().unwrap(); + drop(cloned_cancellation); + drop(cancellation); + connection + .send_notification( + UntypedMessage::new("after-routed-response", json!({})).unwrap(), + ) + .unwrap(); + connection.incoming_closed().await; + stop_tx.send(()).unwrap(); + }; + let peer = async move { + let Some(TransportFrame::Single(RawJsonRpcMessage::Request(request))) = + peer.rx.next().await + else { + panic!("expected the detached request"); + }; + assert_eq!(request.method.as_ref(), "detached"); + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(boundary))) = + peer.rx.next().await + else { + panic!("expected the detach barrier"); + }; + assert_eq!( + boundary.method.as_ref(), + "after-detach", + "detach must not emit automatic cancellation" + ); + detached_tx.send(request.id.clone()).unwrap(); + let Some(TransportFrame::Single(message)) = peer.rx.next().await else { + panic!("expected explicit cancellation after detach"); + }; + assert!( + serde_json::to_value(&message).unwrap().get("id").is_none(), + "cancellation must have no outer request id" + ); + let RawJsonRpcMessage::Notification(cancellation) = message else { + panic!("cancellation must be a notification"); + }; + assert_eq!(cancellation.method.as_ref(), "$/cancel_request"); + assert_eq!( + serde_json::to_value(cancellation.params).unwrap(), + json!({"requestId": request.id}) + ); + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(boundary))) = + peer.rx.next().await + else { + panic!("expected the cancellation barrier"); + }; + assert_eq!( + boundary.method.as_ref(), + "after-cancellation", + "cloned handles must emit exactly one cancellation" + ); + peer.tx + .unbounded_send(TransportFrame::Batch( + TransportBatch::from_messages([ + RawJsonRpcMessage::response(request.id, result), + RawJsonRpcMessage::notification("following".into(), json!({})).unwrap(), + ]) + .unwrap(), + )) + .unwrap(); + let Some(TransportFrame::Single(RawJsonRpcMessage::Notification(boundary))) = + peer.rx.next().await + else { + panic!("expected the routed-response barrier"); + }; + assert_eq!( + boundary.method.as_ref(), + "after-routed-response", + "cancelling after response routing must not emit another cancellation" + ); + drop(peer.tx); + assert!( + peer.rx.next().await.is_none(), + "detached request or retained handles emitted another message" + ); + }; + tokio::time::timeout(Duration::from_secs(10), async { + let ((), (), driver) = futures::join!(caller, peer, driver); + driver.unwrap().unwrap(); + }) + .await + .expect("detached request lost independent cancellation control"); + } + } +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn cancellation_handles_share_once_only_state_with_sent_request_and_drop() { let ExternalConnection {