Skip to main content

agent_client_protocol/session/
v2.rs

1use std::{future::Future, path::Path};
2
3use futures::{
4    channel::oneshot,
5    future::{self, Either},
6};
7
8use crate::{
9    Agent, Client, ConnectionTo, DynamicHandlerGuard, JsonRpcRequest, Responder, SentRequest,
10    V2ConnectionTo,
11    jsonrpc::run::{NullRun, RunWithConnectionTo},
12    role::{HasPeer, acp::ProxySessionMessages},
13    schema::v2,
14};
15
16#[cfg(feature = "unstable_mcp_over_acp")]
17use crate::{jsonrpc::run::ChainRun, mcp_server::McpServer};
18
19use super::{
20    SessionCleanup, close_session_registrations, drive_session_cleanup,
21    run_attached_session_runner, session_cleanup, session_runner_connection, wait_session_cleanup,
22};
23
24async fn run_pending_session_setup<Counterpart, Run>(
25    connection: ConnectionTo<Counterpart>,
26    run: Run,
27    started_tx: oneshot::Sender<Result<(), crate::Error>>,
28    promotion_rx: oneshot::Receiver<()>,
29    cleanup: SessionCleanup,
30) -> Result<(), crate::Error>
31where
32    Counterpart: HasPeer<Agent>,
33    Run: RunWithConnectionTo<Counterpart>,
34{
35    let (connection, error_scope) = session_runner_connection(connection, &cleanup);
36    let mut run = Box::pin(run.run_with_connection_to(connection));
37    let first_poll =
38        future::poll_fn(|cx| std::task::Poll::Ready(std::future::Future::poll(run.as_mut(), cx)))
39            .await;
40    // A composed runner may have observed an immediate error but remain Pending
41    // while its sibling drives owned cleanup. It still failed before publication.
42    if let Some(error) = error_scope.error() {
43        close_session_registrations(&cleanup);
44        if first_poll.is_pending() {
45            drop(drive_session_cleanup(run, cleanup).await);
46        } else {
47            wait_session_cleanup(cleanup).await;
48        }
49        drop(started_tx.send(Err(error)));
50        return Ok(());
51    }
52    let readiness = match &first_poll {
53        std::task::Poll::Ready(result) => result.clone(),
54        std::task::Poll::Pending => Ok(()),
55    };
56    drop(started_tx.send(readiness));
57
58    match first_poll {
59        std::task::Poll::Ready(Ok(())) => {
60            if promotion_rx.await.is_err() {
61                close_session_registrations(&cleanup);
62                wait_session_cleanup(cleanup).await;
63            }
64            Ok(())
65        }
66        // The request has not been published yet, so report an immediate
67        // startup failure through its readiness result without failing the
68        // whole connection.
69        std::task::Poll::Ready(Err(_)) => {
70            close_session_registrations(&cleanup);
71            wait_session_cleanup(cleanup).await;
72            Ok(())
73        }
74        std::task::Poll::Pending => match future::select(run, promotion_rx).await {
75            Either::Left((result, promotion_rx)) => {
76                // A Pending first poll releases session setup for publication.
77                // From that point onward the agent may already be using an
78                // attachment, so runner failures are connection-fatal just as
79                // they are for other connection runners.
80                if result.is_err() || promotion_rx.await.is_err() {
81                    close_session_registrations(&cleanup);
82                    wait_session_cleanup(cleanup).await;
83                }
84                result
85            }
86            Either::Right((Ok(()), run)) => {
87                run_attached_session_runner(run, cleanup, error_scope).await
88            }
89            Either::Right((Err(_), run)) => {
90                close_session_registrations(&cleanup);
91                run_attached_session_runner(run, cleanup, error_scope).await
92            }
93        },
94    }
95}
96
97fn send_session_setup<Counterpart, Request, Run>(
98    connection: V2ConnectionTo<Counterpart>,
99    request: Request,
100    dynamic_handler_registrations: Vec<DynamicHandlerGuard<Counterpart>>,
101    run: Run,
102    ordered: bool,
103) -> SentRequest<Request::Response>
104where
105    Counterpart: HasPeer<Agent>,
106    Request: JsonRpcRequest + 'static,
107    Request::Response: 'static,
108    Run: RunWithConnectionTo<Counterpart> + 'static,
109{
110    let raw_connection = connection.raw_connection().clone();
111    if dynamic_handler_registrations.is_empty() {
112        drop(run);
113        if ordered {
114            raw_connection.send_ordered_request_to(Agent, request)
115        } else {
116            raw_connection.send_request_to(Agent, request)
117        }
118    } else {
119        let handlers_ready = raw_connection.dynamic_handler_barrier();
120        let (runner_started_tx, runner_started_rx) = oneshot::channel();
121        let (promotion_tx, promotion_rx) = oneshot::channel();
122        let cleanup = session_cleanup(&dynamic_handler_registrations);
123        let runner_started = match raw_connection.spawn(run_pending_session_setup(
124            raw_connection.clone(),
125            run,
126            runner_started_tx,
127            promotion_rx,
128            cleanup,
129        )) {
130            Ok(()) => Either::Left(async move {
131                runner_started_rx.await.map_err(|error| {
132                    crate::util::internal_error(format!(
133                        "session setup runner stopped before its initial poll: {error}"
134                    ))
135                })?
136            }),
137            Err(error) => Either::Right(future::ready(Err(error))),
138        };
139        let readiness = async move {
140            future::try_join(handlers_ready, runner_started).await?;
141            Ok(())
142        };
143        let response_hook = move |_response: &Request::Response| {
144            promotion_tx.send(()).map_err(|()| {
145                crate::util::internal_error("session setup runner stopped before setup completed")
146            })?;
147            dynamic_handler_registrations
148                .into_iter()
149                .for_each(DynamicHandlerGuard::detach);
150            Ok(())
151        };
152
153        if ordered {
154            raw_connection.send_ordered_request_to_with_response_hook_after(
155                Agent,
156                request,
157                readiness,
158                response_hook,
159            )
160        } else {
161            raw_connection.send_request_to_with_response_hook_after(
162                Agent,
163                request,
164                readiness,
165                response_hook,
166            )
167        }
168    }
169}
170
171impl<Counterpart> V2ConnectionTo<Counterpart>
172where
173    Counterpart: HasPeer<Agent>,
174{
175    /// Build a draft protocol v2 `session/new` request.
176    pub fn build_session(&self, cwd: impl AsRef<Path>) -> V2SessionBuilder<Counterpart> {
177        V2SessionBuilder::new(self, v2::NewSessionRequest::new(cwd.as_ref()))
178    }
179
180    /// Build a draft protocol v2 session using the current working directory.
181    ///
182    /// Returns an error if the current directory cannot be determined.
183    pub fn build_session_cwd(&self) -> Result<V2SessionBuilder<Counterpart>, crate::Error> {
184        let cwd = std::env::current_dir().map_err(|error| {
185            crate::Error::internal_error().data(format!("cannot get current directory: {error}"))
186        })?;
187        Ok(self.build_session(cwd))
188    }
189
190    /// Build a draft protocol v2 session from an existing `session/new` request.
191    pub fn build_session_from(
192        &self,
193        request: v2::NewSessionRequest,
194    ) -> V2SessionBuilder<Counterpart> {
195        V2SessionBuilder::new(self, request)
196    }
197
198    /// Build an unstable draft protocol v2 `session/fork` request.
199    ///
200    /// This helper is available with the `unstable_session_fork` feature. Call
201    /// [`V2ForkSessionBuilder::start_session`] to publish the request and
202    /// obtain a command handle for the newly created fork.
203    #[cfg(feature = "unstable_session_fork")]
204    pub fn fork_session(
205        &self,
206        session_id: impl Into<v2::SessionId>,
207        cwd: impl AsRef<Path>,
208    ) -> V2ForkSessionBuilder<Counterpart> {
209        self.fork_session_from(v2::ForkSessionRequest::new(session_id, cwd.as_ref()))
210    }
211
212    /// Build an unstable draft protocol v2 `session/fork` request from an
213    /// existing request.
214    ///
215    /// This helper is available with the `unstable_session_fork` feature. Call
216    /// [`V2ForkSessionBuilder::start_session`] to publish the request.
217    #[cfg(feature = "unstable_session_fork")]
218    pub fn fork_session_from(
219        &self,
220        request: v2::ForkSessionRequest,
221    ) -> V2ForkSessionBuilder<Counterpart> {
222        V2ForkSessionBuilder::new(self, request)
223    }
224
225    /// Build a draft protocol v2 `session/resume` request.
226    ///
227    /// Use [`Self::resume_session_from`] to request history replay or set
228    /// other optional resume parameters. Call
229    /// [`V2ResumeSessionBuilder::start_session`] to publish the request.
230    pub fn resume_session(
231        &self,
232        session_id: impl Into<v2::SessionId>,
233        cwd: impl AsRef<Path>,
234    ) -> V2ResumeSessionBuilder<Counterpart> {
235        self.resume_session_from(v2::ResumeSessionRequest::new(session_id, cwd.as_ref()))
236    }
237
238    /// Build a draft protocol v2 `session/resume` request from an existing request.
239    ///
240    /// Register typed session update handlers before connecting. When the
241    /// request asks for replay, the agent sends those updates before the
242    /// [`v2::ResumeSessionResponse`]. Call
243    /// [`V2ResumeSessionBuilder::start_session`] to publish the request.
244    pub fn resume_session_from(
245        &self,
246        request: v2::ResumeSessionRequest,
247    ) -> V2ResumeSessionBuilder<Counterpart> {
248        V2ResumeSessionBuilder::new(self, request)
249    }
250}
251
252/// Builder for a draft protocol v2 `session/new` request.
253///
254/// Protocol v2 acknowledges `session/prompt` independently from inbound
255/// session updates. Register typed [`v2::UpdateSessionNotification`] and
256/// session request handlers on [`crate::Builder`] before connecting, then use
257/// [`Self::start_session`] to create the command-only [`V2Session`] handle or
258/// `on_proxy_session_start` to forward setup through a proxy.
259///
260/// With both the `unstable_protocol_v2` and `unstable_mcp_over_acp` features,
261/// `with_mcp_server` attaches an MCP server to the new session.
262#[must_use = "call `start_session` or `on_proxy_session_start` to send the `session/new` request"]
263#[derive(Debug)]
264pub struct V2SessionBuilder<Counterpart, Run = NullRun>
265where
266    Counterpart: HasPeer<Agent>,
267    Run: RunWithConnectionTo<Counterpart>,
268{
269    connection: V2ConnectionTo<Counterpart>,
270    request: v2::NewSessionRequest,
271    dynamic_handler_registrations: Vec<DynamicHandlerGuard<Counterpart>>,
272    run: Run,
273}
274
275impl<Counterpart> V2SessionBuilder<Counterpart, NullRun>
276where
277    Counterpart: HasPeer<Agent>,
278{
279    fn new(connection: &V2ConnectionTo<Counterpart>, request: v2::NewSessionRequest) -> Self {
280        Self {
281            connection: connection.clone(),
282            request,
283            dynamic_handler_registrations: Vec::new(),
284            run: NullRun,
285        }
286    }
287}
288
289impl<Counterpart, Run> V2SessionBuilder<Counterpart, Run>
290where
291    Counterpart: HasPeer<Agent>,
292    Run: RunWithConnectionTo<Counterpart>,
293{
294    /// Attach an MCP server to this new protocol v2 session.
295    ///
296    /// This method is available when both `unstable_protocol_v2` and
297    /// `unstable_mcp_over_acp` are enabled. MCP routes are installed and their
298    /// runner tasks receive an initial poll before `session/new` is published,
299    /// allowing the agent to connect while handling session setup. A
300    /// successful attachment remains active for the lifetime of the connection.
301    #[cfg(feature = "unstable_mcp_over_acp")]
302    pub fn with_mcp_server<McpRun>(
303        mut self,
304        mcp_server: McpServer<Counterpart, McpRun>,
305    ) -> Result<V2SessionBuilder<Counterpart, ChainRun<Run, McpRun>>, crate::Error>
306    where
307        McpRun: RunWithConnectionTo<Counterpart>,
308    {
309        let (handler, mcp_run) = mcp_server.into_v2_handler_and_runner();
310        self.dynamic_handler_registrations
311            .push(handler.into_dynamic_handler(&mut self.request.mcp_servers, &self.connection)?);
312        Ok(V2SessionBuilder {
313            connection: self.connection,
314            request: self.request,
315            dynamic_handler_registrations: self.dynamic_handler_registrations,
316            run: ChainRun::new(self.run, mcp_run),
317        })
318    }
319
320    fn send_new_session(self, ordered: bool) -> SentRequest<v2::NewSessionResponse>
321    where
322        Run: 'static,
323    {
324        let Self {
325            connection,
326            request,
327            dynamic_handler_registrations,
328            run,
329        } = self;
330        send_session_setup(
331            connection,
332            request,
333            dynamic_handler_registrations,
334            run,
335            ordered,
336        )
337    }
338
339    /// Send `session/new` and return its independently consumable request.
340    ///
341    /// The successful result contains both a cloneable command handle and the
342    /// complete [`v2::NewSessionResponse`]. Consume the returned request with
343    /// [`SentRequest::block_task`], [`SentRequest::on_receiving_result`], or
344    /// another explicit [`SentRequest`] completion mode.
345    ///
346    /// Attached MCP routes are installed and their runner tasks begin
347    /// executing before the request is published. A valid success response
348    /// promotes them to the connection lifetime, independently from how this
349    /// request handle is consumed. Setup errors clean up the pending
350    /// attachment.
351    pub fn start_session(self) -> SentRequest<OpenedV2Session<Counterpart, v2::NewSessionResponse>>
352    where
353        Run: 'static,
354    {
355        let session_connection = self.connection.clone();
356        self.send_new_session(false).map(move |response| {
357            let session = V2Session {
358                session_id: response.session_id.clone(),
359                connection: session_connection,
360            };
361            Ok(OpenedV2Session { session, response })
362        })
363    }
364
365    /// Start a protocol v2 session through a proxy and forward its response.
366    ///
367    /// The downstream request is ordered and inherits cancellation from the
368    /// upstream request. On success, this helper installs session routing before
369    /// later inbound traffic is processed, forwards the complete response, and
370    /// spawns `op` with an [`OpenedV2Session`] containing the command-only
371    /// session handle plus the complete setup response. Inbound updates and
372    /// interactive requests remain independent connection traffic.
373    ///
374    /// The callback runs outside the ordered response barrier, so it may wait
375    /// for later connection traffic without deadlocking the dispatch loop.
376    pub fn on_proxy_session_start<F, Fut>(
377        self,
378        responder: Responder<v2::NewSessionResponse>,
379        op: F,
380    ) -> Result<(), crate::Error>
381    where
382        Counterpart: HasPeer<Client>,
383        Run: 'static,
384        F: FnOnce(OpenedV2Session<Counterpart, v2::NewSessionResponse>) -> Fut + Send + 'static,
385        Fut: Future<Output = Result<(), crate::Error>> + Send,
386    {
387        let session_connection = self.connection.clone();
388        self.send_new_session(true)
389            .forward_cancellation_from(responder.cancellation())
390            .on_receiving_ok_result(responder, async move |response, responder| {
391                let session_id = response.session_id.clone();
392                let raw_connection = session_connection.raw_connection();
393                let route = match raw_connection.add_dynamic_handler(ProxySessionMessages::new(
394                    crate::schema::v1::SessionId::from(session_id.0.clone()),
395                )) {
396                    Ok(route) => route,
397                    Err(error) => return responder.respond_with_error(error),
398                };
399
400                let opened = OpenedV2Session {
401                    session: V2Session {
402                        session_id,
403                        connection: session_connection.clone(),
404                    },
405                    response: response.clone(),
406                };
407                responder.respond(response)?;
408                route.detach();
409                raw_connection.spawn(async move { op(opened).await })
410            })
411    }
412}
413
414/// Builder for an unstable draft protocol v2 `session/fork` request.
415///
416/// A successful fork creates a new independent session whose ID comes from the
417/// [`v2::ForkSessionResponse`]. Register typed
418/// [`v2::UpdateSessionNotification`] and interactive request handlers on
419/// [`crate::Builder`] before connecting, then use [`Self::start_session`] to
420/// obtain a command handle for the fork or `on_proxy_session_start` to forward
421/// setup through a proxy.
422///
423/// This type is available with the `unstable_session_fork` feature. With
424/// `unstable_mcp_over_acp` as well, `with_mcp_server` attaches an MCP server to
425/// the forked session.
426#[cfg(feature = "unstable_session_fork")]
427#[must_use = "call `start_session` or `on_proxy_session_start` to send the `session/fork` request"]
428#[derive(Debug)]
429pub struct V2ForkSessionBuilder<Counterpart, Run = NullRun>
430where
431    Counterpart: HasPeer<Agent>,
432    Run: RunWithConnectionTo<Counterpart>,
433{
434    connection: V2ConnectionTo<Counterpart>,
435    request: v2::ForkSessionRequest,
436    dynamic_handler_registrations: Vec<DynamicHandlerGuard<Counterpart>>,
437    run: Run,
438}
439
440#[cfg(feature = "unstable_session_fork")]
441impl<Counterpart> V2ForkSessionBuilder<Counterpart, NullRun>
442where
443    Counterpart: HasPeer<Agent>,
444{
445    fn new(connection: &V2ConnectionTo<Counterpart>, request: v2::ForkSessionRequest) -> Self {
446        Self {
447            connection: connection.clone(),
448            request,
449            dynamic_handler_registrations: Vec::new(),
450            run: NullRun,
451        }
452    }
453}
454
455#[cfg(feature = "unstable_session_fork")]
456impl<Counterpart, Run> V2ForkSessionBuilder<Counterpart, Run>
457where
458    Counterpart: HasPeer<Agent>,
459    Run: RunWithConnectionTo<Counterpart>,
460{
461    /// Attach an MCP server to this forked protocol v2 session.
462    ///
463    /// This method is available when `unstable_mcp_over_acp` is enabled in
464    /// addition to `unstable_protocol_v2` and `unstable_session_fork`. MCP
465    /// routes are installed and their runner tasks receive an initial poll
466    /// before `session/fork` is published, allowing the agent to connect while
467    /// handling session setup. A successful attachment remains active for the
468    /// lifetime of the connection.
469    #[cfg(feature = "unstable_mcp_over_acp")]
470    pub fn with_mcp_server<McpRun>(
471        mut self,
472        mcp_server: McpServer<Counterpart, McpRun>,
473    ) -> Result<V2ForkSessionBuilder<Counterpart, ChainRun<Run, McpRun>>, crate::Error>
474    where
475        McpRun: RunWithConnectionTo<Counterpart>,
476    {
477        let (handler, mcp_run) = mcp_server.into_v2_handler_and_runner();
478        self.dynamic_handler_registrations
479            .push(handler.into_dynamic_handler(&mut self.request.mcp_servers, &self.connection)?);
480        Ok(V2ForkSessionBuilder {
481            connection: self.connection,
482            request: self.request,
483            dynamic_handler_registrations: self.dynamic_handler_registrations,
484            run: ChainRun::new(self.run, mcp_run),
485        })
486    }
487
488    fn send_fork_session(self, ordered: bool) -> SentRequest<v2::ForkSessionResponse>
489    where
490        Run: 'static,
491    {
492        let Self {
493            connection,
494            request,
495            dynamic_handler_registrations,
496            run,
497        } = self;
498        send_session_setup(
499            connection,
500            request,
501            dynamic_handler_registrations,
502            run,
503            ordered,
504        )
505    }
506
507    /// Send `session/fork` and return its independently consumable request.
508    ///
509    /// The successful result contains both a cloneable command handle for the
510    /// newly created fork and the complete [`v2::ForkSessionResponse`]. Consume
511    /// the returned request with [`SentRequest::block_task`],
512    /// [`SentRequest::on_receiving_result`], or another explicit [`SentRequest`]
513    /// completion mode.
514    ///
515    /// Attached MCP routes are installed and their runner tasks begin
516    /// executing before the request is published. A valid success response
517    /// promotes them to the connection lifetime independently from how this
518    /// request handle is consumed; setup errors clean up the pending
519    /// attachment.
520    pub fn start_session(self) -> SentRequest<OpenedV2Session<Counterpart, v2::ForkSessionResponse>>
521    where
522        Run: 'static,
523    {
524        let session_connection = self.connection.clone();
525        self.send_fork_session(false).map(move |response| {
526            let session = V2Session {
527                session_id: response.session_id.clone(),
528                connection: session_connection,
529            };
530            Ok(OpenedV2Session { session, response })
531        })
532    }
533
534    /// Fork a protocol v2 session through a proxy and forward its response.
535    ///
536    /// The downstream request is ordered and inherits cancellation from the
537    /// upstream request. On success, this helper obtains the new session ID
538    /// from the response, installs session routing before later inbound traffic
539    /// is processed, forwards the complete response, and spawns `op` with an
540    /// [`OpenedV2Session`] containing the fork's command handle and response.
541    /// Inbound updates and interactive requests remain independent connection
542    /// traffic.
543    ///
544    /// The callback runs outside the ordered response barrier, so it may wait
545    /// for later connection traffic without deadlocking the dispatch loop.
546    pub fn on_proxy_session_start<F, Fut>(
547        self,
548        responder: Responder<v2::ForkSessionResponse>,
549        op: F,
550    ) -> Result<(), crate::Error>
551    where
552        Counterpart: HasPeer<Client>,
553        Run: 'static,
554        F: FnOnce(OpenedV2Session<Counterpart, v2::ForkSessionResponse>) -> Fut + Send + 'static,
555        Fut: Future<Output = Result<(), crate::Error>> + Send,
556    {
557        let session_connection = self.connection.clone();
558        self.send_fork_session(true)
559            .forward_cancellation_from(responder.cancellation())
560            .on_receiving_ok_result(responder, async move |response, responder| {
561                let session_id = response.session_id.clone();
562                let raw_connection = session_connection.raw_connection();
563                let route = match raw_connection.add_dynamic_handler(ProxySessionMessages::new(
564                    crate::schema::v1::SessionId::from(session_id.0.clone()),
565                )) {
566                    Ok(route) => route,
567                    Err(error) => return responder.respond_with_error(error),
568                };
569
570                let opened = OpenedV2Session {
571                    session: V2Session {
572                        session_id,
573                        connection: session_connection.clone(),
574                    },
575                    response: response.clone(),
576                };
577                responder.respond(response)?;
578                route.detach();
579                raw_connection.spawn(async move { op(opened).await })
580            })
581    }
582}
583
584/// Builder for a draft protocol v2 `session/resume` request.
585///
586/// Replay updates arrive before the resume response. Direct clients must
587/// register typed [`v2::UpdateSessionNotification`] and interactive request
588/// handlers on [`crate::Builder`] before connecting. Proxies should use
589/// [`Self::on_proxy_session_start`], which makes downstream session routing
590/// ready before publishing the resume request.
591///
592/// With both the `unstable_protocol_v2` and `unstable_mcp_over_acp` features,
593/// `with_mcp_server` attaches an MCP server to the resumed session and makes it
594/// ready before the agent can send replay or its response.
595#[must_use = "call `start_session` or `on_proxy_session_start` to send the `session/resume` request"]
596#[derive(Debug)]
597pub struct V2ResumeSessionBuilder<Counterpart, Run = NullRun>
598where
599    Counterpart: HasPeer<Agent>,
600    Run: RunWithConnectionTo<Counterpart>,
601{
602    connection: V2ConnectionTo<Counterpart>,
603    request: v2::ResumeSessionRequest,
604    dynamic_handler_registrations: Vec<DynamicHandlerGuard<Counterpart>>,
605    run: Run,
606}
607
608impl<Counterpart> V2ResumeSessionBuilder<Counterpart, NullRun>
609where
610    Counterpart: HasPeer<Agent>,
611{
612    fn new(connection: &V2ConnectionTo<Counterpart>, request: v2::ResumeSessionRequest) -> Self {
613        Self {
614            connection: connection.clone(),
615            request,
616            dynamic_handler_registrations: Vec::new(),
617            run: NullRun,
618        }
619    }
620}
621
622impl<Counterpart, Run> V2ResumeSessionBuilder<Counterpart, Run>
623where
624    Counterpart: HasPeer<Agent>,
625    Run: RunWithConnectionTo<Counterpart>,
626{
627    /// Attach an MCP server to this resumed protocol v2 session.
628    ///
629    /// This method is available when both `unstable_protocol_v2` and
630    /// `unstable_mcp_over_acp` are enabled. MCP routes are installed and their
631    /// runner tasks receive an initial poll before `session/resume` is
632    /// published, allowing the agent to use the server during replay and
633    /// session setup. A successful attachment remains active for the lifetime
634    /// of the connection.
635    #[cfg(feature = "unstable_mcp_over_acp")]
636    pub fn with_mcp_server<McpRun>(
637        mut self,
638        mcp_server: McpServer<Counterpart, McpRun>,
639    ) -> Result<V2ResumeSessionBuilder<Counterpart, ChainRun<Run, McpRun>>, crate::Error>
640    where
641        McpRun: RunWithConnectionTo<Counterpart>,
642    {
643        let (handler, mcp_run) = mcp_server.into_v2_handler_and_runner();
644        self.dynamic_handler_registrations
645            .push(handler.into_dynamic_handler(&mut self.request.mcp_servers, &self.connection)?);
646        Ok(V2ResumeSessionBuilder {
647            connection: self.connection,
648            request: self.request,
649            dynamic_handler_registrations: self.dynamic_handler_registrations,
650            run: ChainRun::new(self.run, mcp_run),
651        })
652    }
653
654    fn send_resume_session(self, ordered: bool) -> SentRequest<v2::ResumeSessionResponse>
655    where
656        Run: 'static,
657    {
658        let Self {
659            connection,
660            request,
661            dynamic_handler_registrations,
662            run,
663        } = self;
664        send_session_setup(
665            connection,
666            request,
667            dynamic_handler_registrations,
668            run,
669            ordered,
670        )
671    }
672
673    /// Send `session/resume` and return its independently consumable request.
674    ///
675    /// The successful result contains both a cloneable command handle and the
676    /// complete [`v2::ResumeSessionResponse`]. Consume the returned request
677    /// with [`SentRequest::block_task`], [`SentRequest::on_receiving_result`],
678    /// or another explicit [`SentRequest`] completion mode.
679    ///
680    /// Replay is delivered through the typed connection handlers before the
681    /// response. Attached MCP routes and runner tasks are ready before the
682    /// request is published. A valid success response promotes them to the
683    /// connection lifetime independently from how this request handle is
684    /// consumed; setup errors clean up the pending attachment.
685    pub fn start_session(
686        self,
687    ) -> SentRequest<OpenedV2Session<Counterpart, v2::ResumeSessionResponse>>
688    where
689        Run: 'static,
690    {
691        let session_id = self.request.session_id.clone();
692        let session_connection = self.connection.clone();
693        self.send_resume_session(false).map(move |response| {
694            let session = V2Session {
695                session_id,
696                connection: session_connection,
697            };
698            Ok(OpenedV2Session { session, response })
699        })
700    }
701
702    /// Resume a protocol v2 session through a proxy and forward its response.
703    ///
704    /// The session ID is known before the request, so this helper installs and
705    /// acknowledges downstream session routing before publishing the ordered
706    /// `session/resume` request. Replay can therefore be forwarded before the
707    /// response as required by the protocol. The downstream request inherits
708    /// cancellation from the upstream request.
709    ///
710    /// On success, the helper forwards the complete response and spawns `op`
711    /// with an [`OpenedV2Session`] containing the command-only session handle
712    /// plus that response. The callback runs outside the ordered response
713    /// barrier, so it may wait for later connection traffic without
714    /// deadlocking the dispatch loop.
715    pub fn on_proxy_session_start<F, Fut>(
716        mut self,
717        responder: Responder<v2::ResumeSessionResponse>,
718        op: F,
719    ) -> Result<(), crate::Error>
720    where
721        Counterpart: HasPeer<Client>,
722        Run: 'static,
723        F: FnOnce(OpenedV2Session<Counterpart, v2::ResumeSessionResponse>) -> Fut + Send + 'static,
724        Fut: Future<Output = Result<(), crate::Error>> + Send,
725    {
726        let session_id = self.request.session_id.clone();
727        let session_connection = self.connection.clone();
728        self.dynamic_handler_registrations.push(
729            session_connection
730                .raw_connection()
731                .add_dynamic_handler(ProxySessionMessages::new(
732                    crate::schema::v1::SessionId::from(session_id.0.clone()),
733                ))?,
734        );
735
736        self.send_resume_session(true)
737            .forward_cancellation_from(responder.cancellation())
738            .on_receiving_ok_result(responder, async move |response, responder| {
739                let opened = OpenedV2Session {
740                    session: V2Session {
741                        session_id,
742                        connection: session_connection.clone(),
743                    },
744                    response: response.clone(),
745                };
746                responder.respond(response)?;
747                session_connection
748                    .raw_connection()
749                    .spawn(async move { op(opened).await })
750            })
751    }
752}
753
754/// A newly available protocol v2 session and its operation-specific response.
755///
756/// Keeping the response separate from [`V2Session`] avoids treating
757/// `session/new` setup data as mutable session state and lets each setup
758/// operation, including `session/resume`, retain its own complete response
759/// type.
760#[derive(Debug)]
761pub struct OpenedV2Session<Link, Response>
762where
763    Link: HasPeer<Agent>,
764{
765    session: V2Session<Link>,
766    response: Response,
767}
768
769impl<Link, Response> OpenedV2Session<Link, Response>
770where
771    Link: HasPeer<Agent>,
772{
773    /// Access the command handle for the opened session.
774    pub fn session(&self) -> &V2Session<Link> {
775        &self.session
776    }
777
778    /// Access the complete response from the operation that opened the session.
779    pub fn response(&self) -> &Response {
780        &self.response
781    }
782
783    /// Split this result into the command handle and complete setup response.
784    pub fn into_parts(self) -> (V2Session<Link>, Response) {
785        (self.session, self.response)
786    }
787
788    /// Consume this result and return only the command handle.
789    pub fn into_session(self) -> V2Session<Link> {
790        self.session
791    }
792}
793
794/// Cloneable command handle for a draft protocol v2 session.
795///
796/// Inbound protocol traffic is intentionally not owned by this value. Receive
797/// authoritative [`v2::UpdateSessionNotification`] values and interactive
798/// requests such as [`v2::RequestPermissionRequest`] through typed handlers
799/// installed on [`crate::Builder`].
800#[derive(Debug, Clone)]
801pub struct V2Session<Link>
802where
803    Link: HasPeer<Agent>,
804{
805    session_id: v2::SessionId,
806    connection: V2ConnectionTo<Link>,
807}
808
809impl<Link> V2Session<Link>
810where
811    Link: HasPeer<Agent>,
812{
813    /// Access the session ID.
814    pub fn session_id(&self) -> &v2::SessionId {
815        &self.session_id
816    }
817
818    /// Access the underlying connection.
819    pub fn connection(&self) -> &V2ConnectionTo<Link> {
820        &self.connection
821    }
822
823    /// Submit a text prompt and return its independent acceptance request.
824    ///
825    /// A successful response only acknowledges that the agent accepted the
826    /// prompt. The accepted user message, output, state changes, and completion
827    /// arrive independently through [`v2::UpdateSessionNotification`].
828    pub fn send_prompt(&self, prompt: impl ToString) -> SentRequest<v2::PromptResponse> {
829        self.send_prompt_blocks(vec![prompt.to_string().into()])
830    }
831
832    /// Submit arbitrary prompt content and return its acceptance request.
833    ///
834    /// The SDK does not track foreground state or gate prompt submission
835    /// locally. Wait for an `idle` state update before another ordinary prompt
836    /// unless using a separately defined admission mechanism.
837    pub fn send_prompt_blocks(
838        &self,
839        prompt: Vec<v2::ContentBlock>,
840    ) -> SentRequest<v2::PromptResponse> {
841        self.connection.send_request_to(
842            Agent,
843            v2::PromptRequest::new(self.session_id.clone(), prompt),
844        )
845    }
846
847    /// Ask the agent to cancel the session's current foreground work.
848    ///
849    /// This is independent from cancelling a prompt's [`SentRequest`].
850    /// Cancellation completes when the agent reports an `idle` state update
851    /// with [`v2::StopReason::Cancelled`]. The client should immediately mark
852    /// unfinished tool calls for the active work as cancelled and remains
853    /// responsible for resolving every pending [`v2::RequestPermissionRequest`]
854    /// with the cancelled outcome.
855    pub fn cancel_active_work(&self) -> Result<(), crate::Error> {
856        self.connection.send_notification_to(
857            Agent,
858            v2::CancelSessionNotification::new(self.session_id.clone()),
859        )
860    }
861
862    /// Set a session configuration option.
863    ///
864    /// The response contains the full current option set. It is not cached on
865    /// this command handle.
866    pub fn set_config_option(
867        &self,
868        config_id: impl Into<v2::SessionConfigId>,
869        value: impl Into<v2::SessionConfigOptionValue>,
870    ) -> SentRequest<v2::SetSessionConfigOptionResponse> {
871        self.connection.send_request_to(
872            Agent,
873            v2::SetSessionConfigOptionRequest::new(self.session_id.clone(), config_id, value),
874        )
875    }
876
877    /// Close the remote session and release its resources.
878    ///
879    /// Existing clones of this local command handle are not invalidated, but
880    /// the agent should reject subsequent commands for the closed session.
881    pub fn close(&self) -> SentRequest<v2::CloseSessionResponse> {
882        self.connection
883            .send_request_to(Agent, v2::CloseSessionRequest::new(self.session_id.clone()))
884    }
885}