Skip to main content

cloud_sdk/transport/
asynchronous.rs

1//! Runtime-neutral asynchronous transport contracts.
2
3use core::{fmt, future::Future};
4
5use super::{
6    AsyncResponseStaging, DeliveryPhase, ResponseCompletion, ResponseWriter, ResponseWriterError,
7    TransportRequest,
8};
9
10/// Conservative delivery classification after an asynchronous future is
11/// cancelled by being dropped.
12pub const ASYNC_CANCELLATION_DELIVERY_PHASE: DeliveryPhase = DeliveryPhase::PossiblySent;
13
14/// Failure while SDK-owned asynchronous response staging is driven to completion.
15#[derive(Clone, Copy, Eq, PartialEq)]
16pub enum AsyncExecutionError<E> {
17    /// The transport failed before successful response completion.
18    Transport(E),
19    /// Response staging or final commitment failed.
20    Response(ResponseWriterError),
21}
22
23impl<E> fmt::Debug for AsyncExecutionError<E> {
24    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
25        match self {
26            Self::Transport(_) => formatter.write_str("Transport([redacted])"),
27            Self::Response(error) => formatter.debug_tuple("Response").field(error).finish(),
28        }
29    }
30}
31
32impl<E> fmt::Display for AsyncExecutionError<E> {
33    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
34        formatter.write_str(match self {
35            Self::Transport(_) => "asynchronous transport failed",
36            Self::Response(_) => "asynchronous response staging failed",
37        })
38    }
39}
40
41impl<E> core::error::Error for AsyncExecutionError<E> {}
42
43/// Local asynchronous transport for `!Send` futures and executors.
44///
45/// Implementations receive a non-committing staging view and return completion
46/// metadata. [`drive_local`] owns the cleanup attempt across the await and
47/// commits only after `Ready(Ok)`. Cancellation clears partial response state
48/// and remains conservatively [`DeliveryPhase::PossiblySent`].
49///
50/// A local implementation may deliberately return a future that is not
51/// `Send`:
52///
53/// ```compile_fail
54/// use core::cell::Cell;
55/// use cloud_sdk::transport::{
56///     AsyncResponseStaging, LocalAsyncTransport, ResponseCompletion,
57///     ResponseMetadata, StatusCode, TransportRequest,
58/// };
59///
60/// struct Local(Cell<()>);
61///
62/// impl LocalAsyncTransport for Local {
63///     type Error = ();
64///
65///     fn send_local<'transport, 'request, 'writer, 'buffer>(
66///         &'transport self,
67///         _request: TransportRequest<'request>,
68///         _response: AsyncResponseStaging<'writer, 'buffer>,
69///     ) -> impl core::future::Future<Output = Result<ResponseCompletion, Self::Error>> + 'writer
70///     where
71///         'transport: 'writer,
72///         'request: 'writer,
73///         'buffer: 'writer,
74///     {
75///         async move {
76///             self.0.get();
77///             Ok(ResponseCompletion::new(
78///                 StatusCode::NO_CONTENT,
79///                 0,
80///                 ResponseMetadata::EMPTY,
81///             ))
82///         }
83///     }
84/// }
85///
86/// fn require_send<T: Send>(_: T) {}
87/// fn reject_send<'a, 'buffer: 'a>(
88///     transport: &'a Local,
89///     request: TransportRequest<'a>,
90///     response: AsyncResponseStaging<'a, 'buffer>,
91/// ) {
92///     require_send(transport.send_local(request, response));
93/// }
94/// ```
95pub trait LocalAsyncTransport {
96    /// Transport-specific failure.
97    type Error;
98
99    /// Stages one response without requiring the returned future to be `Send`.
100    fn send_local<'transport, 'request, 'writer, 'buffer>(
101        &'transport self,
102        request: TransportRequest<'request>,
103        response: AsyncResponseStaging<'writer, 'buffer>,
104    ) -> impl Future<Output = Result<ResponseCompletion, Self::Error>> + 'writer
105    where
106        'transport: 'writer,
107        'request: 'writer,
108        'buffer: 'writer;
109}
110
111/// Drives one local transport attempt and commits only after `Ready(Ok)`.
112pub async fn drive_local<'transport, 'request, 'writer, 'buffer, T>(
113    transport: &'transport T,
114    request: TransportRequest<'request>,
115    response: &'writer mut ResponseWriter<'buffer>,
116) -> Result<(), AsyncExecutionError<T::Error>>
117where
118    T: LocalAsyncTransport + ?Sized,
119    'transport: 'writer,
120    'request: 'writer,
121    'buffer: 'writer,
122{
123    let mut attempt = response
124        .begin_attempt()
125        .map_err(AsyncExecutionError::Response)?;
126    let completion = transport
127        .send_local(request, attempt.staging())
128        .await
129        .map_err(AsyncExecutionError::Transport)?;
130    attempt
131        .commit_completion(completion)
132        .map_err(AsyncExecutionError::Response)
133}
134
135/// Cross-thread asynchronous transport over caller-owned buffers.
136///
137/// Implementations can stage body and headers but cannot commit a response.
138/// Callers use [`drive_async`], which owns cleanup across the await and commits
139/// returned completion metadata only after `Ready(Ok)`.
140pub trait AsyncTransport {
141    /// Transport-specific failure.
142    type Error;
143
144    /// Stages one complete response without committing it.
145    fn send<'transport, 'request, 'writer, 'buffer>(
146        &'transport self,
147        request: TransportRequest<'request>,
148        response: AsyncResponseStaging<'writer, 'buffer>,
149    ) -> impl Future<Output = Result<ResponseCompletion, Self::Error>> + Send + 'writer
150    where
151        'transport: 'writer,
152        'request: 'writer,
153        'buffer: 'writer;
154}
155
156/// Drives one cross-thread async attempt and commits only after `Ready(Ok)`.
157pub async fn drive_async<'transport, 'request, 'writer, 'buffer, T>(
158    transport: &'transport T,
159    request: TransportRequest<'request>,
160    response: &'writer mut ResponseWriter<'buffer>,
161) -> Result<(), AsyncExecutionError<T::Error>>
162where
163    T: AsyncTransport + ?Sized,
164    'transport: 'writer,
165    'request: 'writer,
166    'buffer: 'writer,
167{
168    let mut attempt = response
169        .begin_attempt()
170        .map_err(AsyncExecutionError::Response)?;
171    let completion = transport
172        .send(request, attempt.staging())
173        .await
174        .map_err(AsyncExecutionError::Transport)?;
175    attempt
176        .commit_completion(completion)
177        .map_err(AsyncExecutionError::Response)
178}
179
180impl<T> LocalAsyncTransport for T
181where
182    T: AsyncTransport + ?Sized,
183{
184    type Error = T::Error;
185
186    async fn send_local<'transport, 'request, 'writer, 'buffer>(
187        &'transport self,
188        request: TransportRequest<'request>,
189        response: AsyncResponseStaging<'writer, 'buffer>,
190    ) -> Result<ResponseCompletion, Self::Error>
191    where
192        'transport: 'writer,
193        'request: 'writer,
194        'buffer: 'writer,
195    {
196        AsyncTransport::send(self, request, response).await
197    }
198}