aboutsummaryrefslogtreecommitdiff
path: root/crates/arti-rpc-client-core/src/ll_conn.rs
blob: 45c7df50f24243a27b96b4b7736cb585ea83d7f2 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
//! Low-level connection implementations.
//!
//! This module defines two main types: [`NonblockingConnection`].
//! (a low-level type for use with external tools
//! that want to implement their own nonblocking IO),
//! and [`BlockingConnection`] (a slightly higher-level type
//! that we use internally when we are asked to provide
//! our own nonblocking IO loop(s)).
//!
//! This module also defines several traits for use by these types.
//!
//! Treats messages as unrelated strings, and validates outgoing messages for correctness.

mod blocking;
mod nonblocking;

use std::io;

#[cfg(unix)]
use std::os::fd::{AsFd as _, BorrowedFd as BorrowedOsHandle};
#[cfg(windows)]
use std::os::windows::io::{AsSocket as _, BorrowedSocket as BorrowedOsHandle};

pub(crate) use blocking::BlockingConnection;
pub(crate) use nonblocking::{NonblockingConnection, PollStatus, WriteHandle};

pub use nonblocking::{EventLoop, SendRequestError};

/// Retry `f` until it returns Ok() or an error whose kind is not `Interrupted`
fn retry_eintr<F, T>(mut f: F) -> io::Result<T>
where
    F: FnMut() -> io::Result<T>,
{
    loop {
        let r = f();
        match r {
            Err(e) if e.kind() == io::ErrorKind::Interrupted => continue,
            _ => return r,
        }
    }
}

/// Any type we can use as a target for [`NonblockingConnection`].
pub(crate) trait Stream: io::Read + io::Write + Send {
    /// If this Stream object is a [`MioStream`], return it as a `mio::event::Source`.
    ///
    /// Otherwise return None.
    fn as_mio_source(&mut self) -> Option<&mut dyn mio::event::Source>;

    /// Discard any mio-specific wrappers on this stream.
    fn remove_mio(self: Box<Self>) -> Box<dyn Stream>;

    /// Return an os-specific handle for using this stream type within a nonblocking event loop.
    ///
    /// (This will be an fd on unix and a SOCKET on windows.)
    fn try_as_handle(&self) -> io::Result<BorrowedOsHandle<'_>>;
}

/// A [`Stream`] that we can use inside a [`BlockingConnection`].
pub(crate) trait MioStream: Stream + mio::event::Source {}

/// Implement Stream and MioStream for a related pair of types.
macro_rules! impl_traits {
    { $stream:ty => $mio_stream:ty } => {
        impl Stream for $stream {
            fn as_mio_source(&mut self) -> Option<&mut dyn mio::event::Source> {
                None
            }
            fn remove_mio(self: Box<Self>) -> Box<dyn Stream> {
                self
            }
            fn try_as_handle(&self) -> io::Result<BorrowedOsHandle<'_>> {
                cfg_if::cfg_if!{
                    if #[cfg(unix)] {
                        Ok(self.as_fd())
                    } else if #[cfg(windows)] {
                        Ok(self.as_socket())
                    }
                }
            }
        }
        impl Stream for $mio_stream {
            fn as_mio_source(&mut self) -> Option<&mut dyn mio::event::Source> {
                Some(self as _)
            }
            fn remove_mio(self: Box<Self>) -> Box<dyn Stream> {
                Box::new(<$stream>::from(*self))
            }
            fn try_as_handle(&self) -> io::Result<BorrowedOsHandle<'_>> {
                cfg_if::cfg_if!{
                    if #[cfg(unix)] {
                        Ok(self.as_fd())
                    } else if #[cfg(windows)] {
                        Ok(self.as_socket())
                    }
                }
            }
        }
        impl MioStream for $mio_stream {
        }
    }
}

impl_traits! { std::net::TcpStream => mio::net::TcpStream }
#[cfg(unix)]
impl_traits! { std::os::unix::net::UnixStream => mio::net::UnixStream }

// We implement "Stream" for Empty so that we can use it to temporarily swap it in
// as a placeholder for a Box<dyn Stream>.
impl Stream for std::io::Empty {
    fn as_mio_source(&mut self) -> Option<&mut dyn mio::event::Source> {
        None
    }

    fn remove_mio(self: Box<Self>) -> Box<dyn Stream> {
        self
    }

    fn try_as_handle(&self) -> io::Result<BorrowedOsHandle<'_>> {
        Err(io::ErrorKind::Unsupported.into())
    }
}