cloud_sdk/transport/
asynchronous.rs1use core::{fmt, future::Future};
4
5use super::{
6 AsyncResponseStaging, DeliveryPhase, ResponseCompletion, ResponseWriter, ResponseWriterError,
7 TransportRequest,
8};
9
10pub const ASYNC_CANCELLATION_DELIVERY_PHASE: DeliveryPhase = DeliveryPhase::PossiblySent;
13
14#[derive(Clone, Copy, Eq, PartialEq)]
16pub enum AsyncExecutionError<E> {
17 Transport(E),
19 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
43pub trait LocalAsyncTransport {
96 type Error;
98
99 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
111pub 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
135pub trait AsyncTransport {
141 type Error;
143
144 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
156pub 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}