Skip to main content

ntex_rt/
builder.rs

1use std::{fmt, future::Future, io, marker::PhantomData, panic, rc::Rc, sync::Arc, time};
2
3use crate::{driver::Runner, signals, system::System, system::SystemConfig};
4
5#[derive(Debug, Clone)]
6/// Builder for an ntex runtime system.
7///
8/// Use [`build`](Self::build) to create a [`SystemRunner`], then run the event
9/// loop or block on a future.
10pub struct Builder {
11    /// Name of the System. Defaults to "ntex" if unset.
12    name: String,
13    /// New thread stack size
14    stack_size: usize,
15    /// Arbiters ping interval
16    ping_interval: usize,
17    /// Arbiter ping response threshold
18    ping_threshold: usize,
19    /// Signal handling
20    signals: bool,
21    /// Panic handling
22    panics: bool,
23    /// Thread pool config
24    pool_limit: usize,
25    pool_recv_timeout: time::Duration,
26    /// testing flag
27    testing: bool,
28}
29
30impl Builder {
31    pub(super) fn new() -> Self {
32        Builder {
33            name: "ntex".into(),
34            stack_size: 0,
35            ping_interval: 2000,
36            ping_threshold: 1000,
37            signals: false,
38            panics: false,
39            testing: false,
40            pool_limit: 256,
41            pool_recv_timeout: time::Duration::from_mins(1),
42        }
43    }
44
45    #[must_use]
46    /// Sets the name of the System.
47    pub fn name<N: AsRef<str>>(mut self, name: N) -> Self {
48        self.name = name.as_ref().into();
49        self
50    }
51
52    #[must_use]
53    /// Enables or disables process signal handling.
54    ///
55    /// By default, signal handling is disabled.
56    pub fn signals(mut self, eanbled: bool) -> Self {
57        self.signals = eanbled;
58        self
59    }
60
61    #[must_use]
62    /// Enables panic handling.
63    ///
64    /// When panic handling is enabled, the application can receive
65    /// `Signal::Panic(PanicSource::App(..))` signals while signal handling
66    /// is enabled. The previously installed panic hook is still called.
67    /// By default, panic handling is disabled.
68    pub fn panic_handling(mut self, eanbled: bool) -> Self {
69        self.panics = eanbled;
70        self
71    }
72
73    #[doc(hidden)]
74    #[must_use]
75    /// Disable signal handling.
76    ///
77    /// By default, signal handling is disabled.
78    pub fn disable_signals(mut self) -> Self {
79        self.signals = false;
80        self
81    }
82
83    #[doc(hidden)]
84    #[must_use]
85    /// Enable signal handling.
86    ///
87    /// By default, signal handling is enabled.
88    pub fn enable_signals(mut self) -> Self {
89        self.signals = true;
90        self
91    }
92
93    #[must_use]
94    /// Sets the size of the stack (in bytes) for the new worker thread.
95    pub fn stack_size(mut self, size: usize) -> Self {
96        self.stack_size = size;
97        self
98    }
99
100    #[must_use]
101    /// Sets the ping interval for spawned arbiters.
102    ///
103    /// The interval is specified in milliseconds and defaults to 2,000.
104    /// Set it to zero to disable pings.
105    pub fn ping_interval(mut self, interval: usize) -> Self {
106        self.ping_interval = interval;
107        self
108    }
109
110    #[must_use]
111    /// Sets the ping response threshold.
112    ///
113    /// If a response takes too long, an attempt is made to create a backtrace
114    /// for the busy arbiter.
115    ///
116    /// The threshold is specified in milliseconds and defaults to 1,000.
117    pub fn ping_threshold(mut self, interval: usize) -> Self {
118        self.ping_threshold = interval;
119        self
120    }
121
122    #[must_use]
123    /// Sets the maximum number of blocking thread-pool workers.
124    ///
125    /// The default is 256.
126    pub fn thread_pool_limit(mut self, value: usize) -> Self {
127        self.pool_limit = value;
128        self
129    }
130
131    #[must_use]
132    /// Configures the system for testing.
133    ///
134    /// This disables signal and panic handling.
135    pub fn testing(mut self) -> Self {
136        self.testing = true;
137        self.signals = false;
138        self.panics = false;
139        self
140    }
141
142    #[must_use]
143    /// Sets how long an idle blocking worker waits before exiting.
144    ///
145    /// The default is 60 seconds.
146    pub fn thread_pool_recv_timeout<T>(mut self, timeout: T) -> Self
147    where
148        time::Duration: From<T>,
149    {
150        self.pool_recv_timeout = timeout.into();
151        self
152    }
153
154    /// Creates a system runner using the specified runtime runner.
155    ///
156    /// # Panics
157    ///
158    /// Panics if the runtime cannot be created.
159    pub fn build<R: Runner>(self, runner: R) -> SystemRunner {
160        let config = SystemConfig {
161            name: self.name.clone(),
162            testing: self.testing,
163            stack_size: self.stack_size,
164            ping_interval: self.ping_interval,
165            ping_threshold: self.ping_threshold,
166            pool_limit: self.pool_limit,
167            pool_recv_timeout: self.pool_recv_timeout,
168            runner: Arc::new(runner),
169        };
170        self.build_with(config)
171    }
172
173    /// Creates a system runner from an existing system configuration.
174    ///
175    /// # Panics
176    ///
177    /// Panics if the runtime cannot be created.
178    pub fn build_with(self, config: SystemConfig) -> SystemRunner {
179        let runner = config.runner.clone();
180
181        // init system arbiter and run configuration method
182        SystemRunner {
183            config,
184            runner,
185            signals: self.signals,
186            panics: self.panics,
187            _t: PhantomData,
188        }
189    }
190}
191
192/// A configured system that has not yet started its event loop.
193#[must_use = "SystemRunner must be run"]
194pub struct SystemRunner {
195    config: SystemConfig,
196    runner: Arc<dyn Runner>,
197    signals: bool,
198    panics: bool,
199    _t: PhantomData<Rc<()>>,
200}
201
202impl SystemRunner {
203    /// Runs the event loop until [`System::stop()`] is called.
204    pub fn run_until_stop(self) -> io::Result<()> {
205        self.run(|| Ok(()))
206    }
207
208    /// Runs `f`, then drives the event loop until [`System::stop()`] is called.
209    pub fn run<F>(self, f: F) -> io::Result<()>
210    where
211        F: FnOnce() -> io::Result<()> + 'static,
212    {
213        log::info!("Starting {:?} system", self.config.name);
214
215        let SystemRunner {
216            config,
217            runner,
218            signals,
219            panics,
220            ..
221        } = self;
222
223        if panics {
224            signals::enable_panic_handling();
225        }
226
227        // run loop
228        crate::driver::block_on_panic(runner.as_ref(), async move {
229            let (system, stop) = System::start(config);
230            let _signals = SignalsGuard(system.clone());
231            if signals {
232                system.enable_signals();
233            }
234
235            f()?;
236
237            let result = stop.await;
238
239            match result {
240                Ok(code) => {
241                    if code != 0 {
242                        Err(io::Error::other(format!("Non-zero exit code: {code}")))
243                    } else {
244                        Ok(())
245                    }
246                }
247                Err(_) => Err(io::Error::other("Closed")),
248            }
249        })
250    }
251
252    #[allow(clippy::missing_panics_doc)]
253    /// Execute a future and wait for result.
254    pub fn block_on<F, R>(self, fut: F) -> R
255    where
256        F: Future<Output = R> + 'static,
257        R: 'static,
258    {
259        let SystemRunner {
260            config,
261            runner,
262            signals,
263            panics,
264            ..
265        } = self;
266
267        if panics {
268            signals::enable_panic_handling();
269        }
270
271        crate::driver::block_on_panic(runner.as_ref(), async move {
272            let (system, _) = System::start(config);
273            let _signals = SignalsGuard(system.clone());
274            if signals {
275                system.enable_signals();
276            }
277
278            let loc = current_location();
279            ntex_error::set_backtrace_start(loc.file(), loc.line() + 2);
280            fut.await
281        })
282    }
283
284    #[cfg(feature = "tokio")]
285    /// Execute a future and wait for result.
286    pub async fn run_local<F, R>(self, fut: F) -> R
287    where
288        F: Future<Output = R> + 'static,
289        R: 'static,
290    {
291        let SystemRunner { config, .. } = self;
292
293        // run loop
294        let result = tok_io::task::LocalSet::new()
295            .run_until(async move {
296                let (system, _) = System::start(config);
297                let _signals = SignalsGuard(system);
298
299                let loc = current_location();
300                ntex_error::set_backtrace_start(loc.file(), loc.line() + 2);
301                fut.await
302            })
303            .await;
304
305        crate::arbiter::run_shutdown_callbacks();
306        unsafe {
307            crate::remove_all_items();
308        }
309        result
310    }
311}
312
313/// Releases the signals for other systems when the system's future completes,
314/// fails or panics.
315struct SignalsGuard(System);
316
317impl Drop for SignalsGuard {
318    fn drop(&mut self) {
319        self.0.disable_signals();
320    }
321}
322
323#[track_caller]
324pub(crate) fn current_location() -> &'static panic::Location<'static> {
325    panic::Location::caller()
326}
327
328impl fmt::Debug for SystemRunner {
329    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
330        f.debug_struct("SystemRunner")
331            .field("config", &self.config)
332            .finish()
333    }
334}
335
336#[cfg(test)]
337mod tests {
338    use std::{cell::Cell, cell::RefCell, time::Duration};
339
340    use super::*;
341    use crate::{Arbiter, testing::TestRunner};
342
343    #[test]
344    fn builder_options() {
345        let builder = System::build()
346            .name("opts")
347            .signals(false)
348            .panic_handling(false)
349            .enable_signals()
350            .disable_signals()
351            .stack_size(2 * 1024 * 1024)
352            .ping_interval(0)
353            .ping_threshold(10)
354            .thread_pool_limit(1)
355            .thread_pool_recv_timeout(Duration::from_millis(100))
356            .testing();
357        assert!(format!("{builder:?}").contains("opts"));
358
359        let runner = builder.build(TestRunner);
360        assert!(format!("{runner:?}").contains("opts"));
361        runner.block_on(async {
362            let sys = System::current();
363            assert_eq!(sys.name(), "opts");
364            assert!(sys.testing());
365            assert!(!sys.signals());
366
367            // arbiter thread with configured stack size
368            let mut arb = Arbiter::new();
369            arb.stop();
370            arb.join().unwrap();
371
372            assert_eq!(sys.spawn_blocking(|| 1).await, Ok(1));
373        });
374    }
375
376    #[test]
377    fn run_exit_codes() {
378        // `f` error is returned
379        let err = System::new("test", TestRunner)
380            .run(|| Err(io::Error::other("init failed")))
381            .unwrap_err();
382        assert_eq!(err.to_string(), "init failed");
383
384        System::new("test", TestRunner)
385            .run(|| {
386                System::current().stop();
387                Ok(())
388            })
389            .unwrap();
390
391        // stopping the system stops arbiters and runs shutdown callbacks
392        let called = Rc::new(Cell::new(false));
393        let arb = Rc::new(RefCell::new(None));
394        let (called2, arb2) = (called.clone(), arb.clone());
395        let err = System::new("test", TestRunner)
396            .run(move || {
397                Arbiter::on_shutdown(move || called2.set(true));
398                *arb2.borrow_mut() = Some(Arbiter::new());
399                System::current().stop_with_code(3);
400                Ok(())
401            })
402            .unwrap_err();
403        assert_eq!(err.to_string(), "Non-zero exit code: 3");
404        assert!(called.get());
405
406        let mut arb = arb.borrow_mut().take().unwrap();
407        arb.join().unwrap();
408        assert!(!arb.is_running());
409    }
410}