Skip to main content

ntex_net/
lib.rs

1//! Network transports and runtime reactor integration for ntex.
2//!
3//! This crate provides TCP and Unix domain socket connection helpers, DNS-aware
4//! connectors, and reactor implementations for native ntex, Tokio, and Compio
5//! runtimes. All transports are exposed as [`Io`] values.
6//!
7//! [`DefaultRuntime`] selects a reactor from the enabled Cargo features and the
8//! current platform.
9//!
10//! # Runtime features
11//!
12//! - `tokio` enables the Tokio reactor.
13//! - `compio` enables the Compio reactor.
14//! - `neon-polling` requires the native `polling` reactor on Unix.
15//! - `neon-uring` requires the native `io-uring` reactor on Linux.
16//!
17//! Without an explicit runtime feature, Linux tries `io-uring` and falls back to
18//! `polling`; other Unix platforms use polling and Windows always uses `IOCP`.
19#![deny(clippy::pedantic)]
20#![allow(
21    clippy::clone_on_copy,
22    clippy::cast_possible_truncation,
23    clippy::missing_fields_in_debug,
24    clippy::must_use_candidate,
25    clippy::missing_errors_doc,
26    clippy::missing_panics_doc,
27    clippy::unused_async_trait_impl
28)]
29use std::{any::Any, io, net, net::SocketAddr, panic};
30
31use ntex_io::Io;
32use ntex_rt::{BlockFuture, Driver, Runner};
33use ntex_service::cfg::SharedCfg;
34
35pub mod channel;
36pub mod connect;
37
38#[cfg(unix)]
39pub mod polling;
40
41#[cfg(target_os = "linux")]
42pub mod uring;
43
44#[cfg(windows)]
45pub mod iocp;
46
47#[cfg(any(unix, windows))]
48mod helpers;
49
50#[cfg(feature = "tokio")]
51pub mod tokio;
52
53#[cfg(feature = "compio")]
54pub mod compio;
55
56/// Runtime reactor used for network connection and transport operations.
57#[allow(clippy::wrong_self_convention)]
58pub trait Reactor: Driver {
59    /// Starts an asynchronous TCP connection.
60    fn tcp_connect(&self, addr: net::SocketAddr, cfg: SharedCfg) -> channel::Receiver<Io>;
61
62    /// Starts an asynchronous Unix domain socket connection.
63    fn unix_connect(&self, addr: std::path::PathBuf, cfg: SharedCfg) -> channel::Receiver<Io>;
64
65    /// Converts a standard-library TCP stream into an [`Io`] value.
66    fn from_tcp_stream(&self, stream: net::TcpStream, cfg: SharedCfg) -> io::Result<Io>;
67
68    #[cfg(unix)]
69    /// Converts a standard-library Unix stream into an [`Io`] value.
70    fn from_unix_stream(&self, _: std::os::unix::net::UnixStream, _: SharedCfg) -> io::Result<Io>;
71}
72
73#[inline]
74/// Opens a TCP connection to a remote host.
75pub async fn tcp_connect(addr: SocketAddr, cfg: SharedCfg) -> io::Result<Io> {
76    with_current(|driver| driver.tcp_connect(addr, cfg)).await
77}
78
79#[inline]
80/// Opens a Unix domain socket connection.
81pub async fn unix_connect<'a, P>(addr: P, cfg: SharedCfg) -> io::Result<Io>
82where
83    P: AsRef<std::path::Path> + 'a,
84{
85    with_current(|driver| driver.unix_connect(addr.as_ref().into(), cfg)).await
86}
87
88#[inline]
89/// Converts a standard-library TCP stream into an [`Io`] value.
90pub fn from_tcp_stream(stream: net::TcpStream, cfg: SharedCfg) -> io::Result<Io> {
91    with_current(|driver| driver.from_tcp_stream(stream, cfg))
92}
93
94#[cfg(unix)]
95#[inline]
96/// Converts a standard-library Unix stream into an [`Io`] value.
97pub fn from_unix_stream(stream: std::os::unix::net::UnixStream, cfg: SharedCfg) -> io::Result<Io> {
98    with_current(|driver| driver.from_unix_stream(stream, cfg))
99}
100
101fn with_current<T, F: FnOnce(&dyn Reactor) -> T>(f: F) -> T {
102    #[cold]
103    fn not_in_ntex_driver() -> ! {
104        panic!("not in a ntex driver")
105    }
106
107    if CURRENT_DRIVER.is_set() {
108        CURRENT_DRIVER.with(|d| f(&**d))
109    } else {
110        not_in_ntex_driver()
111    }
112}
113
114#[allow(clippy::borrowed_box)]
115/// Sets the current reactor and runs the provided closure.
116///
117/// # Panics
118///
119/// Panics if a reactor is already active on the current thread.
120pub fn with_reactor<R, F: FnOnce() -> R>(r: &Box<dyn Reactor>, f: F) -> R {
121    #[cold]
122    fn reactor_is_set() -> ! {
123        panic!("reactor is already set");
124    }
125
126    if CURRENT_DRIVER.is_set() {
127        reactor_is_set()
128    }
129    CURRENT_DRIVER.set(r, f)
130}
131
132scoped_tls::scoped_thread_local!(static CURRENT_DRIVER: Box<dyn Reactor>);
133
134/// The default runtime.
135///
136/// Automatically selects the runtime implementation based on the
137/// configured features and the platform on which it runs.
138#[derive(Copy, Clone, Debug)]
139pub struct DefaultRuntime;
140
141impl Runner for DefaultRuntime {
142    #[allow(unused_variables, clippy::too_many_lines)]
143    fn block_on(&self, fut: BlockFuture) -> Result<(), Box<dyn Any + Send>> {
144        #[cfg(feature = "tokio")]
145        {
146            let driver: Box<dyn Reactor> = Box::new(self::tokio::Reactor);
147
148            with_reactor(&driver, || crate::tokio::block_on(fut));
149            Ok(())
150        }
151
152        #[cfg(all(feature = "compio", not(feature = "tokio")))]
153        {
154            let driver: Box<dyn Reactor> = Box::new(self::compio::Reactor);
155
156            with_reactor(&driver, || crate::compio::block_on(fut));
157            Ok(())
158        }
159
160        #[cfg(all(windows, not(feature = "tokio"), not(feature = "compio")))]
161        {
162            let driver: Box<dyn Reactor> =
163                Box::new(crate::iocp::Reactor::new().expect("Cannot construct driver"));
164
165            with_reactor(&driver, || {
166                panic::catch_unwind(panic::AssertUnwindSafe(|| {
167                    let rt = ntex_rt::Runtime::new(driver.handle());
168                    rt.block_on(fut, &*driver);
169                }))
170            })
171        }
172
173        #[cfg(all(unix, not(feature = "tokio"), not(feature = "compio")))]
174        {
175            #[cfg(feature = "neon-polling")]
176            {
177                let driver: Box<dyn Reactor> = Box::new(
178                    crate::polling::Reactor::new().expect("Cannot construct polling reactor"),
179                );
180
181                with_reactor(&driver, || {
182                    panic::catch_unwind(panic::AssertUnwindSafe(|| {
183                        let rt = ntex_rt::Runtime::new(driver.handle());
184                        rt.block_on(fut, &*driver);
185                    }))
186                })
187            }
188
189            #[cfg(all(target_os = "linux", feature = "neon-uring"))]
190            {
191                let driver: Box<dyn Reactor> = Box::new(
192                    crate::uring::Reactor::new(2048).expect("Cannot construct io-uring reactor"),
193                );
194
195                with_reactor(&driver, || {
196                    panic::catch_unwind(panic::AssertUnwindSafe(|| {
197                        let rt = ntex_rt::Runtime::new(driver.handle());
198                        rt.block_on(fut, &*driver);
199                    }))
200                })
201            }
202
203            #[cfg(all(not(feature = "neon-uring"), not(feature = "neon-polling")))]
204            {
205                #[cfg(target_os = "linux")]
206                let driver: Box<dyn Reactor> = if let Ok(reactor) = crate::uring::Reactor::new(2048)
207                {
208                    Box::new(reactor)
209                } else {
210                    Box::new(
211                        crate::polling::Reactor::new().expect("Cannot construct io-uring reactor"),
212                    )
213                };
214
215                #[cfg(not(target_os = "linux"))]
216                let driver: Box<dyn Reactor> = Box::new(
217                    crate::polling::Reactor::new().expect("Cannot construct polling reactor"),
218                );
219
220                with_reactor(&driver, || {
221                    panic::catch_unwind(panic::AssertUnwindSafe(|| {
222                        let rt = ntex_rt::Runtime::new(driver.handle());
223                        rt.block_on(fut, &*driver);
224                    }))
225                })
226            }
227        }
228    }
229}