use core::{fmt, future::Future};
use super::{
AsyncResponseStaging, DeliveryPhase, ResponseCompletion, ResponseWriter, ResponseWriterError,
TransportRequest,
};
pub const ASYNC_CANCELLATION_DELIVERY_PHASE: DeliveryPhase = DeliveryPhase::PossiblySent;
#[derive(Clone, Copy, Eq, PartialEq)]
pub enum AsyncExecutionError<E> {
Transport(E),
Response(ResponseWriterError),
}
impl<E> fmt::Debug for AsyncExecutionError<E> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Transport(_) => formatter.write_str("Transport([redacted])"),
Self::Response(error) => formatter.debug_tuple("Response").field(error).finish(),
}
}
}
impl<E> fmt::Display for AsyncExecutionError<E> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(match self {
Self::Transport(_) => "asynchronous transport failed",
Self::Response(_) => "asynchronous response staging failed",
})
}
}
impl<E> core::error::Error for AsyncExecutionError<E> {}
pub trait LocalAsyncTransport {
type Error;
fn send_local<'transport, 'request, 'writer, 'buffer>(
&'transport self,
request: TransportRequest<'request>,
response: AsyncResponseStaging<'writer, 'buffer>,
) -> impl Future<Output = Result<ResponseCompletion, Self::Error>> + 'writer
where
'transport: 'writer,
'request: 'writer,
'buffer: 'writer;
}
pub async fn drive_local<'transport, 'request, 'writer, 'buffer, T>(
transport: &'transport T,
request: TransportRequest<'request>,
response: &'writer mut ResponseWriter<'buffer>,
) -> Result<(), AsyncExecutionError<T::Error>>
where
T: LocalAsyncTransport + ?Sized,
'transport: 'writer,
'request: 'writer,
'buffer: 'writer,
{
let mut attempt = response
.begin_attempt()
.map_err(AsyncExecutionError::Response)?;
let completion = transport
.send_local(request, attempt.staging())
.await
.map_err(AsyncExecutionError::Transport)?;
attempt
.commit_completion(completion)
.map_err(AsyncExecutionError::Response)
}
pub trait AsyncTransport {
type Error;
fn send<'transport, 'request, 'writer, 'buffer>(
&'transport self,
request: TransportRequest<'request>,
response: AsyncResponseStaging<'writer, 'buffer>,
) -> impl Future<Output = Result<ResponseCompletion, Self::Error>> + Send + 'writer
where
'transport: 'writer,
'request: 'writer,
'buffer: 'writer;
}
pub async fn drive_async<'transport, 'request, 'writer, 'buffer, T>(
transport: &'transport T,
request: TransportRequest<'request>,
response: &'writer mut ResponseWriter<'buffer>,
) -> Result<(), AsyncExecutionError<T::Error>>
where
T: AsyncTransport + ?Sized,
'transport: 'writer,
'request: 'writer,
'buffer: 'writer,
{
let mut attempt = response
.begin_attempt()
.map_err(AsyncExecutionError::Response)?;
let completion = transport
.send(request, attempt.staging())
.await
.map_err(AsyncExecutionError::Transport)?;
attempt
.commit_completion(completion)
.map_err(AsyncExecutionError::Response)
}
impl<T> LocalAsyncTransport for T
where
T: AsyncTransport + ?Sized,
{
type Error = T::Error;
async fn send_local<'transport, 'request, 'writer, 'buffer>(
&'transport self,
request: TransportRequest<'request>,
response: AsyncResponseStaging<'writer, 'buffer>,
) -> Result<ResponseCompletion, Self::Error>
where
'transport: 'writer,
'request: 'writer,
'buffer: 'writer,
{
AsyncTransport::send(self, request, response).await
}
}