Skip to main content

async_rs/traits/
reactor.rs

1//! A collection of traits to define a common interface across reactors
2
3use crate::{sys::AsSysFd, traits::AsyncToSocketAddrs};
4use futures_core::Stream;
5use futures_io::{AsyncRead, AsyncWrite};
6use std::{
7    io::{self, Read, Write},
8    net::SocketAddr,
9    ops::Deref,
10    time::{Duration, Instant},
11};
12
13/// A common interface for performing actions on a reactor
14pub trait Reactor {
15    /// The type representing a TCP stream (after tcp_connect) for this reactor
16    type TcpStream: AsyncRead + AsyncWrite + Send + Unpin + 'static;
17
18    /// A future that completes after a requested duration has elapsed (see [`Reactor::sleep`])
19    type Sleep: Future + Send + 'static;
20
21    /// Register a synchronous handle, returning an asynchronous one
22    ///
23    /// Whether the handle has to be in non-blocking mode already depends on the reactor: the
24    /// `async-io` based ones set it themselves, while the tokio one requires the caller to have
25    /// done so. Nothing checks, and getting it wrong is quiet: on a readiness notification that
26    /// does not pan out, or a write into a full send buffer, the read or write syscall blocks the
27    /// executor thread instead of reporting `WouldBlock`, stalling every task on it with no error
28    /// to go on. Set it yourself to stay portable across reactors.
29    fn register<H: Read + Write + AsSysFd + Send + 'static>(
30        &self,
31        socket: H,
32    ) -> io::Result<impl AsyncRead + AsyncWrite + Send + Unpin + 'static>
33    where
34        Self: Sized;
35
36    /// Sleep for the given duration
37    fn sleep(&self, dur: Duration) -> Self::Sleep
38    where
39        Self: Sized;
40
41    /// Stream that yields immediately, then at every given interval
42    ///
43    /// The noop reactor never yields an item.
44    fn interval(&self, dur: Duration) -> impl Stream<Item = Instant> + Send + 'static
45    where
46        Self: Sized;
47
48    /// Create a TcpStream by connecting to a remote host
49    fn tcp_connect<A: AsyncToSocketAddrs + Send>(
50        &self,
51        addrs: A,
52    ) -> impl Future<Output = io::Result<Self::TcpStream>> + Send
53    where
54        Self: Sync + Sized,
55    {
56        async move {
57            let mut err = None;
58            for addr in addrs.to_socket_addrs().await? {
59                match self.tcp_connect_addr(addr).await {
60                    Ok(stream) => return Ok(stream),
61                    Err(e) => err = Some(e),
62                }
63            }
64            Err(err.unwrap_or_else(|| {
65                io::Error::new(io::ErrorKind::AddrNotAvailable, "couldn't resolve host")
66            }))
67        }
68    }
69
70    /// Create a TcpStream by connecting to a specific pre-resolved address
71    fn tcp_connect_addr(
72        &self,
73        addr: SocketAddr,
74    ) -> impl Future<Output = io::Result<Self::TcpStream>> + Send + 'static
75    where
76        Self: Sized;
77}
78
79impl<R: Deref> Reactor for R
80where
81    R::Target: Reactor + Sized,
82{
83    type TcpStream = <<R as Deref>::Target as Reactor>::TcpStream;
84    type Sleep = <<R as Deref>::Target as Reactor>::Sleep;
85
86    fn register<H: Read + Write + AsSysFd + Send + 'static>(
87        &self,
88        socket: H,
89    ) -> io::Result<impl AsyncRead + AsyncWrite + Send + Unpin + 'static> {
90        self.deref().register(socket)
91    }
92
93    fn sleep(&self, dur: Duration) -> Self::Sleep {
94        self.deref().sleep(dur)
95    }
96
97    fn interval(&self, dur: Duration) -> impl Stream<Item = Instant> + Send + 'static {
98        self.deref().interval(dur)
99    }
100
101    fn tcp_connect_addr(
102        &self,
103        addr: SocketAddr,
104    ) -> impl Future<Output = io::Result<Self::TcpStream>> + Send + 'static {
105        self.deref().tcp_connect_addr(addr)
106    }
107}