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}