Skip to main content

agent_client_protocol/
component.rs

1//! ConnectTo abstraction for agents and proxies.
2//!
3//! This module provides the [`ConnectTo`] trait that defines the interface for things
4//! that can be run as part of a conductor's chain - agents, proxies, or any ACP-speaking component.
5//!
6//! ## Usage
7//!
8//! Components connect to other components, creating a chain of message processors.
9//! The type parameter `R` is the role that this component connects to (its counterpart).
10//!
11//! To implement a component, implement the `connect_to` method:
12//!
13//! ```rust
14//! use agent_client_protocol::{Agent, Client, ConnectTo, Result};
15//!
16//! struct MyAgent;
17//!
18//! // An agent connects to clients
19//! impl ConnectTo<Client> for MyAgent {
20//!     async fn connect_to(self, client: impl ConnectTo<Agent>) -> Result<()> {
21//!         Agent.builder()
22//!             .name("my-agent")
23//!             .connect_to(client)
24//!             .await
25//!     }
26//! }
27//! ```
28
29use futures::future::BoxFuture;
30use std::{
31    fmt::Debug,
32    future::Future,
33    marker::PhantomData,
34    pin::Pin,
35    task::{Context, Poll},
36};
37
38use crate::{Channel, Result, role::Role};
39
40// Presence of the control records the cooperative contract, even after its
41// one-shot action has run. Moving a requested driver must not erase that fact.
42pub(crate) struct FinishControl {
43    hook: Option<Box<dyn FnOnce() + Send + 'static>>,
44}
45
46impl FinishControl {
47    fn new(hook: impl FnOnce() + Send + 'static) -> Self {
48        Self {
49            hook: Some(Box::new(hook)),
50        }
51    }
52
53    pub(crate) fn request(&mut self) {
54        if let Some(hook) = self.hook.take() {
55            hook();
56        }
57    }
58}
59
60/// Drives owned endpoint work.
61///
62/// A driver owns the endpoint: successful completion means no further
63/// output is expected, and adapters must drain output already accepted before
64/// terminating. Errors may abort the connection without guaranteed output
65/// drain; an adapter may still preserve queued error replies before terminating.
66///
67/// Poll the driver concurrently with channel traffic. Endpoints without owned
68/// work return `None` from [`ConnectTo::into_channel_and_future`], not a driver:
69/// their channel halves independently determine their lifetime.
70///
71/// Use [`new`](Self::new) for opaque work or
72/// [`with_finish`](Self::with_finish) for a transport that can finish gracefully.
73/// Use [`map_future`](Self::map_future) to decorate existing work without losing
74/// its finish capability.
75#[must_use = "connection drivers must be polled to make progress"]
76pub struct ConnectionDriver {
77    future: BoxFuture<'static, Result<()>>,
78    finish: Option<FinishControl>,
79}
80
81impl ConnectionDriver {
82    /// Create a driver that owns the endpoint's lifetime.
83    ///
84    /// This driver has no cooperative finish hook. A finite foreground may drop
85    /// it after handing off accepted output, rather than wait for arbitrary
86    /// work to finish. Reactive serving still awaits owned work after input EOF.
87    ///
88    /// Custom transports that need to flush before a finite foreground returns
89    /// should use [`with_finish`](Self::with_finish) instead.
90    pub fn new(future: impl Future<Output = Result<()>> + Send + 'static) -> Self {
91        Self {
92            future: Box::pin(future),
93            finish: None,
94        }
95    }
96
97    /// Create owned work that supports cooperative graceful completion.
98    ///
99    /// The finish hook only requests completion; it must be nonblocking and
100    /// should signal the future to stop accepting output, drain what it has
101    /// already accepted, flush and close its write half, then return. It must
102    /// not require independently open remote input to reach EOF. The future
103    /// remains responsible for reporting I/O and flush errors.
104    ///
105    /// SDK consumers invoke the hook after handing off their accepted output,
106    /// then continue polling the driver until completion. There is no implicit
107    /// timeout: if the adapter cannot finish, the enclosing connection remains
108    /// pending and may be cancelled by its caller.
109    ///
110    /// The hook is invoked at most once. Dropping the driver drops its owned
111    /// future without requesting graceful completion. Dropping only the hook
112    /// does not invoke it or necessarily stop the work.
113    ///
114    /// # Example
115    ///
116    /// A custom adapter can use any signal understood by its future. For
117    /// example, a one-shot channel separates the finish request from completion:
118    ///
119    /// ```
120    /// use agent_client_protocol::ConnectionDriver;
121    /// use futures::{channel::oneshot, FutureExt};
122    ///
123    /// let (finish_tx, finish_rx) = oneshot::channel();
124    /// let mut driver = ConnectionDriver::with_finish(
125    ///     async move {
126    ///         if finish_rx.await.is_err() {
127    ///             // Losing the hook must not masquerade as a finish request.
128    ///             futures::future::pending::<()>().await;
129    ///         }
130    ///         // Seal the adapter's outgoing queue, drain it, and flush/close
131    ///         // the physical writer here before returning.
132    ///         Ok(())
133    ///     },
134    ///     move || { let _ = finish_tx.send(()); },
135    /// );
136    ///
137    /// assert!((&mut driver).now_or_never().is_none());
138    /// assert!(driver.request_finish());
139    /// assert!(driver.request_finish()); // Supported, but the hook runs only once.
140    /// futures::executor::block_on(driver).unwrap();
141    /// ```
142    pub fn with_finish(
143        future: impl Future<Output = Result<()>> + Send + 'static,
144        finish: impl FnOnce() + Send + 'static,
145    ) -> Self {
146        Self {
147            future: Box::pin(future),
148            finish: Some(FinishControl::new(finish)),
149        }
150    }
151
152    /// Decorate the owned future while preserving its finish capability.
153    ///
154    /// This is useful for tracing, error annotation, or completion cleanup.
155    /// Wrapping this driver in [`new`](Self::new) instead would hide its finish
156    /// control from the outer driver.
157    ///
158    /// `map` is called immediately and receives the boxed future, not the
159    /// driver. Its returned future must uphold the same completion contract:
160    /// keep driving the original work and do not report success before accepted
161    /// output is drained. An already-requested finish remains requested, and
162    /// opaque work remains opaque.
163    ///
164    /// ```
165    /// use agent_client_protocol::ConnectionDriver;
166    /// use futures::FutureExt;
167    ///
168    /// let driver = ConnectionDriver::new(async { Ok(()) });
169    /// let decorated = driver.map_future(|work| {
170    ///     work.inspect(|result| eprintln!("transport completed: {result:?}"))
171    /// });
172    /// futures::executor::block_on(decorated).unwrap();
173    /// ```
174    pub fn map_future<F>(self, map: impl FnOnce(BoxFuture<'static, Result<()>>) -> F) -> Self
175    where
176        F: Future<Output = Result<()>> + Send + 'static,
177    {
178        Self {
179            future: Box::pin(map(self.future)),
180            finish: self.finish,
181        }
182    }
183
184    /// Request graceful completion, without waiting for it.
185    ///
186    /// Returns `true` if this driver supports cooperative finish, including
187    /// when finish was already requested. Repeated requests are idempotent:
188    /// the hook runs at most once and the driver retains its graceful-finish
189    /// contract across wrapping or ownership handoff.
190    ///
191    /// Returns `false` for opaque work constructed with [`new`](Self::new);
192    /// this method does not cancel that work. A `true` return does not prove
193    /// flushing is complete: continue polling or await the driver to observe
194    /// completion and any errors.
195    #[must_use]
196    pub fn request_finish(&mut self) -> bool {
197        if let Some(finish) = self.finish.as_mut() {
198            finish.request();
199            true
200        } else {
201            false
202        }
203    }
204
205    pub(crate) fn take_finish(&mut self) -> Option<FinishControl> {
206        self.finish.take()
207    }
208}
209
210impl Future for ConnectionDriver {
211    type Output = Result<()>;
212
213    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
214        self.future.as_mut().poll(cx)
215    }
216}
217
218impl Debug for ConnectionDriver {
219    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
220        f.debug_struct("ConnectionDriver")
221            .field("finishable", &self.finish.is_some())
222            .finish_non_exhaustive()
223    }
224}
225
226/// A component that can exchange JSON-RPC messages to an endpoint playing the role `R`
227/// (e.g., an ACP [`Agent`](`crate::role::acp::Agent`) or an MCP [`Server`](`crate::role::mcp::Server`)).
228///
229/// This trait represents anything that can communicate via JSON-RPC messages over channels -
230/// agents, proxies, in-process connections, or any ACP-speaking component.
231///
232/// The type parameter `R` is the role that this component connects to (its counterpart).
233/// For example:
234/// - An agent implements `ConnectTo<Client>` to connect to clients
235/// - A proxy implements `ConnectTo<Conductor>` to connect to conductors
236/// - Transports like `Channel` implement `ConnectTo<R>` for every `R` because they are role-agnostic
237///
238/// # Component Types
239///
240/// The trait is implemented by several built-in types representing different communication patterns:
241///
242/// - **[`Lines`]**: A component communicating over asynchronous line streams
243/// - **[`ByteStreams`]**: A component communicating over byte streams (stdin/stdout, sockets, etc.)
244/// - **[`Channel`]**: A component communicating via in-process message channels (for testing or direct connections)
245/// - **Custom components**: Proxies, transformers, or any ACP-aware service
246#[cfg_attr(
247    all(feature = "process", not(target_family = "wasm")),
248    doc = "- **[`AcpAgent`]**: An external agent running in a separate process with stdio communication"
249)]
250///
251/// # Two Ways to Connect
252///
253/// Components can be used in two ways:
254///
255/// 1. **`connect_to(client)`** - Connect directly to another component (most components implement this)
256/// 2. **`into_channel_and_future()`** - Obtain a channel endpoint and optional owned driver
257///
258/// Most components only need to implement `connect_to(client)`. The
259/// `into_channel_and_future()` method has a default implementation that creates an intermediate
260/// channel and calls `connect_to`.
261///
262/// # Implementation Example
263///
264/// ```rust
265/// use agent_client_protocol::{Agent, Client, ConnectTo, Result};
266///
267/// struct MyAgent;
268///
269/// impl ConnectTo<Client> for MyAgent {
270///     async fn connect_to(self, client: impl ConnectTo<Agent>) -> Result<()> {
271///         Agent.builder()
272///             .name("my-agent")
273///             .connect_to(client)
274///             .await
275///     }
276/// }
277/// ```
278///
279/// # Heterogeneous Collections
280///
281/// For storing different component types in the same collection, use [`DynConnectTo`]:
282///
283/// ```rust
284/// use agent_client_protocol::{Channel, Client, DynConnectTo};
285///
286/// let (first, _first_peer) = Channel::duplex();
287/// let (second, _second_peer) = Channel::duplex();
288/// let components: Vec<DynConnectTo<Client>> = vec![
289///     DynConnectTo::new(first),
290///     DynConnectTo::new(second),
291/// ];
292/// assert_eq!(components.len(), 2);
293/// ```
294///
295/// [`ByteStreams`]: crate::ByteStreams
296/// [`Lines`]: crate::Lines
297/// [`Builder`]: crate::Builder
298#[cfg_attr(
299    all(feature = "process", not(target_family = "wasm")),
300    doc = "[`AcpAgent`]: crate::AcpAgent"
301)]
302pub trait ConnectTo<R: Role>: Send + 'static {
303    /// Connect this component to another component.
304    ///
305    /// Most components implement this method to set up their connection and
306    /// exchange messages with the provided component.
307    ///
308    /// # Arguments
309    ///
310    /// * `client` - The component to connect to (implements `ConnectTo<R::Counterpart>`)
311    ///
312    /// # Returns
313    ///
314    /// A future that resolves when the connection ends, either successfully
315    /// or with an error. The future must be `Send`.
316    ///
317    /// A component that buffers outbound messages should not return `Ok(())`
318    /// merely because its client completed: it should first finish messages the
319    /// client already transferred to it. This lets wrappers preserve graceful
320    /// drain guarantees through to the physical transport sink. Errors may
321    /// still terminate the connection immediately.
322    fn connect_to(
323        self,
324        client: impl ConnectTo<R::Counterpart>,
325    ) -> impl Future<Output = Result<()>> + Send;
326
327    /// Convert this component into a channel endpoint and optional owned driver.
328    ///
329    /// The returned [`Channel`] is the canonical frame-aware boundary. It carries
330    /// complete [`TransportFrame`](crate::TransportFrame) values so default
331    /// adapters preserve batch grouping.
332    ///
333    /// This method returns:
334    /// - A `Channel` that can be used to communicate with this component
335    /// - `Some(ConnectionDriver)` when the component owns work to drive
336    /// - `None` when the channel halves alone own the endpoint's lifetime
337    ///
338    /// The default implementation creates an intermediate channel pair and calls `connect_to`
339    /// on one endpoint while returning the other endpoint for the caller to use.
340    ///
341    /// Base cases like `Channel` and `ByteStreams` override this to avoid unnecessary copying.
342    ///
343    /// # Returns
344    ///
345    /// A tuple of `(Channel, Option<ConnectionDriver>)`. Owned drivers must be
346    /// polled concurrently with channel traffic. Successful owned completion
347    /// ends the endpoint after draining accepted output. `None` is not EOF:
348    /// preserve both independent channel half-closes.
349    ///
350    /// Absence must be handled explicitly; the optional driver is not awaitable:
351    ///
352    /// ```compile_fail,E0277
353    /// use agent_client_protocol::{Channel, ConnectTo, UntypedRole};
354    ///
355    /// # async fn example() -> agent_client_protocol::Result<()> {
356    /// let (channel, _peer) = Channel::duplex();
357    /// let (_channel, driver) = ConnectTo::<UntypedRole>::into_channel_and_future(channel);
358    /// driver.await?;
359    /// # Ok(())
360    /// # }
361    /// ```
362    ///
363    /// Once present, the owned driver itself is awaitable:
364    ///
365    /// ```no_run
366    /// use agent_client_protocol::{Channel, ConnectionDriver, Result};
367    ///
368    /// async fn drive_owned_work((_channel, driver): (Channel, Option<ConnectionDriver>)) -> Result<()> {
369    ///     if let Some(driver) = driver {
370    ///         // In a real adapter, also poll the channel traffic concurrently.
371    ///         driver.await?;
372    ///     }
373    ///     Ok(())
374    /// }
375    /// ```
376    fn into_channel_and_future(self) -> (Channel, Option<ConnectionDriver>)
377    where
378        Self: Sized,
379    {
380        let (channel_a, channel_b) = Channel::duplex();
381        let future = ConnectionDriver::new(self.connect_to(channel_b));
382        (channel_a, Some(future))
383    }
384}
385
386/// Type-erased connect trait for object-safe dynamic dispatch.
387///
388/// This trait is internal and used by [`DynConnectTo`]. Users should implement
389/// [`ConnectTo`] instead, which is automatically converted to `ErasedConnectTo`
390/// via a blanket implementation.
391trait ErasedConnectTo<R: Role>: Send {
392    fn type_name(&self) -> &'static str;
393
394    fn connect_to_erased(
395        self: Box<Self>,
396        client: Box<dyn ErasedConnectTo<R::Counterpart>>,
397    ) -> BoxFuture<'static, Result<()>>;
398
399    fn into_channel_and_future_erased(self: Box<Self>) -> (Channel, Option<ConnectionDriver>);
400}
401
402/// Blanket implementation: any `ConnectTo<R>` can be type-erased.
403impl<C: ConnectTo<R>, R: Role> ErasedConnectTo<R> for C {
404    fn type_name(&self) -> &'static str {
405        std::any::type_name::<C>()
406    }
407
408    fn connect_to_erased(
409        self: Box<Self>,
410        client: Box<dyn ErasedConnectTo<R::Counterpart>>,
411    ) -> BoxFuture<'static, Result<()>> {
412        Box::pin(async move {
413            (*self)
414                .connect_to(DynConnectTo {
415                    inner: client,
416                    _marker: PhantomData,
417                })
418                .await
419        })
420    }
421
422    fn into_channel_and_future_erased(self: Box<Self>) -> (Channel, Option<ConnectionDriver>) {
423        (*self).into_channel_and_future()
424    }
425}
426
427/// A dynamically-typed component for heterogeneous collections.
428///
429/// This type wraps any [`ConnectTo`] implementation and provides dynamic dispatch,
430/// allowing you to store different component types in the same collection.
431///
432/// The type parameter `R` is the role that all components in the
433/// collection connect to (their counterpart).
434///
435/// # Examples
436///
437/// ```rust
438/// use agent_client_protocol::{Channel, Client, DynConnectTo};
439///
440/// let (first, _first_peer) = Channel::duplex();
441/// let (second, _second_peer) = Channel::duplex();
442/// let components: Vec<DynConnectTo<Client>> = vec![
443///     DynConnectTo::new(first),
444///     DynConnectTo::new(second),
445/// ];
446/// assert_eq!(components.len(), 2);
447/// ```
448pub struct DynConnectTo<R: Role> {
449    inner: Box<dyn ErasedConnectTo<R>>,
450    _marker: PhantomData<R>,
451}
452
453impl<R: Role> DynConnectTo<R> {
454    /// Create a new `DynConnectTo` from any type implementing [`ConnectTo`].
455    pub fn new<C: ConnectTo<R>>(component: C) -> Self {
456        Self {
457            inner: Box::new(component),
458            _marker: PhantomData,
459        }
460    }
461
462    /// Returns the type name of the wrapped component.
463    #[must_use]
464    pub fn type_name(&self) -> &'static str {
465        self.inner.type_name()
466    }
467}
468
469impl<R: Role> ConnectTo<R> for DynConnectTo<R> {
470    async fn connect_to(self, client: impl ConnectTo<R::Counterpart>) -> Result<()> {
471        self.inner
472            .connect_to_erased(Box::new(client) as Box<dyn ErasedConnectTo<R::Counterpart>>)
473            .await
474    }
475
476    fn into_channel_and_future(self) -> (Channel, Option<ConnectionDriver>) {
477        self.inner.into_channel_and_future_erased()
478    }
479}
480
481impl<R: Role> Debug for DynConnectTo<R> {
482    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
483        f.debug_struct("DynConnectTo")
484            .field("type_name", &self.type_name())
485            .finish()
486    }
487}
488
489#[cfg(test)]
490mod tests {
491    use super::*;
492    use crate::role::UntypedRole;
493    use futures::FutureExt as _;
494
495    struct OwnedWork(BoxFuture<'static, Result<()>>);
496
497    impl ConnectTo<UntypedRole> for OwnedWork {
498        async fn connect_to(self, _client: impl ConnectTo<UntypedRole>) -> Result<()> {
499            self.0.await
500        }
501    }
502
503    #[test]
504    fn raw_channel_has_no_owned_work() {
505        let (channel, _other) = Channel::duplex();
506        let (_, driver) = ConnectTo::<UntypedRole>::into_channel_and_future(channel);
507        assert!(driver.is_none());
508    }
509
510    #[test]
511    fn owned_driver_preserves_errors_and_polls_unpinned() {
512        let error = crate::Error::internal_error().data("driver failure");
513        let mut driver = ConnectionDriver::new(futures::future::ready(Err(error.clone())));
514        assert_eq!(futures::executor::block_on(&mut driver), Err(error));
515    }
516
517    #[test]
518    fn finish_request_is_idempotent_and_does_not_mean_completion() {
519        let (finish_tx, finish_rx) = futures::channel::oneshot::channel();
520        let (flushed_tx, flushed_rx) = futures::channel::oneshot::channel();
521        let mut driver = ConnectionDriver::with_finish(
522            async move {
523                finish_rx.await.unwrap();
524                flushed_rx.await.unwrap()
525            },
526            move || finish_tx.send(()).unwrap(),
527        );
528
529        assert!((&mut driver).now_or_never().is_none());
530        assert!(driver.request_finish());
531        assert!(driver.request_finish());
532        assert!((&mut driver).now_or_never().is_none());
533
534        let error = crate::Error::internal_error().data("custom flush failed");
535        flushed_tx.send(Err(error.clone())).unwrap();
536        assert_eq!(futures::executor::block_on(driver), Err(error));
537    }
538
539    #[test]
540    fn opaque_driver_cannot_be_cooperatively_finished() {
541        let mut driver = ConnectionDriver::new(futures::future::pending());
542        assert!(!driver.request_finish());
543        assert!((&mut driver).now_or_never().is_none());
544    }
545
546    #[test]
547    fn future_decoration_preserves_finish_and_completion_errors() {
548        use std::sync::{
549            Arc,
550            atomic::{AtomicUsize, Ordering},
551        };
552
553        for request_before_wrapping in [false, true] {
554            let (finish_tx, finish_rx) = futures::channel::oneshot::channel();
555            let (flush_tx, flush_rx) = futures::channel::oneshot::channel();
556            let calls = Arc::new(AtomicUsize::new(0));
557            let hook_calls = calls.clone();
558            let mut driver = ConnectionDriver::with_finish(
559                async move {
560                    finish_rx.await.unwrap();
561                    flush_rx.await.unwrap()
562                },
563                move || {
564                    hook_calls.fetch_add(1, Ordering::SeqCst);
565                    finish_tx.send(()).unwrap();
566                },
567            );
568            if request_before_wrapping {
569                assert!(driver.request_finish());
570            }
571
572            let observed = Arc::new(AtomicUsize::new(0));
573            let observe_completion = observed.clone();
574            let mut decorated = driver.map_future(|work| {
575                work.inspect(move |_| {
576                    observe_completion.fetch_add(1, Ordering::SeqCst);
577                })
578            });
579            assert!(decorated.request_finish());
580            assert!(decorated.request_finish());
581            assert_eq!(calls.load(Ordering::SeqCst), 1);
582            assert!((&mut decorated).now_or_never().is_none());
583            assert_eq!(observed.load(Ordering::SeqCst), 0);
584
585            let error = crate::Error::internal_error().data("decorated flush failed");
586            flush_tx.send(Err(error.clone())).unwrap();
587            assert_eq!(futures::executor::block_on(decorated), Err(error));
588            assert_eq!(observed.load(Ordering::SeqCst), 1);
589        }
590    }
591
592    #[test]
593    fn future_decoration_does_not_make_opaque_work_cooperative() {
594        let driver = ConnectionDriver::new(futures::future::pending());
595        let mut decorated = driver.map_future(|work| work);
596
597        assert!(!decorated.request_finish());
598        assert!((&mut decorated).now_or_never().is_none());
599    }
600
601    #[test]
602    fn dropping_driver_does_not_invoke_finish_hook() {
603        let invoked = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
604        let hook_invoked = invoked.clone();
605        let driver = ConnectionDriver::with_finish(futures::future::pending(), move || {
606            hook_invoked.store(true, std::sync::atomic::Ordering::Release);
607        });
608
609        drop(driver);
610        assert!(!invoked.load(std::sync::atomic::Ordering::Acquire));
611    }
612
613    #[test]
614    fn default_conversion_owns_real_work_until_completion() {
615        let (done_tx, done_rx) = futures::channel::oneshot::channel();
616        let component = OwnedWork(async move { done_rx.await.unwrap() }.boxed());
617        let (_channel, driver) = component.into_channel_and_future();
618        let mut driver = driver.expect("default conversion always owns its connect_to work");
619        assert!((&mut driver).now_or_never().is_none());
620
621        let error = crate::Error::internal_error().data("owned work failed");
622        done_tx.send(Err(error.clone())).unwrap();
623        assert_eq!(futures::executor::block_on(driver), Err(error));
624    }
625
626    #[test]
627    fn dropping_optional_owned_driver_cancels_unpolled_work() {
628        let (done_tx, done_rx) = futures::channel::oneshot::channel::<Result<()>>();
629        let component = OwnedWork(async move { done_rx.await.unwrap() }.boxed());
630        let (_channel, driver) = component.into_channel_and_future();
631        assert!(driver.is_some());
632        assert!(!done_tx.is_canceled());
633        drop(driver);
634        assert!(done_tx.is_canceled());
635    }
636
637    #[test]
638    fn type_erasure_preserves_owned_work_and_finish_metadata() {
639        let outgoing = futures::sink::unfold((), |(), _line: String| {
640            futures::future::ready(Ok::<_, std::io::Error>(()))
641        });
642        // Independent physical input remains open: only a preserved explicit
643        // finish handle can complete this driver without read EOF.
644        let incoming = futures::stream::pending::<std::io::Result<String>>();
645        let component = DynConnectTo::<UntypedRole>::new(crate::Lines::new(outgoing, incoming));
646        let (_channel, driver) = component.into_channel_and_future();
647        let mut driver = driver.expect("erasure must retain ownership");
648        assert!((&mut driver).now_or_never().is_none());
649        assert!(
650            driver.request_finish(),
651            "erasure must retain finish coordination"
652        );
653        futures::executor::block_on(driver).unwrap();
654    }
655
656    #[test]
657    fn type_erasure_preserves_passive_lifetime() {
658        let (channel, _other) = Channel::duplex();
659        let (_, driver) = DynConnectTo::<UntypedRole>::new(channel).into_channel_and_future();
660        assert!(driver.is_none());
661    }
662
663    #[test]
664    fn dyn_connect_to_reports_static_type_name_and_correct_debug_label() {
665        let (channel, _other) = Channel::duplex();
666        let component = DynConnectTo::<UntypedRole>::new(channel);
667
668        let type_name: &'static str = component.type_name();
669        assert_eq!(type_name, std::any::type_name::<Channel>());
670        assert_eq!(
671            format!("{component:?}"),
672            format!("DynConnectTo {{ type_name: {type_name:?} }}")
673        );
674    }
675}