1#![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#[allow(clippy::wrong_self_convention)]
58pub trait Reactor: Driver {
59 fn tcp_connect(&self, addr: net::SocketAddr, cfg: SharedCfg) -> channel::Receiver<Io>;
61
62 fn unix_connect(&self, addr: std::path::PathBuf, cfg: SharedCfg) -> channel::Receiver<Io>;
64
65 fn from_tcp_stream(&self, stream: net::TcpStream, cfg: SharedCfg) -> io::Result<Io>;
67
68 #[cfg(unix)]
69 fn from_unix_stream(&self, _: std::os::unix::net::UnixStream, _: SharedCfg) -> io::Result<Io>;
71}
72
73#[inline]
74pub 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]
80pub 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]
89pub 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]
96pub 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)]
115pub 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#[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}