Skip to main content

moirai_async/io/
compat.rs

1// These are used only by the `tokio-compat` trait-bridge impls below.
2#[cfg(feature = "tokio-compat")]
3use std::io;
4#[cfg(feature = "tokio-compat")]
5use std::pin::Pin;
6#[cfg(feature = "tokio-compat")]
7use std::task::{Context, Poll};
8
9#[cfg(feature = "tokio-compat")]
10use crate::io::traits::{AsyncRead, AsyncWrite};
11
12#[cfg(feature = "tokio-compat")]
13use tokio_dep as tokio;
14
15/// Wrapper providing Tokio's I/O traits compatibility.
16#[repr(transparent)]
17pub struct TokioCompat<T> {
18    inner: T,
19}
20
21impl<T> TokioCompat<T> {
22    /// Create a new Tokio compatibility wrapper.
23    pub fn new(inner: T) -> Self {
24        Self { inner }
25    }
26
27    /// Extract the inner type.
28    pub fn into_inner(self) -> T {
29        self.inner
30    }
31}
32
33impl<T> From<T> for TokioCompat<T> {
34    fn from(inner: T) -> Self {
35        Self::new(inner)
36    }
37}
38
39/// Wrapper providing Moirai's native I/O traits compatibility for Tokio types.
40#[repr(transparent)]
41pub struct MoiraiCompat<T> {
42    inner: T,
43}
44
45impl<T> MoiraiCompat<T> {
46    /// Create a new Moirai compatibility wrapper.
47    pub fn new(inner: T) -> Self {
48        Self { inner }
49    }
50
51    /// Extract the inner type.
52    pub fn into_inner(self) -> T {
53        self.inner
54    }
55}
56
57impl<T> From<T> for MoiraiCompat<T> {
58    fn from(inner: T) -> Self {
59        Self::new(inner)
60    }
61}
62
63#[cfg(feature = "tokio-compat")]
64impl<T: AsyncRead + Unpin> tokio::io::AsyncRead for TokioCompat<T> {
65    fn poll_read(
66        mut self: Pin<&mut Self>,
67        cx: &mut Context<'_>,
68        buf: &mut tokio::io::ReadBuf<'_>,
69    ) -> Poll<io::Result<()>> {
70        let unfilled = buf.initialize_unfilled();
71        match Pin::new(&mut self.inner).poll_read(cx, unfilled) {
72            Poll::Ready(Ok(n)) => {
73                buf.advance(n);
74                Poll::Ready(Ok(()))
75            }
76            Poll::Ready(Err(e)) => Poll::Ready(Err(e)),
77            Poll::Pending => Poll::Pending,
78        }
79    }
80}
81
82#[cfg(feature = "tokio-compat")]
83impl<T: AsyncWrite + Unpin> tokio::io::AsyncWrite for TokioCompat<T> {
84    fn poll_write(
85        mut self: Pin<&mut Self>,
86        cx: &mut Context<'_>,
87        buf: &[u8],
88    ) -> Poll<io::Result<usize>> {
89        Pin::new(&mut self.inner).poll_write(cx, buf)
90    }
91
92    fn poll_flush(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<io::Result<()>> {
93        Pin::new(&mut self.inner).poll_flush(cx)
94    }
95
96    fn poll_shutdown(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<io::Result<()>> {
97        Pin::new(&mut self.inner).poll_shutdown(cx)
98    }
99}
100
101#[cfg(feature = "tokio-compat")]
102impl<T: tokio::io::AsyncRead + Unpin> AsyncRead for MoiraiCompat<T> {
103    fn poll_read(
104        mut self: Pin<&mut Self>,
105        cx: &mut Context<'_>,
106        buf: &mut [u8],
107    ) -> Poll<io::Result<usize>> {
108        let mut read_buf = tokio::io::ReadBuf::new(buf);
109        match Pin::new(&mut self.inner).poll_read(cx, &mut read_buf) {
110            Poll::Ready(Ok(())) => Poll::Ready(Ok(read_buf.filled().len())),
111            Poll::Ready(Err(e)) => Poll::Ready(Err(e)),
112            Poll::Pending => Poll::Pending,
113        }
114    }
115}
116
117#[cfg(feature = "tokio-compat")]
118impl<T: tokio::io::AsyncWrite + Unpin> AsyncWrite for MoiraiCompat<T> {
119    fn poll_write(
120        mut self: Pin<&mut Self>,
121        cx: &mut Context<'_>,
122        buf: &[u8],
123    ) -> Poll<io::Result<usize>> {
124        Pin::new(&mut self.inner).poll_write(cx, buf)
125    }
126
127    fn poll_flush(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<io::Result<()>> {
128        Pin::new(&mut self.inner).poll_flush(cx)
129    }
130
131    fn poll_shutdown(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<io::Result<()>> {
132        Pin::new(&mut self.inner).poll_shutdown(cx)
133    }
134}