pub struct ConnectionDriver { /* private fields */ }Expand description
Drives owned endpoint work.
A driver owns the endpoint: successful completion means no further output is expected, and adapters must drain output already accepted before terminating. Errors may abort the connection without guaranteed output drain; an adapter may still preserve queued error replies before terminating.
Poll the driver concurrently with channel traffic. Endpoints without owned
work return None from ConnectTo::into_channel_and_future, not a driver:
their channel halves independently determine their lifetime.
Use new for opaque work or
with_finish for a transport that can finish gracefully.
Use map_future to decorate existing work without losing
its finish capability.
Implementations§
Source§impl ConnectionDriver
impl ConnectionDriver
Sourcepub fn new(future: impl Future<Output = Result<()>> + Send + 'static) -> Self
pub fn new(future: impl Future<Output = Result<()>> + Send + 'static) -> Self
Create a driver that owns the endpoint’s lifetime.
This driver has no cooperative finish hook. A finite foreground may drop it after handing off accepted output, rather than wait for arbitrary work to finish. Reactive serving still awaits owned work after input EOF.
Custom transports that need to flush before a finite foreground returns
should use with_finish instead.
Sourcepub fn with_finish(
future: impl Future<Output = Result<()>> + Send + 'static,
finish: impl FnOnce() + Send + 'static,
) -> Self
pub fn with_finish( future: impl Future<Output = Result<()>> + Send + 'static, finish: impl FnOnce() + Send + 'static, ) -> Self
Create owned work that supports cooperative graceful completion.
The finish hook only requests completion; it must be nonblocking and should signal the future to stop accepting output, drain what it has already accepted, flush and close its write half, then return. It must not require independently open remote input to reach EOF. The future remains responsible for reporting I/O and flush errors.
SDK consumers invoke the hook after handing off their accepted output, then continue polling the driver until completion. There is no implicit timeout: if the adapter cannot finish, the enclosing connection remains pending and may be cancelled by its caller.
The hook is invoked at most once. Dropping the driver drops its owned future without requesting graceful completion. Dropping only the hook does not invoke it or necessarily stop the work.
§Example
A custom adapter can use any signal understood by its future. For example, a one-shot channel separates the finish request from completion:
use agent_client_protocol::ConnectionDriver;
use futures::{channel::oneshot, FutureExt};
let (finish_tx, finish_rx) = oneshot::channel();
let mut driver = ConnectionDriver::with_finish(
async move {
if finish_rx.await.is_err() {
// Losing the hook must not masquerade as a finish request.
futures::future::pending::<()>().await;
}
// Seal the adapter's outgoing queue, drain it, and flush/close
// the physical writer here before returning.
Ok(())
},
move || { let _ = finish_tx.send(()); },
);
assert!((&mut driver).now_or_never().is_none());
assert!(driver.request_finish());
assert!(driver.request_finish()); // Supported, but the hook runs only once.
futures::executor::block_on(driver).unwrap();Sourcepub fn map_future<F>(
self,
map: impl FnOnce(BoxFuture<'static, Result<()>>) -> F,
) -> Self
pub fn map_future<F>( self, map: impl FnOnce(BoxFuture<'static, Result<()>>) -> F, ) -> Self
Decorate the owned future while preserving its finish capability.
This is useful for tracing, error annotation, or completion cleanup.
Wrapping this driver in new instead would hide its finish
control from the outer driver.
map is called immediately and receives the boxed future, not the
driver. Its returned future must uphold the same completion contract:
keep driving the original work and do not report success before accepted
output is drained. An already-requested finish remains requested, and
opaque work remains opaque.
use agent_client_protocol::ConnectionDriver;
use futures::FutureExt;
let driver = ConnectionDriver::new(async { Ok(()) });
let decorated = driver.map_future(|work| {
work.inspect(|result| eprintln!("transport completed: {result:?}"))
});
futures::executor::block_on(decorated).unwrap();Sourcepub fn request_finish(&mut self) -> bool
pub fn request_finish(&mut self) -> bool
Request graceful completion, without waiting for it.
Returns true if this driver supports cooperative finish, including
when finish was already requested. Repeated requests are idempotent:
the hook runs at most once and the driver retains its graceful-finish
contract across wrapping or ownership handoff.
Returns false for opaque work constructed with new;
this method does not cancel that work. A true return does not prove
flushing is complete: continue polling or await the driver to observe
completion and any errors.
Trait Implementations§
Source§impl Debug for ConnectionDriver
impl Debug for ConnectionDriver
Source§impl Future for ConnectionDriver
impl Future for ConnectionDriver
Auto Trait Implementations§
impl !RefUnwindSafe for ConnectionDriver
impl !Sync for ConnectionDriver
impl !UnwindSafe for ConnectionDriver
impl Freeze for ConnectionDriver
impl Send for ConnectionDriver
impl Unpin for ConnectionDriver
impl UnsafeUnpin for ConnectionDriver
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<T> FutureExt for T
impl<T> FutureExt for T
Source§fn map<U, F>(self, f: F) -> Map<Self, F> ⓘ
fn map<U, F>(self, f: F) -> Map<Self, F> ⓘ
Source§fn map_into<U>(self) -> MapInto<Self, U> ⓘ
fn map_into<U>(self) -> MapInto<Self, U> ⓘ
Source§fn then<Fut, F>(self, f: F) -> Then<Self, Fut, F> ⓘ
fn then<Fut, F>(self, f: F) -> Then<Self, Fut, F> ⓘ
f. Read moreSource§fn left_future<B>(self) -> Either<Self, B> ⓘ
fn left_future<B>(self) -> Either<Self, B> ⓘ
Source§fn right_future<A>(self) -> Either<A, Self> ⓘ
fn right_future<A>(self) -> Either<A, Self> ⓘ
Source§fn into_stream(self) -> IntoStream<Self>where
Self: Sized,
fn into_stream(self) -> IntoStream<Self>where
Self: Sized,
Source§fn flatten(self) -> Flatten<Self> ⓘ
fn flatten(self) -> Flatten<Self> ⓘ
Source§fn flatten_stream(self) -> FlattenStream<Self>
fn flatten_stream(self) -> FlattenStream<Self>
Source§fn fuse(self) -> Fuse<Self> ⓘwhere
Self: Sized,
fn fuse(self) -> Fuse<Self> ⓘwhere
Self: Sized,
poll will never again be called once it has
completed. This method can be used to turn any Future into a
FusedFuture. Read moreSource§fn inspect<F>(self, f: F) -> Inspect<Self, F> ⓘ
fn inspect<F>(self, f: F) -> Inspect<Self, F> ⓘ
Source§fn catch_unwind(self) -> CatchUnwind<Self> ⓘwhere
Self: Sized + UnwindSafe,
fn catch_unwind(self) -> CatchUnwind<Self> ⓘwhere
Self: Sized + UnwindSafe,
std only.std, or crate features alloc and spin only.Source§fn remote_handle(self) -> (Remote<Self>, RemoteHandle<Self::Output>)where
Self: Sized,
fn remote_handle(self) -> (Remote<Self>, RemoteHandle<Self::Output>)where
Self: Sized,
channel and std only.() on completion and sends
its output to another future on a separate task. Read moreSource§fn boxed<'a>(self) -> Pin<Box<dyn Future<Output = Self::Output> + Send + 'a>>
fn boxed<'a>(self) -> Pin<Box<dyn Future<Output = Self::Output> + Send + 'a>>
alloc only.Source§fn boxed_local<'a>(self) -> Pin<Box<dyn Future<Output = Self::Output> + 'a>>where
Self: Sized + 'a,
fn boxed_local<'a>(self) -> Pin<Box<dyn Future<Output = Self::Output> + 'a>>where
Self: Sized + 'a,
alloc only.Source§fn unit_error(self) -> UnitError<Self> ⓘwhere
Self: Sized,
fn unit_error(self) -> UnitError<Self> ⓘwhere
Self: Sized,
Future<Output = T> into a
TryFuture<Ok = T, Error = ()>.Source§fn never_error(self) -> NeverError<Self> ⓘwhere
Self: Sized,
fn never_error(self) -> NeverError<Self> ⓘwhere
Self: Sized,
Future<Output = T> into a
TryFuture<Ok = T, Error = Never>.Source§impl<F1> FutureExt for F1where
F1: Future,
impl<F1> FutureExt for F1where
F1: Future,
Source§fn join<F2>(self, other: F2) -> Join2<F1, <F2 as IntoFuture>::IntoFuture>
fn join<F2>(self, other: F2) -> Join2<F1, <F2 as IntoFuture>::IntoFuture>
Source§fn race<T, S2>(self, other: S2) -> Race2<T, F1, <S2 as IntoFuture>::IntoFuture>
fn race<T, S2>(self, other: S2) -> Race2<T, F1, <S2 as IntoFuture>::IntoFuture>
Source§fn wait_until<D>(
self,
deadline: D,
) -> WaitUntil<Self, <D as IntoFuture>::IntoFuture> ⓘwhere
Self: Sized,
D: IntoFuture,
fn wait_until<D>(
self,
deadline: D,
) -> WaitUntil<Self, <D as IntoFuture>::IntoFuture> ⓘwhere
Self: Sized,
D: IntoFuture,
Source§impl<F> FutureExt for F
impl<F> FutureExt for F
Source§fn catch_unwind(self) -> CatchUnwind<Self> ⓘwhere
Self: Sized + UnwindSafe,
fn catch_unwind(self) -> CatchUnwind<Self> ⓘwhere
Self: Sized + UnwindSafe,
std only.Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<F> IntoFuture for Fwhere
F: Future,
impl<F> IntoFuture for Fwhere
F: Future,
Source§type IntoFuture = F
type IntoFuture = F
Source§fn into_future(self) -> <F as IntoFuture>::IntoFuture
fn into_future(self) -> <F as IntoFuture>::IntoFuture
Source§impl<T> IntoMaybeUndefined<T> for T
impl<T> IntoMaybeUndefined<T> for T
Source§fn into_maybe_undefined(self) -> MaybeUndefined<T>
fn into_maybe_undefined(self) -> MaybeUndefined<T>
Source§impl<T> IntoOption<T> for T
impl<T> IntoOption<T> for T
Source§fn into_option(self) -> Option<T>
fn into_option(self) -> Option<T>
Source§impl<Fut> TryFutureExt for Fut
impl<Fut> TryFutureExt for Fut
Source§fn flatten_sink<Item>(self) -> FlattenSink<Self, Self::Ok>
fn flatten_sink<Item>(self) -> FlattenSink<Self, Self::Ok>
sink only.