Skip to main content

agent_client_protocol/concepts/
ordering.rs

1//! Message ordering, concurrency, and the dispatch loop.
2//!
3//! Understanding how agent-client-protocol processes messages is key to writing correct code.
4//! This chapter explains the dispatch loop and the ordering guarantees you
5//! can rely on.
6//!
7//! # The Dispatch Loop
8//!
9//! Each connection has a central **dispatch loop** that processes incoming
10//! messages one at a time. When a message arrives, it is passed to your
11//! handlers in order until one claims it.
12//!
13//! The key property: **the dispatch loop waits for each handler to complete
14//! before processing the next message.** This gives you sequential ordering
15//! guarantees within a single connection.
16//!
17//! # Ordered Callbacks Hold the Loop
18//!
19//! Request and notification callbacks registered with [`on_receive_request`]
20//! and [`on_receive_notification`] run inside the dispatch loop. The loop is
21//! blocked until the callback completes.
22//!
23//! Registering [`on_receiving_result`] or [`on_receiving_ok_result`] returns
24//! immediately. If registration happens before a peer response is routed
25//! during its original dispatch, response handling then holds an ordering
26//! barrier until the callback completes. A pending-request failure delivered
27//! without an incoming response, or a response routed later, does not carry
28//! that barrier.
29//!
30//! To select ordering before sending, use [`prepare_request`] followed by
31//! [`PreparedRequest::on_receiving_result`]. The callback task and ordering
32//! marker are registered before publication, so a concurrent connection cannot
33//! route a fast response before ordering has been selected. Eager
34//! `send_request(...).on_receiving_result(...)` retains the conditional guarantee
35//! above.
36//!
37//! Session-start helpers are two-phase: [`on_session_start`] and
38//! [`on_proxy_session_start`] perform framework-owned session setup under this
39//! ordering guarantee, then invoke the user callback in a spawned task so it
40//! can consume later session traffic. No user callback code runs under the
41//! session-setup ordering guarantee.
42//!
43//! While a callback holds the dispatch loop, this means:
44//! - No other messages are processed while your callback runs
45//! - You can safely do setup before "releasing" control back to the loop
46//! - Messages are processed in the order they arrive
47//!
48//! # Deadlock Risk
49//!
50//! Because callbacks can hold the dispatch loop, it's easy to create deadlocks.
51//! The most common pattern:
52//!
53//! ```ignore
54//! // DEADLOCK: This blocks the loop waiting for a response,
55//! // but the response can't arrive because the loop is blocked!
56//! builder.on_receive_request(async |request: MyRequest, responder, cx| {
57//!     let response = cx.send_request(SomeRequest { ... })
58//!         .block_task()  // <-- Waits for response
59//!         .await?;       // <-- But response can never arrive!
60//!     responder.respond(response)
61//! }, on_receive_request!());
62//! ```
63//!
64//! The response can never arrive because the dispatch loop is blocked waiting
65//! for your callback to complete.
66//!
67//! # `block_task` vs `on_receiving_result`
68//!
69//! When you send a request, you get a [`SentRequest`] with two ways to handle it:
70//!
71//! ## `block_task()` - Does not hold dispatch while you process
72//!
73//! Use this from a task that already runs outside the dispatch loop, such as a
74//! foreground `connect_with` future or a spawned task:
75//!
76//! ```
77//! # use agent_client_protocol::{Client, Agent, ConnectTo};
78//! # use agent_client_protocol_test::MyRequest;
79//! # async fn example(transport: impl ConnectTo<Client>) -> Result<(), agent_client_protocol::Error> {
80//! # Client.builder().connect_with(transport, async |cx| {
81//! cx.spawn({
82//!     let cx = cx.clone();
83//!     async move {
84//!         // Safe: we're in a spawned task, not blocking the dispatch loop
85//!         let response = cx.send_request(MyRequest {})
86//!             .block_task()
87//!             .await?;
88//!         // Process response...
89//!         Ok(())
90//!     }
91//! })?;
92//! # Ok(())
93//! # }).await?;
94//! # Ok(())
95//! # }
96//! ```
97//!
98//! The dispatch loop continues immediately after delivering the response.
99//! Your code receives the response and can take as long as it wants.
100//!
101//! ## `on_receiving_result()` - A peer response callback can hold the loop
102//!
103//! Use a prepared request when you need race-free selection of ordered handling:
104//!
105//! ```
106//! # use agent_client_protocol::{Client, Agent, ConnectTo};
107//! # use agent_client_protocol_test::MyRequest;
108//! # async fn example(transport: impl ConnectTo<Client>) -> Result<(), agent_client_protocol::Error> {
109//! # Client.builder().connect_with(transport, async |cx| {
110//! cx.prepare_request(MyRequest {})
111//!     .on_receiving_result(async |result| {
112//!         // A peer response routed in its original dispatch holds this barrier
113//!         let response = result?;
114//!         // Do something with response...
115//!         Ok(())
116//!     })?;
117//! # Ok(())
118//! # }).await?;
119//! # Ok(())
120//! # }
121//! ```
122//!
123//! Preparation does not send the request. Calling the callback method installs
124//! ordered handling, then publishes it. The dispatch loop waits for your
125//! callback before processing the next message. A pending-request failure
126//! delivered without an incoming response (such as EOF), or one that an
127//! interceptor retains and routes after its original dispatch, does not carry
128//! the barrier. With eager sending, an already-routed response has no barrier
129//! either.
130//!
131//! `prepare_request(...).block_task()` and `.detach()` send synchronously when
132//! called but do not select ordering. An unconsumed prepared request can be
133//! retained or dropped without sending traffic or holding dispatch.
134//!
135//! An ordered callback must not wait for later inbound traffic on the same
136//! connection. Spawn that follow-up work and return from the callback so the
137//! loop can continue.
138//!
139//! # Escaping the Loop: `spawn`
140//!
141//! Use [`spawn`] to run work outside the dispatch loop:
142//!
143//! ```ignore
144//! builder.on_receive_request(async |request: MyRequest, responder, cx| {
145//!     cx.spawn(async move {
146//!         // This runs outside the loop - other messages may be processed
147//!         let response = cx.send_request(SomeRequest { ... })
148//!             .block_task()
149//!             .await?;
150//!         // ...
151//!         Ok(())
152//!     })?;
153//!     responder.respond(MyResponse { ... })  // Return immediately
154//! }, on_receive_request!());
155//! ```
156//!
157//! # Blocking Session Methods
158//!
159//! `SessionBuilder::run_until` does not spawn its caller. It waits for the
160//! session response on the current task, so call it only when that task already
161//! runs outside the dispatch loop—for example, from the foreground future
162//! passed to `connect_with` or from a task created with [`spawn`]. Awaiting it
163//! from a message handler deadlocks just like awaiting `block_task()` there.
164//!
165//! ```
166//! # use agent_client_protocol::{Client, Agent, ConnectTo};
167//! # async fn example(transport: impl ConnectTo<Client>) -> Result<(), agent_client_protocol::Error> {
168//! # Client.builder().connect_with(transport, async |cx| {
169//! cx.build_session_cwd()?
170//!     .block_task()
171//!     .run_until(async |mut session| {
172//!         // Safe: connect_with's foreground future runs outside the dispatch loop
173//!         session.send_prompt("Hello")?;
174//!         let response = session.read_to_string().await?;
175//!         Ok(())
176//!     })
177//!     .await?;
178//! # Ok(())
179//! # }).await?;
180//! # Ok(())
181//! # }
182//! ```
183//!
184//! # Summary
185//!
186//! | Pattern | Blocks Loop? | Use When |
187//! |---------|--------------|----------|
188//! | `on_receive_*` callback | Yes | Handle one incoming message |
189//! | `on_receiving_*` callback | For a timely peer response | Bounded response work that needs ordering |
190//! | `on_session_start` / `on_proxy_session_start` | Setup only | Install session routing, then run session work concurrently |
191//! | `block_task()` | If awaited in a handler | Wait for a response from outside the dispatch loop |
192//! | `spawn(...)` | No | Long-running work, don't need ordering |
193//! | `block_task().run_until(...)` | If called in a handler | Session-scoped work from outside the dispatch loop |
194//!
195//! # Next Steps
196//!
197//! - [Proxies and Conductors](super::proxies) - Building message interceptors
198//!
199//! [`on_receive_request`]: crate::Builder::on_receive_request
200//! [`on_receive_notification`]: crate::Builder::on_receive_notification
201//! [`on_receiving_result`]: crate::SentRequest::on_receiving_result
202//! [`on_receiving_ok_result`]: crate::SentRequest::on_receiving_ok_result
203//! [`on_session_start`]: crate::SessionBuilder::on_session_start
204//! [`on_proxy_session_start`]: crate::SessionBuilder::on_proxy_session_start
205//! [`SentRequest`]: crate::SentRequest
206//! [`PreparedRequest::on_receiving_result`]: crate::PreparedRequest::on_receiving_result
207//! [`prepare_request`]: crate::ConnectionTo::prepare_request
208//! [`spawn`]: crate::ConnectionTo::spawn