Skip to main content

agent_client_protocol/
session.rs

1use std::{future::Future, marker::PhantomData, path::Path, sync::Arc};
2
3use futures::channel::{mpsc, oneshot};
4use futures::future::{self, Either};
5
6use crate::{
7    Agent, Client, ConnectionTo, Dispatch, HandleDispatchFrom, Handled, JsonRpcRequest, Responder,
8    Role,
9    jsonrpc::{
10        DynamicHandlerCleanup, DynamicHandlerGuard,
11        run::{NullRun, RunWithConnectionTo, RunnerErrorScope},
12    },
13    role::{HasPeer, acp::ProxySessionMessages},
14    schema::v1::{
15        ContentBlock, ContentChunk, LoadSessionRequest, LoadSessionResponse, Meta,
16        NewSessionRequest, NewSessionResponse, PromptRequest, PromptResponse, ResumeSessionRequest,
17        ResumeSessionResponse, SessionConfigOption, SessionId, SessionModeState,
18        SessionNotification, SessionUpdate, StopReason,
19    },
20    util::{MatchDispatch, MatchDispatchFrom},
21};
22
23#[cfg(feature = "unstable_mcp_over_acp")]
24use crate::{jsonrpc::run::ChainRun, mcp_server::McpServer};
25
26#[cfg(feature = "unstable_protocol_v2")]
27mod v2;
28#[cfg(feature = "unstable_protocol_v2")]
29pub use v2::*;
30
31type SessionCleanup = Vec<Arc<dyn DynamicHandlerCleanup>>;
32
33fn session_cleanup<R: Role>(guards: &[DynamicHandlerGuard<R>]) -> SessionCleanup {
34    guards
35        .iter()
36        .filter_map(DynamicHandlerGuard::cleanup)
37        .collect()
38}
39
40fn close_session_registrations(cleanup: &SessionCleanup) {
41    for registration in cleanup {
42        registration.close();
43    }
44}
45
46async fn wait_session_cleanup(cleanup: SessionCleanup) {
47    future::join_all(cleanup.iter().map(|registration| registration.wait())).await;
48}
49
50fn session_runner_connection<R: Role>(
51    connection: ConnectionTo<R>,
52    cleanup: &SessionCleanup,
53) -> (ConnectionTo<R>, RunnerErrorScope) {
54    let closing = cleanup.clone();
55    let scope = RunnerErrorScope::new(
56        move || close_session_registrations(&closing),
57        wait_session_cleanup(cleanup.clone()),
58    );
59    (connection.with_runner_error_scope(scope.clone()), scope)
60}
61
62/// Keep the actual (possibly borrowed) runner alive until only these
63/// registrations finish. Never seal connection-wide protected admission here.
64async fn drive_session_cleanup(
65    run: impl Future<Output = Result<(), crate::Error>>,
66    cleanup: SessionCleanup,
67) -> Result<(), crate::Error> {
68    match future::select(
69        Box::pin(wait_session_cleanup(cleanup.clone())),
70        Box::pin(run),
71    )
72    .await
73    {
74        Either::Left(((), _run)) => Ok(()),
75        Either::Right((result, waiting)) => {
76            if result.is_err() {
77                close_session_registrations(&cleanup);
78            }
79            waiting.await;
80            result
81        }
82    }
83}
84
85async fn run_attached_session_runner(
86    run: impl Future<Output = Result<(), crate::Error>>,
87    cleanup: SessionCleanup,
88    error_scope: RunnerErrorScope,
89) -> Result<(), crate::Error> {
90    let result = if cleanup.is_empty() {
91        run.await
92    } else {
93        drive_session_cleanup(run, cleanup).await
94    };
95    // Local cleanup can win the race against the retained chain's final poll.
96    // Once attached, an observed runner failure must still reach the task actor.
97    error_scope.error().map_or(result, Err)
98}
99
100async fn run_session_scope<T>(
101    run: impl Future<Output = Result<(), crate::Error>>,
102    op: impl Future<Output = Result<T, crate::Error>>,
103    cleanup: SessionCleanup,
104    error_scope: RunnerErrorScope,
105) -> Result<T, crate::Error> {
106    let result = match future::select(Box::pin(run), Box::pin(op)).await {
107        Either::Left((run_result, op)) => {
108            let result = match run_result {
109                Ok(()) => op.await,
110                Err(error) => {
111                    drop(op);
112                    Err(error)
113                }
114            };
115            close_session_registrations(&cleanup);
116            wait_session_cleanup(cleanup).await;
117            result
118        }
119        Either::Right((result, run)) => {
120            close_session_registrations(&cleanup);
121            let runner_result = drive_session_cleanup(run, cleanup).await;
122            // Foreground failure remains authoritative over cleanup failures.
123            result.and_then(|value| runner_result.map(|()| value))
124        }
125    };
126    // The chain may have observed an error while still Pending for cleanup.
127    // Successful foreground completion must not hide that recorded failure.
128    // An explicit foreground error remains authoritative.
129    result.and_then(|value| error_scope.error().map_or(Ok(value), Err))
130}
131
132/// Marker type indicating the session builder will block the current task.
133#[derive(Debug)]
134pub struct Blocking;
135impl SessionBlockState for Blocking {}
136
137/// Marker type indicating the session builder will not block the current task.
138#[derive(Debug)]
139pub struct NonBlocking;
140impl SessionBlockState for NonBlocking {}
141
142/// Trait for marker types that indicate blocking vs non-blocking API.
143/// See [`SessionBuilder::block_task`].
144pub trait SessionBlockState: Send + 'static + Sync + std::fmt::Debug {}
145
146impl<Counterpart: Role> ConnectionTo<Counterpart>
147where
148    Counterpart: HasPeer<Agent>,
149{
150    /// Stable protocol v1 session builder for a new session request.
151    ///
152    /// With `unstable_protocol_v2`, a `Client.v2()` callback receives a
153    /// `V2ConnectionTo` with its own v2 `build_session` helper.
154    pub fn build_session(&self, cwd: impl AsRef<Path>) -> SessionBuilder<Counterpart, NullRun> {
155        SessionBuilder::new(self, NewSessionRequest::new(cwd.as_ref()))
156    }
157
158    /// Stable protocol v1 session builder using the current working directory.
159    ///
160    /// This is a convenience wrapper around [`build_session`](Self::build_session)
161    /// that uses [`std::env::current_dir`] to get the working directory.
162    ///
163    /// Returns an error if the current directory cannot be determined.
164    pub fn build_session_cwd(&self) -> Result<SessionBuilder<Counterpart, NullRun>, crate::Error> {
165        let cwd = std::env::current_dir().map_err(|e| {
166            crate::Error::internal_error().data(format!("cannot get current directory: {e}"))
167        })?;
168        Ok(self.build_session(cwd))
169    }
170
171    /// Stable protocol v1 session builder starting from an existing request.
172    ///
173    /// Use this when you've intercepted a `session.new` request and want to
174    /// modify it (e.g., inject MCP servers) before forwarding.
175    pub fn build_session_from(
176        &self,
177        request: NewSessionRequest,
178    ) -> SessionBuilder<Counterpart, NullRun> {
179        SessionBuilder::new(self, request)
180    }
181
182    /// Stable protocol v1 session builder that loads an existing session.
183    ///
184    /// The returned builder installs session routing before publishing
185    /// `session/load`, so replay notifications sent before the response are
186    /// available through the restored [`ActiveSession`].
187    ///
188    /// Call this only when the initialization response advertises
189    /// `agentCapabilities.loadSession`.
190    pub fn load_session(
191        &self,
192        session_id: impl Into<SessionId>,
193        cwd: impl AsRef<Path>,
194    ) -> RestoreSessionBuilder<Counterpart, LoadSessionRequest> {
195        self.load_session_from(LoadSessionRequest::new(session_id, cwd.as_ref()))
196    }
197
198    /// Stable protocol v1 session builder from an existing `session/load`
199    /// request.
200    ///
201    /// Use this to send a typed request assembled or intercepted elsewhere
202    /// without rebuilding it.
203    pub fn load_session_from(
204        &self,
205        request: LoadSessionRequest,
206    ) -> RestoreSessionBuilder<Counterpart, LoadSessionRequest> {
207        RestoreSessionBuilder::new(self, request)
208    }
209
210    /// Stable protocol v1 session builder that resumes an existing session.
211    ///
212    /// This is the `session/resume` counterpart of
213    /// [`load_session`](Self::load_session), but continues without replaying
214    /// conversation history. Call this only when the initialization response
215    /// advertises `agentCapabilities.sessionCapabilities.resume`.
216    pub fn resume_session(
217        &self,
218        session_id: impl Into<SessionId>,
219        cwd: impl AsRef<Path>,
220    ) -> RestoreSessionBuilder<Counterpart, ResumeSessionRequest> {
221        self.resume_session_from(ResumeSessionRequest::new(session_id, cwd.as_ref()))
222    }
223
224    /// Stable protocol v1 session builder from an existing `session/resume`
225    /// request.
226    ///
227    /// Use this to send a typed request assembled or intercepted elsewhere
228    /// without rebuilding it.
229    pub fn resume_session_from(
230        &self,
231        request: ResumeSessionRequest,
232    ) -> RestoreSessionBuilder<Counterpart, ResumeSessionRequest> {
233        RestoreSessionBuilder::new(self, request)
234    }
235
236    /// Given a session response received from the agent,
237    /// attach a handler to process messages related to this session
238    /// and let you access them.
239    ///
240    /// Normally you would not use this method directly but would
241    /// instead use [`Self::build_session`] and then [`SessionBuilder::start_session`].
242    ///
243    /// The vector `dynamic_handler_registrations` contains any dynamic
244    /// handle registrations associated with this session (e.g., from MCP servers).
245    /// You can simply pass `Default::default()` if not applicable.
246    pub(crate) fn attach_session<'runner>(
247        &self,
248        response: NewSessionResponse,
249        mcp_handler_registrations: Vec<DynamicHandlerGuard<Counterpart>>,
250    ) -> Result<ActiveSession<'runner, Counterpart>, crate::Error> {
251        let NewSessionResponse {
252            session_id,
253            modes,
254            config_options,
255            meta,
256            ..
257        } = response;
258
259        let prepared = self.prepare_session_routing(&session_id)?;
260        Ok(prepared.into_active_session(
261            self.clone(),
262            session_id,
263            modes,
264            config_options,
265            meta,
266            mcp_handler_registrations,
267        ))
268    }
269
270    /// Install the update channel and handler for `session_id`.
271    ///
272    /// Restore requests call this before request publication. Dropping the
273    /// returned value deactivates and removes the route.
274    fn prepare_session_routing(
275        &self,
276        session_id: &SessionId,
277    ) -> Result<PreparedSession<Counterpart>, crate::Error> {
278        let (update_tx, update_rx) = mpsc::unbounded();
279        let handler = ActiveSessionHandler::new(session_id.clone(), update_tx.clone());
280        let session_handler_registration = self.add_dynamic_handler(handler)?;
281
282        Ok(PreparedSession {
283            update_rx,
284            update_tx,
285            session_handler_registration,
286        })
287    }
288}
289
290/// Session-routing state installed before a restore request is published.
291struct PreparedSession<Counterpart: Role>
292where
293    Counterpart: HasPeer<Agent>,
294{
295    update_rx: mpsc::UnboundedReceiver<SessionMessage>,
296    update_tx: mpsc::UnboundedSender<SessionMessage>,
297    session_handler_registration: DynamicHandlerGuard<Counterpart>,
298}
299
300impl<Counterpart> PreparedSession<Counterpart>
301where
302    Counterpart: HasPeer<Agent>,
303{
304    fn into_active_session<'runner>(
305        self,
306        connection: ConnectionTo<Counterpart>,
307        session_id: SessionId,
308        modes: Option<SessionModeState>,
309        config_options: Option<Vec<SessionConfigOption>>,
310        meta: Option<Meta>,
311        mcp_handler_registrations: Vec<DynamicHandlerGuard<Counterpart>>,
312    ) -> ActiveSession<'runner, Counterpart> {
313        ActiveSession {
314            session_id,
315            modes,
316            config_options,
317            meta,
318            update_rx: self.update_rx,
319            update_tx: self.update_tx,
320            connection,
321            session_handler_registration: self.session_handler_registration,
322            mcp_handler_registrations,
323            _runner: PhantomData,
324        }
325    }
326}
327
328/// Internal behavior shared by the two stable restore operations.
329trait RestoreRequest: JsonRpcRequest {
330    fn session_id(&self) -> &SessionId;
331    fn response_modes(response: &Self::Response) -> Option<SessionModeState>;
332    fn response_config_options(response: &Self::Response) -> Option<Vec<SessionConfigOption>>;
333    fn response_meta(response: &Self::Response) -> Option<Meta>;
334}
335
336impl RestoreRequest for LoadSessionRequest {
337    fn session_id(&self) -> &SessionId {
338        &self.session_id
339    }
340
341    fn response_modes(response: &Self::Response) -> Option<SessionModeState> {
342        response.modes.clone()
343    }
344
345    fn response_config_options(response: &Self::Response) -> Option<Vec<SessionConfigOption>> {
346        response.config_options.clone()
347    }
348
349    fn response_meta(response: &Self::Response) -> Option<Meta> {
350        response.meta.clone()
351    }
352}
353
354impl RestoreRequest for ResumeSessionRequest {
355    fn session_id(&self) -> &SessionId {
356        &self.session_id
357    }
358
359    fn response_modes(response: &Self::Response) -> Option<SessionModeState> {
360        response.modes.clone()
361    }
362
363    fn response_config_options(response: &Self::Response) -> Option<Vec<SessionConfigOption>> {
364        response.config_options.clone()
365    }
366
367    fn response_meta(response: &Self::Response) -> Option<Meta> {
368        response.meta.clone()
369    }
370}
371
372/// Stable protocol v1 builder for `session/load` or `session/resume`.
373///
374/// Use [`ConnectionTo::load_session`] or [`ConnectionTo::resume_session`] to
375/// construct this builder. Use the matching `_from` method to send an existing
376/// typed request without rebuilding it.
377///
378/// The `BlockState` parameter mirrors [`SessionBuilder`]:
379/// - [`NonBlocking`] exposes `on_session_start` on each concrete operation.
380/// - [`Blocking`], selected with [`Self::block_task`], exposes
381///   `start_session`.
382///
383/// Session routing is acknowledged before the request can reach the peer.
384/// Dropping a pending blocking start removes that routing and applies the
385/// standard [`SentRequest`](crate::SentRequest) drop-time cancellation
386/// behavior. Error responses remove the route before later entries in the same
387/// transport frame are dispatched.
388#[must_use = "use `start_session` or `on_session_start` to restore the session"]
389#[derive(Debug)]
390pub struct RestoreSessionBuilder<Counterpart, Request, BlockState = NonBlocking>
391where
392    Counterpart: HasPeer<Agent>,
393    BlockState: SessionBlockState,
394{
395    connection: ConnectionTo<Counterpart>,
396    request: Request,
397    block_state: PhantomData<BlockState>,
398}
399
400impl<Counterpart, Request> RestoreSessionBuilder<Counterpart, Request, NonBlocking>
401where
402    Counterpart: HasPeer<Agent>,
403{
404    fn new(connection: &ConnectionTo<Counterpart>, request: Request) -> Self {
405        Self {
406            connection: connection.clone(),
407            request,
408            block_state: PhantomData,
409        }
410    }
411
412    /// Mark this restore builder as able to block the current task.
413    ///
414    /// Do not use the resulting blocking methods inside a message handler.
415    pub fn block_task(self) -> RestoreSessionBuilder<Counterpart, Request, Blocking> {
416        RestoreSessionBuilder {
417            connection: self.connection,
418            request: self.request,
419            block_state: PhantomData,
420        }
421    }
422}
423
424fn restored_session<Counterpart, Request>(
425    connection: ConnectionTo<Counterpart>,
426    session_id: SessionId,
427    prepared: PreparedSession<Counterpart>,
428    response: Request::Response,
429) -> RestoredSession<'static, Counterpart, Request::Response>
430where
431    Counterpart: HasPeer<Agent>,
432    Request: RestoreRequest,
433{
434    let session = prepared.into_active_session(
435        connection,
436        session_id,
437        Request::response_modes(&response),
438        Request::response_config_options(&response),
439        Request::response_meta(&response),
440        Vec::new(),
441    );
442
443    RestoredSession { session, response }
444}
445
446fn on_restore_session_start<Counterpart, Request, F, Fut>(
447    builder: RestoreSessionBuilder<Counterpart, Request>,
448    op: F,
449) -> Result<(), crate::Error>
450where
451    Counterpart: HasPeer<Agent>,
452    Request: RestoreRequest,
453    F: FnOnce(RestoredSession<'static, Counterpart, Request::Response>) -> Fut + Send + 'static,
454    Fut: Future<Output = Result<(), crate::Error>> + Send,
455{
456    ensure_v1_session_protocol(&builder.connection)?;
457
458    let RestoreSessionBuilder {
459        connection,
460        request,
461        block_state: _,
462    } = builder;
463    let session_id = request.session_id().clone();
464    let prepared = connection.prepare_session_routing(&session_id)?;
465    let routing_ready = connection.dynamic_handler_barrier();
466
467    connection
468        .send_ordered_request_to_after(Agent, request, routing_ready)
469        .on_receiving_result({
470            let connection = connection.clone();
471            async move |result| {
472                let response = result?;
473                let restored = restored_session::<_, Request>(
474                    connection.clone(),
475                    session_id,
476                    prepared,
477                    response,
478                );
479                connection.spawn(async move { op(restored).await })
480            }
481        })
482}
483
484async fn start_restored_session<Counterpart, Request>(
485    builder: RestoreSessionBuilder<Counterpart, Request, Blocking>,
486) -> Result<RestoredSession<'static, Counterpart, Request::Response>, crate::Error>
487where
488    Counterpart: HasPeer<Agent>,
489    Request: RestoreRequest,
490{
491    ensure_v1_session_protocol(&builder.connection)?;
492
493    let RestoreSessionBuilder {
494        connection,
495        request,
496        block_state: _,
497    } = builder;
498    let session_id = request.session_id().clone();
499    let prepared = connection.prepare_session_routing(&session_id)?;
500    let routing_ready = connection.dynamic_handler_barrier();
501    let session_connection = connection.clone();
502
503    connection
504        .send_ordered_request_to_after(Agent, request, routing_ready)
505        .block_task_with_ordered_result(move |result| {
506            let response = result?;
507            Ok(restored_session::<_, Request>(
508                session_connection,
509                session_id,
510                prepared,
511                response,
512            ))
513        })
514        .await
515}
516
517impl<Counterpart> RestoreSessionBuilder<Counterpart, LoadSessionRequest>
518where
519    Counterpart: HasPeer<Agent>,
520{
521    /// Restore with `session/load` in the background and run `op` once its
522    /// exact response and active session are available.
523    ///
524    /// This returns immediately and is safe to call from a message handler.
525    /// Replay notifications can arrive before the response and are retained by
526    /// the returned session.
527    pub fn on_session_start<F, Fut>(self, op: F) -> Result<(), crate::Error>
528    where
529        F: FnOnce(RestoredSession<'static, Counterpart, LoadSessionResponse>) -> Fut
530            + Send
531            + 'static,
532        Fut: Future<Output = Result<(), crate::Error>> + Send,
533    {
534        on_restore_session_start(self, op)
535    }
536}
537
538impl<Counterpart> RestoreSessionBuilder<Counterpart, ResumeSessionRequest>
539where
540    Counterpart: HasPeer<Agent>,
541{
542    /// Restore with `session/resume` in the background and run `op` once its
543    /// exact response and active session are available.
544    ///
545    /// This returns immediately and is safe to call from a message handler.
546    /// The returned session receives subsequent session traffic.
547    pub fn on_session_start<F, Fut>(self, op: F) -> Result<(), crate::Error>
548    where
549        F: FnOnce(RestoredSession<'static, Counterpart, ResumeSessionResponse>) -> Fut
550            + Send
551            + 'static,
552        Fut: Future<Output = Result<(), crate::Error>> + Send,
553    {
554        on_restore_session_start(self, op)
555    }
556}
557
558impl<Counterpart> RestoreSessionBuilder<Counterpart, LoadSessionRequest, Blocking>
559where
560    Counterpart: HasPeer<Agent>,
561{
562    /// Publish `session/load`, wait on the current task, and return an
563    /// [`ActiveSession`] together with the exact [`LoadSessionResponse`].
564    ///
565    /// Requires [`block_task`](RestoreSessionBuilder::block_task). Dropping
566    /// this future while it is pending cancels the request and removes the
567    /// provisional session route.
568    pub async fn start_session(
569        self,
570    ) -> Result<RestoredSession<'static, Counterpart, LoadSessionResponse>, crate::Error> {
571        start_restored_session(self).await
572    }
573}
574
575impl<Counterpart> RestoreSessionBuilder<Counterpart, ResumeSessionRequest, Blocking>
576where
577    Counterpart: HasPeer<Agent>,
578{
579    /// Publish `session/resume`, wait on the current task, and return an
580    /// [`ActiveSession`] together with the exact [`ResumeSessionResponse`].
581    ///
582    /// Requires [`block_task`](RestoreSessionBuilder::block_task). Dropping
583    /// this future while it is pending cancels the request and removes the
584    /// provisional session route.
585    pub async fn start_session(
586        self,
587    ) -> Result<RestoredSession<'static, Counterpart, ResumeSessionResponse>, crate::Error> {
588        start_restored_session(self).await
589    }
590}
591
592/// A restored stable-v1 session and the exact operation response that opened
593/// it.
594///
595/// The session ID comes from the load or resume request because stable-v1
596/// restore responses do not repeat it. Keeping the response separate preserves
597/// every operation-specific field without reconstructing it from session
598/// state.
599pub struct RestoredSession<'runner, Link, Response>
600where
601    Link: HasPeer<Agent>,
602{
603    session: ActiveSession<'runner, Link>,
604    response: Response,
605}
606
607impl<'runner, Link, Response> RestoredSession<'runner, Link, Response>
608where
609    Link: HasPeer<Agent>,
610{
611    /// Access the active session.
612    pub fn session(&self) -> &ActiveSession<'runner, Link> {
613        &self.session
614    }
615
616    /// Mutably access the active session, for example to consume replay.
617    pub fn session_mut(&mut self) -> &mut ActiveSession<'runner, Link> {
618        &mut self.session
619    }
620
621    /// Access the complete load or resume response.
622    pub fn response(&self) -> &Response {
623        &self.response
624    }
625
626    /// Split the restored value into its active session and exact response.
627    pub fn into_parts(self) -> (ActiveSession<'runner, Link>, Response) {
628        (self.session, self.response)
629    }
630
631    /// Consume this value and return only the active session.
632    pub fn into_session(self) -> ActiveSession<'runner, Link> {
633        self.session
634    }
635}
636
637impl<Link, Response> std::fmt::Debug for RestoredSession<'_, Link, Response>
638where
639    Link: HasPeer<Agent>,
640    Response: std::fmt::Debug,
641{
642    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
643        formatter
644            .debug_struct("RestoredSession")
645            .field("session_id", self.session.session_id())
646            .field("response", &self.response)
647            .finish()
648    }
649}
650
651/// Stable protocol v1 session builder for a new session request.
652/// Allows you to add MCP servers or set other details for this session.
653///
654/// The `BlockState` type parameter tracks whether blocking methods are available:
655/// - `NonBlocking` (default): Only [`on_session_start`](Self::on_session_start) is available
656/// - `Blocking` (after calling [`block_task`](Self::block_task)):
657///   [`run_until`](Self::run_until) and [`start_session`](Self::start_session) become available
658#[must_use = "use `start_session`, `run_until`, or `on_session_start` to start the session"]
659#[derive(Debug)]
660pub struct SessionBuilder<
661    Counterpart,
662    Run: RunWithConnectionTo<Counterpart> = NullRun,
663    BlockState: SessionBlockState = NonBlocking,
664> where
665    Counterpart: HasPeer<Agent>,
666{
667    connection: ConnectionTo<Counterpart>,
668    request: NewSessionRequest,
669    dynamic_handler_registrations: Vec<DynamicHandlerGuard<Counterpart>>,
670    run: Run,
671    block_state: PhantomData<BlockState>,
672}
673
674impl<Counterpart> SessionBuilder<Counterpart, NullRun, NonBlocking>
675where
676    Counterpart: HasPeer<Agent>,
677{
678    fn new(connection: &ConnectionTo<Counterpart>, request: NewSessionRequest) -> Self {
679        SessionBuilder {
680            connection: connection.clone(),
681            request,
682            dynamic_handler_registrations: Vec::default(),
683            run: NullRun,
684            block_state: PhantomData,
685        }
686    }
687}
688
689impl<Counterpart, R, BlockState> SessionBuilder<Counterpart, R, BlockState>
690where
691    Counterpart: HasPeer<Agent>,
692    R: RunWithConnectionTo<Counterpart>,
693    BlockState: SessionBlockState,
694{
695    /// Attach an MCP server to this new session.
696    #[cfg(feature = "unstable_mcp_over_acp")]
697    pub fn with_mcp_server<McpRun>(
698        mut self,
699        mcp_server: McpServer<Counterpart, McpRun>,
700    ) -> Result<SessionBuilder<Counterpart, ChainRun<R, McpRun>, BlockState>, crate::Error>
701    where
702        McpRun: RunWithConnectionTo<Counterpart>,
703    {
704        let (handler, mcp_run) = mcp_server.into_handler_and_runner();
705        self.dynamic_handler_registrations
706            .push(handler.into_dynamic_handler(&mut self.request, &self.connection)?);
707        Ok(SessionBuilder {
708            connection: self.connection,
709            request: self.request,
710            dynamic_handler_registrations: self.dynamic_handler_registrations,
711            run: ChainRun::new(self.run, mcp_run),
712            block_state: self.block_state,
713        })
714    }
715
716    /// Spawn a task that runs the provided closure once the session starts.
717    ///
718    /// Unlike [`start_session`](Self::start_session), this method returns immediately
719    /// without blocking the current task. The session handshake and closure execution
720    /// happen in a spawned background task.
721    ///
722    /// The closure receives an `ActiveSession<'static, _>` and runs in a
723    /// spawned task. If it returns an error, the error propagates to the
724    /// connection's task handling.
725    ///
726    /// # Example
727    ///
728    /// ```ignore
729    /// # use agent_client_protocol::{Client, Agent, ConnectTo};
730    /// # use agent_client_protocol::mcp_server::McpServer;
731    /// # use agent_client_protocol_rmcp::McpServerExt;
732    /// # async fn example(transport: impl ConnectTo<Client>) -> Result<(), agent_client_protocol::Error> {
733    /// # Client.builder().connect_with(transport, async |cx| {
734    /// # let mcp = McpServer::<Agent, _>::builder("tools").build();
735    /// cx.build_session_cwd()?
736    ///     .with_mcp_server(mcp)?
737    ///     .on_session_start(async |mut session| {
738    ///         // Do something with the session
739    ///         session.send_prompt("Hello")?;
740    ///         let response = session.read_to_string().await?;
741    ///         Ok(())
742    ///     })?;
743    /// // Returns immediately, session runs in background
744    /// # Ok(())
745    /// # }).await?;
746    /// # Ok(())
747    /// # }
748    /// ```
749    ///
750    /// # Ordering
751    ///
752    /// Session runners are scheduled and routing setup is installed before the
753    /// dispatch loop processes the next message when the session response is
754    /// routed during its original dispatch. No user callback code runs under
755    /// that ordering guarantee: the callback is invoked in a spawned task, so
756    /// it may wait for later session traffic without deadlocking the connection.
757    /// A response interceptor that retains the response and routes it later
758    /// cannot retroactively order session setup before messages the dispatch
759    /// loop has already processed.
760    pub fn on_session_start<F, Fut>(self, op: F) -> Result<(), crate::Error>
761    where
762        R: 'static,
763        F: FnOnce(ActiveSession<'static, Counterpart>) -> Fut + Send + 'static,
764        Fut: Future<Output = Result<(), crate::Error>> + Send,
765    {
766        ensure_v1_session_protocol(&self.connection)?;
767
768        let Self {
769            connection,
770            request,
771            dynamic_handler_registrations,
772            run,
773            block_state: _,
774        } = self;
775
776        let cleanup = session_cleanup(&dynamic_handler_registrations);
777        let (runner_connection, error_scope) =
778            session_runner_connection(connection.clone(), &cleanup);
779        connection.spawn(run_attached_session_runner(
780            run.run_with_connection_to(runner_connection),
781            cleanup,
782            error_scope,
783        ))?;
784
785        connection
786            .send_ordered_request_to(Agent, request)
787            .on_receiving_result({
788                let connection = connection.clone();
789                async move |result| {
790                    let response = result?;
791
792                    let active_session =
793                        connection.attach_session(response, dynamic_handler_registrations)?;
794
795                    connection.spawn(async move { op(active_session).await })
796                }
797            })
798    }
799
800    /// Spawn a proxy session and run a closure with the session ID.
801    ///
802    /// A **proxy session** starts the session with the agent and then automatically
803    /// proxies all session updates (prompts, tool calls, etc.) from the agent back
804    /// to the client. You don't need to handle any messages yourself - the proxy
805    /// takes care of forwarding everything. This is useful when you want to inject
806    /// and/or filter prompts coming from the client but otherwise not be involved
807    /// in the session.
808    ///
809    /// Unlike [`start_session_proxy`](Self::start_session_proxy), this method returns
810    /// immediately without blocking the current task. The session handshake, client
811    /// response, and proxy setup all happen in a spawned background task.
812    ///
813    /// The closure receives the `SessionId` once the session is established. Use it for logging
814    /// or eventual tracking; it runs concurrently with later connection traffic. Register
815    /// ID-independent state that later handlers must observe before calling this helper. For
816    /// ID-keyed bookkeeping, install a gate or placeholder first, make later handlers await it,
817    /// and populate it from the closure.
818    ///
819    /// # Example
820    ///
821    /// ```ignore
822    /// # use agent_client_protocol::{Proxy, Client, Conductor, ConnectTo};
823    /// # use agent_client_protocol::schema::v1::NewSessionRequest;
824    /// # use agent_client_protocol::mcp_server::McpServer;
825    /// # use agent_client_protocol_rmcp::McpServerExt;
826    /// # async fn example(transport: impl ConnectTo<Proxy>) -> Result<(), agent_client_protocol::Error> {
827    /// Proxy.builder()
828    ///     .on_receive_request_from(Client, async |request: NewSessionRequest, responder, cx| {
829    ///         let mcp = McpServer::<Conductor, _>::builder("tools").build();
830    ///         cx.build_session_from(request)
831    ///             .with_mcp_server(mcp)?
832    ///             .on_proxy_session_start(responder, async |session_id| {
833    ///                 // Session started
834    ///                 Ok(())
835    ///             })
836    ///     }, agent_client_protocol::on_receive_request!())
837    ///     .connect_to(transport)
838    ///     .await?;
839    /// # Ok(())
840    /// # }
841    /// ```
842    ///
843    /// # Ordering
844    ///
845    /// The client response is queued, proxy routing is installed, and session runners are
846    /// scheduled before the dispatch loop processes the next message when the session response
847    /// is routed during its original dispatch. This is a local ordering guarantee, not a
848    /// guarantee that the response reaches the client before later wire traffic. No user callback
849    /// code runs under the barrier: the callback is invoked in a spawned task, so it may wait for
850    /// later connection traffic. A response interceptor that retains the response and routes it
851    /// later cannot retroactively order this setup before messages the loop already processed.
852    pub fn on_proxy_session_start<F, Fut>(
853        self,
854        responder: Responder<NewSessionResponse>,
855        op: F,
856    ) -> Result<(), crate::Error>
857    where
858        F: FnOnce(SessionId) -> Fut + Send + 'static,
859        Fut: Future<Output = Result<(), crate::Error>> + Send,
860        Counterpart: HasPeer<Client>,
861        R: 'static,
862    {
863        ensure_v1_session_protocol(&self.connection)?;
864
865        let Self {
866            connection,
867            request,
868            dynamic_handler_registrations,
869            run,
870            block_state: _,
871        } = self;
872
873        let cleanup = session_cleanup(&dynamic_handler_registrations);
874        let (runner_connection, error_scope) =
875            session_runner_connection(connection.clone(), &cleanup);
876        connection.spawn(run_attached_session_runner(
877            run.run_with_connection_to(runner_connection),
878            cleanup,
879            error_scope,
880        ))?;
881
882        // Send the "new session" request to the agent.
883        let sent = connection.send_ordered_request_to(Agent, request);
884        let sent = sent.forward_cancellation_from(responder.cancellation());
885
886        sent.on_receiving_ok_result(responder, {
887            let connection = connection.clone();
888            async move |response, responder| {
889                // Extract the session-id from the response and forward
890                // the response back to the client
891                let session_id = response.session_id.clone();
892                responder.respond(response)?;
893
894                // Install a dynamic handler to proxy messages from this session
895                connection
896                    .add_dynamic_handler(ProxySessionMessages::new(session_id.clone()))?
897                    .detach();
898
899                // Keep dynamic handlers live for the connection.
900                dynamic_handler_registrations
901                    .into_iter()
902                    .for_each(DynamicHandlerGuard::detach);
903
904                connection.spawn(async move { op(session_id).await })
905            }
906        })
907    }
908}
909
910impl<Counterpart, R> SessionBuilder<Counterpart, R, NonBlocking>
911where
912    Counterpart: HasPeer<Agent>,
913    R: RunWithConnectionTo<Counterpart>,
914{
915    /// Mark this session builder as being able to block the current task.
916    ///
917    /// After calling this, you can use [`run_until`](Self::run_until) or
918    /// [`start_session`](Self::start_session) which block the current task.
919    ///
920    /// This should not be used from inside a message handler like
921    /// [`Builder::on_receive_request`](`crate::Builder::on_receive_request`) or [`HandleDispatchFrom`]
922    /// implementations.
923    pub fn block_task(self) -> SessionBuilder<Counterpart, R, Blocking> {
924        SessionBuilder {
925            connection: self.connection,
926            request: self.request,
927            dynamic_handler_registrations: self.dynamic_handler_registrations,
928            run: self.run,
929            block_state: PhantomData,
930        }
931    }
932}
933
934impl<Counterpart, R> SessionBuilder<Counterpart, R, Blocking>
935where
936    Counterpart: HasPeer<Agent>,
937    R: RunWithConnectionTo<Counterpart>,
938{
939    /// Run this session synchronously. The current task will be blocked
940    /// and `op` will be executed with the active session information.
941    /// This is useful when you have MCP servers that are borrowed from your local
942    /// stack frame.
943    ///
944    /// The `ActiveSession` passed to `op` has a non-`'static` lifetime, which
945    /// prevents calling [`ActiveSession::proxy_remaining_messages`] (since the
946    /// session's background runners would terminate when `op` returns).
947    ///
948    /// Requires calling [`block_task`](Self::block_task) first.
949    pub async fn run_until<T>(
950        self,
951        op: impl for<'runner> AsyncFnOnce(
952            ActiveSession<'runner, Counterpart>,
953        ) -> Result<T, crate::Error>,
954    ) -> Result<T, crate::Error> {
955        let Self {
956            connection,
957            request,
958            dynamic_handler_registrations,
959            run,
960            block_state: _,
961        } = self;
962
963        let cleanup = session_cleanup(&dynamic_handler_registrations);
964        let (runner_connection, error_scope) =
965            session_runner_connection(connection.clone(), &cleanup);
966        run_session_scope(
967            run.run_with_connection_to(runner_connection),
968            async move {
969                ensure_v1_session_protocol(&connection)?;
970                let response = connection
971                    .send_request_to(Agent, request)
972                    .block_task()
973                    .await?;
974                let active_session =
975                    connection.attach_session(response, dynamic_handler_registrations)?;
976                op(active_session).await
977            },
978            cleanup,
979            error_scope,
980        )
981        .await
982    }
983
984    /// Send the request to create the session and return a handle.
985    /// This is an alternative to [`Self::run_until`] that avoids rightward
986    /// drift but at the cost of requiring MCP servers that are `Send` and
987    /// don't access data from the surrounding scope.
988    ///
989    /// Returns an `ActiveSession<'static, _>` because the session's runners are spawned into
990    /// background tasks that live for the connection lifetime.
991    ///
992    /// Requires calling [`block_task`](Self::block_task) first.
993    pub async fn start_session(self) -> Result<ActiveSession<'static, Counterpart>, crate::Error>
994    where
995        R: 'static,
996    {
997        ensure_v1_session_protocol(&self.connection)?;
998
999        let Self {
1000            connection,
1001            request,
1002            dynamic_handler_registrations,
1003            run,
1004            block_state: _,
1005        } = self;
1006
1007        let (active_session_tx, active_session_rx) = oneshot::channel();
1008
1009        let cleanup = session_cleanup(&dynamic_handler_registrations);
1010        let (runner_connection, error_scope) =
1011            session_runner_connection(connection.clone(), &cleanup);
1012        connection.spawn(run_attached_session_runner(
1013            run.run_with_connection_to(runner_connection),
1014            cleanup,
1015            error_scope,
1016        ))?;
1017
1018        connection.clone().spawn(async move {
1019            let response = connection
1020                .send_request_to(Agent, request)
1021                .block_task()
1022                .await?;
1023
1024            let active_session =
1025                connection.attach_session(response, dynamic_handler_registrations)?;
1026
1027            active_session_tx
1028                .send(active_session)
1029                .map_err(|_| crate::Error::internal_error())?;
1030
1031            Ok(())
1032        })?;
1033
1034        active_session_rx
1035            .await
1036            .map_err(|_| crate::Error::internal_error())
1037    }
1038
1039    /// Start a proxy session that forwards all messages between client and agent.
1040    ///
1041    /// A **proxy session** starts the session with the agent and then automatically
1042    /// proxies all session updates (prompts, tool calls, etc.) from the agent back
1043    /// to the client. You don't need to handle any messages yourself - the proxy
1044    /// takes care of forwarding everything. This is useful when you want to inject
1045    /// and/or filter prompts coming from the client but otherwise not be involved
1046    /// in the session.
1047    ///
1048    /// This is a convenience method that combines [`start_session`](Self::start_session),
1049    /// responding to the client, and [`ActiveSession::proxy_remaining_messages`].
1050    ///
1051    /// For more control (e.g., to send some messages before proxying), use
1052    /// [`start_session`](Self::start_session) instead and call
1053    /// [`proxy_remaining_messages`](ActiveSession::proxy_remaining_messages) manually.
1054    ///
1055    /// Requires calling [`block_task`](Self::block_task) first.
1056    pub async fn start_session_proxy(
1057        self,
1058        responder: Responder<NewSessionResponse>,
1059    ) -> Result<SessionId, crate::Error>
1060    where
1061        Counterpart: HasPeer<Client>,
1062        R: 'static,
1063    {
1064        let active_session = self.start_session().await?;
1065        let session_id = active_session.session_id().clone();
1066        responder.respond(active_session.response())?;
1067        active_session.proxy_remaining_messages()?;
1068        Ok(session_id)
1069    }
1070}
1071
1072/// Stable protocol v1 active session that lets you send prompts and receive updates.
1073///
1074/// The `'runner` lifetime represents the span during which session support runners
1075/// (such as MCP servers) are active. When created via [`SessionBuilder::start_session`],
1076/// this is `'static` because the runners are spawned into background tasks.
1077/// When created via [`SessionBuilder::run_until`], this is tied to the
1078/// closure scope, preventing [`Self::proxy_remaining_messages`] from being called
1079/// (since the runners would stop when the closure returns).
1080#[derive(Debug)]
1081pub struct ActiveSession<'runner, Link>
1082where
1083    Link: HasPeer<Agent>,
1084{
1085    session_id: SessionId,
1086    update_rx: mpsc::UnboundedReceiver<SessionMessage>,
1087    update_tx: mpsc::UnboundedSender<SessionMessage>,
1088    modes: Option<SessionModeState>,
1089    config_options: Option<Vec<SessionConfigOption>>,
1090    meta: Option<serde_json::Map<String, serde_json::Value>>,
1091    connection: ConnectionTo<Link>,
1092
1093    /// Registration for the handler that routes session messages to `update_rx`.
1094    /// This is separate from MCP handlers so it can be dropped independently
1095    /// when switching to proxy mode.
1096    session_handler_registration: DynamicHandlerGuard<Link>,
1097
1098    /// Registrations for MCP server handlers.
1099    /// These will be dropped once the active-session struct is dropped
1100    /// which will cause them to be deregistered.
1101    mcp_handler_registrations: Vec<DynamicHandlerGuard<Link>>,
1102
1103    /// Phantom lifetime representing the session-runner lifetime.
1104    _runner: PhantomData<&'runner ()>,
1105}
1106
1107/// Incoming stable protocol v1 message from the agent.
1108#[non_exhaustive]
1109#[derive(Debug)]
1110#[allow(
1111    clippy::large_enum_variant,
1112    reason = "Dispatch messages vastly outnumber StopReason; boxing would add a heap allocation"
1113)]
1114pub enum SessionMessage {
1115    /// Periodic updates with new content, tool requests, etc.
1116    /// Use [`MatchDispatch`] to match on the message type.
1117    SessionMessage(Dispatch),
1118
1119    /// When a prompt completes, the stop reason.
1120    StopReason(StopReason),
1121}
1122
1123impl<Link> ActiveSession<'_, Link>
1124where
1125    Link: HasPeer<Agent>,
1126{
1127    /// Access the session ID.
1128    pub fn session_id(&self) -> &SessionId {
1129        &self.session_id
1130    }
1131
1132    /// Access modes available in this session.
1133    pub fn modes(&self) -> Option<&SessionModeState> {
1134        self.modes.as_ref()
1135    }
1136
1137    /// Access the initial session configuration options returned by the agent.
1138    pub fn config_options(&self) -> Option<&[SessionConfigOption]> {
1139        self.config_options.as_deref()
1140    }
1141
1142    /// Access meta data from session response.
1143    pub fn meta(&self) -> Option<&serde_json::Map<String, serde_json::Value>> {
1144        self.meta.as_ref()
1145    }
1146
1147    /// Build a `NewSessionResponse` from the session information.
1148    ///
1149    /// Useful when you need to forward the session response to a client
1150    /// after doing some processing.
1151    pub fn response(&self) -> NewSessionResponse {
1152        NewSessionResponse::new(self.session_id.clone())
1153            .modes(self.modes.clone())
1154            .config_options(self.config_options.clone())
1155            .meta(self.meta.clone())
1156    }
1157
1158    /// Access the underlying connection context used to communicate with the agent.
1159    pub fn connection(&self) -> &ConnectionTo<Link> {
1160        &self.connection
1161    }
1162
1163    /// Send a prompt to the agent. You can then read messages sent in response.
1164    pub fn send_prompt(&mut self, prompt: impl ToString) -> Result<(), crate::Error> {
1165        let update_tx = self.update_tx.clone();
1166        self.connection
1167            .send_ordered_request_to(
1168                Agent,
1169                PromptRequest::new(self.session_id.clone(), vec![prompt.to_string().into()]),
1170            )
1171            .on_receiving_result(async move |result| {
1172                let PromptResponse { stop_reason, .. } = result?;
1173
1174                update_tx
1175                    .unbounded_send(SessionMessage::StopReason(stop_reason))
1176                    .map_err(crate::util::internal_error)?;
1177
1178                Ok(())
1179            })
1180    }
1181
1182    /// Read an update from the agent in response to the prompt.
1183    pub async fn read_update(&mut self) -> Result<SessionMessage, crate::Error> {
1184        use futures::StreamExt;
1185        let message =
1186            self.update_rx.next().await.ok_or_else(|| {
1187                crate::util::internal_error("session channel closed unexpectedly")
1188            })?;
1189
1190        Ok(message)
1191    }
1192
1193    /// Read all updates until the end of the turn and create a string.
1194    /// Ignores non-text updates.
1195    pub async fn read_to_string(&mut self) -> Result<String, crate::Error> {
1196        let mut output = String::new();
1197        loop {
1198            let update = self.read_update().await?;
1199            tracing::trace!(?update, "read_to_string update");
1200            match update {
1201                SessionMessage::SessionMessage(dispatch) => MatchDispatch::new(dispatch)
1202                    .if_notification(async |notif: SessionNotification| match notif.update {
1203                        SessionUpdate::AgentMessageChunk(ContentChunk {
1204                            content: ContentBlock::Text(text),
1205                            ..
1206                        }) => {
1207                            output.push_str(&text.text);
1208                            Ok(())
1209                        }
1210                        _ => Ok(()),
1211                    })
1212                    .await
1213                    .otherwise_ignore()?,
1214                SessionMessage::StopReason(_stop_reason) => break,
1215            }
1216        }
1217        Ok(output)
1218    }
1219}
1220
1221impl<Link> ActiveSession<'static, Link>
1222where
1223    Link: HasPeer<Agent>,
1224{
1225    /// Proxy all remaining messages for this session between client and agent.
1226    ///
1227    /// Use this when you want to inject MCP servers into a session but don't need
1228    /// to actively interact with it after setup. The session messages will be proxied
1229    /// between client and agent automatically.
1230    ///
1231    /// This consumes the `ActiveSession` since you're giving up active control.
1232    ///
1233    /// This method is only available on `ActiveSession<'static, _>` (from
1234    /// [`SessionBuilder::start_session`]) because it requires the session's runners to outlive
1235    /// the method call.
1236    ///
1237    /// # Message Ordering Guarantees
1238    ///
1239    /// This method ensures proper handoff from active session mode to proxy mode
1240    /// without losing or reordering messages:
1241    ///
1242    /// 1. **Stop the session handler** - Drop the registration that routes messages
1243    ///    to `update_rx`. After this, no new messages will be queued.
1244    /// 2. **Close the channel** - Drop `update_tx` so we can detect when the channel
1245    ///    is fully drained.
1246    /// 3. **Drain queued messages** - Forward any messages that were already queued
1247    ///    in `update_rx` to the client, preserving order.
1248    /// 4. **Install proxy handler** - Now that all queued messages are forwarded,
1249    ///    install the proxy handler to handle future messages.
1250    ///
1251    /// This sequence prevents the race condition where messages could be delivered
1252    /// out of order or lost during the transition.
1253    pub fn proxy_remaining_messages(self) -> Result<(), crate::Error>
1254    where
1255        Link: HasPeer<Client>,
1256    {
1257        // Destructure self to get ownership of all fields
1258        let ActiveSession {
1259            session_id,
1260            mut update_rx,
1261            update_tx,
1262            connection,
1263            session_handler_registration,
1264            mcp_handler_registrations,
1265            // These fields are not needed for proxying
1266            modes: _,
1267            config_options: _,
1268            meta: _,
1269            _runner,
1270        } = self;
1271
1272        // Step 1: Drop the session handler registration.
1273        // This unregisters the handler that was routing messages to update_rx.
1274        // After this point, no new messages will be added to the channel.
1275        drop(session_handler_registration);
1276
1277        // Step 2: Drop the sender side of the channel.
1278        // This allows us to detect when the channel is fully drained
1279        // (recv will return None when empty and sender is dropped).
1280        drop(update_tx);
1281
1282        // Step 3: Drain any messages that were already queued and forward to client.
1283        // These messages arrived before we dropped the handler but haven't been
1284        // consumed yet. We must forward them to maintain message ordering.
1285        while let Ok(message) = update_rx.try_recv() {
1286            match message {
1287                SessionMessage::SessionMessage(dispatch) => {
1288                    // Forward the message to the client
1289                    connection.send_proxied_message_to(Client, dispatch)?;
1290                }
1291                SessionMessage::StopReason(_) => {
1292                    // StopReason is internal bookkeeping, not forwarded
1293                }
1294            }
1295        }
1296
1297        // Step 4: Install the proxy handler for future messages.
1298        // Now that all queued messages have been forwarded, the proxy handler
1299        // can take over. Any new messages will go directly through the proxy.
1300        connection
1301            .add_dynamic_handler(ProxySessionMessages::new(session_id))?
1302            .detach();
1303
1304        // Keep MCP server handlers alive for the lifetime of the proxy
1305        for registration in mcp_handler_registrations {
1306            registration.detach();
1307        }
1308
1309        Ok(())
1310    }
1311}
1312
1313struct ActiveSessionHandler {
1314    session_id: SessionId,
1315    update_tx: mpsc::UnboundedSender<SessionMessage>,
1316}
1317
1318impl ActiveSessionHandler {
1319    pub fn new(session_id: SessionId, update_tx: mpsc::UnboundedSender<SessionMessage>) -> Self {
1320        Self {
1321            session_id,
1322            update_tx,
1323        }
1324    }
1325}
1326
1327impl<Counterpart: Role> HandleDispatchFrom<Counterpart> for ActiveSessionHandler
1328where
1329    Counterpart: HasPeer<Agent>,
1330{
1331    async fn handle_dispatch_from(
1332        &mut self,
1333        message: Dispatch,
1334        cx: ConnectionTo<Counterpart>,
1335    ) -> Result<Handled<Dispatch>, crate::Error> {
1336        // If this is a message for our session, grab it.
1337        tracing::trace!(
1338            ?message,
1339            handler_session_id = ?self.session_id,
1340            "ActiveSessionHandler::handle_dispatch"
1341        );
1342        MatchDispatchFrom::new(message, &cx)
1343            .if_dispatch_from(Agent, async |message| {
1344                if let Some(session_id) = message.get_session_id()? {
1345                    tracing::trace!(
1346                        message_session_id = ?session_id,
1347                        handler_session_id = ?self.session_id,
1348                        "ActiveSessionHandler::handle_dispatch"
1349                    );
1350                    if session_id == self.session_id {
1351                        self.update_tx
1352                            .unbounded_send(SessionMessage::SessionMessage(message))
1353                            .map_err(crate::util::internal_error)?;
1354                        return Ok(Handled::Yes);
1355                    }
1356                }
1357
1358                // Otherwise, pass it through.
1359                Ok(Handled::No {
1360                    message,
1361                    retry: false,
1362                })
1363            })
1364            .await
1365            .done()
1366    }
1367
1368    fn describe_chain(&self) -> impl std::fmt::Debug {
1369        format!("ActiveSessionHandler({})", self.session_id)
1370    }
1371}
1372
1373#[cfg(not(feature = "unstable_protocol_v2"))]
1374#[allow(
1375    clippy::unnecessary_wraps,
1376    reason = "signature matches the feature-enabled protocol guard"
1377)]
1378fn ensure_v1_session_protocol<Counterpart: Role>(
1379    _connection: &ConnectionTo<Counterpart>,
1380) -> Result<(), crate::Error> {
1381    Ok(())
1382}
1383
1384#[cfg(feature = "unstable_protocol_v2")]
1385fn ensure_v1_session_protocol<Counterpart: Role>(
1386    connection: &ConnectionTo<Counterpart>,
1387) -> Result<(), crate::Error> {
1388    if connection.acp_protocol_version() != Some(crate::schema::ProtocolVersion::V2) {
1389        return Ok(());
1390    }
1391
1392    Err(crate::Error::invalid_request().data(
1393        "stable session builders use ACP protocol v1 types, but this is a protocol v2 connection; \
1394         use the `V2ConnectionTo` supplied to `Client.v2()` callbacks",
1395    ))
1396}