Skip to main content

agent_client_protocol_cookbook/
ordered_application_dispatch.rs

1//! Pattern: Bridge ordered SDK dispatch onto an application's executor.
2//!
3//! A notification handler finishing means the SDK has delivered that update,
4//! not that a UI task has applied it. If handlers enqueue updates elsewhere,
5//! enqueue response results and connection closure on **the same FIFO queue**.
6//! Its single consumer applies each update before exposing the following
7//! response to application code.
8//!
9//! This is useful for v2 resume: replay precedes the response on the wire, but
10//! awaiting `block_task()` separately from an update queue does not drain that
11//! queue. An [`on_receiving_result`] callback can enqueue a response marker
12//! behind the replay and before later inbound traffic. Use `prepare_request`
13//! to select ordered handling before sending, even when calling from outside
14//! the connection future.
15//!
16//! # Example
17//!
18//! This observer resumes one session and delivers events until the peer closes
19//! its incoming transport. Poll it on the application's executor; `apply` need
20//! not be `Send` and is called only by the foreground consumer, never by an SDK
21//! handler. The handlers capture queue senders, not application state or an
22//! owning connection. Add interactive request handlers as shown in
23//! [`connecting_as_client`](crate::connecting_as_client) for a complete client.
24//!
25//! ```
26//! use std::path::PathBuf;
27//! use agent_client_protocol::{Agent, Client, ConnectTo, Error, V2ConnectionTo};
28//! use agent_client_protocol::schema::{ProtocolVersion, v2};
29//! use futures::{channel::mpsc, StreamExt as _};
30//!
31//! enum ApplicationEvent {
32//!     Update(Box<v2::UpdateSessionNotification>),
33//!     ResumeFinished(Result<v2::ResumeSessionResponse, Error>),
34//!     Closed,
35//! }
36//!
37//! async fn observe_session(
38//!     transport: impl ConnectTo<Client> + 'static,
39//!     session_id: v2::SessionId,
40//!     cwd: PathBuf,
41//!     mut apply: impl FnMut(ApplicationEvent),
42//! ) -> Result<(), Error> {
43//!     let (events_tx, mut events_rx) = mpsc::unbounded();
44//!     let updates_tx = events_tx.clone();
45//!     let closed_tx = events_tx.clone();
46//!
47//!     Client.v2()
48//!         .on_receive_notification(
49//!             async move |update: v2::UpdateSessionNotification,
50//!                         _connection: V2ConnectionTo<Agent>| {
51//!                 // A dropped receiver means the application stopped observing.
52//!                 drop(updates_tx.unbounded_send(ApplicationEvent::Update(Box::new(update))));
53//!                 Ok(())
54//!             },
55//!             agent_client_protocol::on_receive_notification!(),
56//!         )
57//!         .on_close(async move |_connection| {
58//!             drop(closed_tx.unbounded_send(ApplicationEvent::Closed));
59//!             Ok(())
60//!         })
61//!         .connect_with(transport, async move |connection| {
62//!             // Initialization does not need an application projection barrier.
63//!             let initialized = connection.send_request(v2::InitializeRequest::new(
64//!                 ProtocolVersion::V2,
65//!                 v2::Implementation::new("ordered-client", "0.1.0"),
66//!             )).block_task().await?;
67//!             if initialized.capabilities.session.is_none() {
68//!                 return Err(Error::invalid_params().data("agent has no session support"));
69//!             }
70//!
71//!             let resume = v2::ResumeSessionRequest::new(session_id.clone(), cwd)
72//!                 .replay_from(v2::ReplayFrom::from(v2::ReplayFromStart::new()));
73//!             connection.prepare_request(resume).on_receiving_result(async move |result| {
74//!                 drop(events_tx.unbounded_send(ApplicationEvent::ResumeFinished(result)));
75//!                 // Do not wait for the consumer or for another inbound response here.
76//!                 Ok(())
77//!             })?;
78//!
79//!             while let Some(event) = events_rx.next().await {
80//!                 if let ApplicationEvent::Update(update) = &event {
81//!                     if update.session_id != session_id {
82//!                         continue;
83//!                     }
84//!                 }
85//!                 let closed = matches!(event, ApplicationEvent::Closed);
86//!                 // Update: apply the patch. ResumeFinished: expose the result only
87//!                 // now, after applying the preceding replay. Closed: invalidate
88//!                 // remaining application operations and mark the connection closed.
89//!                 apply(event);
90//!                 if closed {
91//!                     break;
92//!                 }
93//!             }
94//!             Ok(())
95//!         })
96//!         .await
97//! }
98//! ```
99//!
100//! # Ordering and liveness
101//!
102//! - Await application work sequentially in the **consumer**, if needed. Merely
103//!   spawning independent UI tasks for every event loses the application-side
104//!   ordering again.
105//! - An ordered callback must not await another inbound response, notification,
106//!   or permission exchange on the same connection: dispatch is waiting for
107//!   that callback to return. Use the foreground consumer or [`ConnectionTo::spawn`]
108//!   for work needing later traffic. If a projection-drained acknowledgement is
109//!   needed, send it from the consumer after applying the marker, not by making
110//!   the SDK callback wait.
111//! - [`on_close`] queues closure after already-dispatched notifications and
112//!   timely ordered **wire responses**. EOF also fails pending requests, but
113//!   those synthetic errors have no response barrier; their callbacks can run
114//!   after `Closed`. Treat `Closed` as terminal for application operations and
115//!   tolerate late completions rather than waiting for one callback per request.
116//!   Transport or handler failures still propagate from `connect_with`.
117//! - Eager `send_request(...).on_receiving_result(...)` can race with response
118//!   routing. Preparation closes that race, but a response routed later through
119//!   a retained `ResponseRouter` still does not impose a barrier on subsequent
120//!   wire traffic. See the SDK's [ordering contract].
121//! - This unbounded queue keeps dispatch nonblocking for clarity. Production
122//!   integrations need an explicit memory/backpressure policy. A bounded queue
123//!   must not wait on a consumer that is itself awaiting later inbound traffic.
124//! - A prompt response is an acceptance event, **not** prompt completion.
125//!   `state_update` describes session-wide foreground state. Do not attribute
126//!   the next `Idle` to a particular prompt or invent a turn boundary; v2 updates
127//!   do not carry a prompt ID.
128//!
129//! [`on_receiving_result`]: agent_client_protocol::PreparedRequest::on_receiving_result
130//! [`on_close`]: agent_client_protocol::Builder::on_close
131//! [`ConnectionTo::spawn`]: agent_client_protocol::ConnectionTo::spawn
132//! [ordering contract]: agent_client_protocol::concepts::ordering