Skip to main content

ntex_rt/
system.rs

1use std::any::{Any, TypeId};
2use std::collections::VecDeque;
3use std::sync::{Arc, atomic::AtomicBool, atomic::AtomicUsize, atomic::Ordering};
4use std::time::{Duration, Instant};
5use std::{cell::RefCell, fmt, future::Future, panic, pin::Pin, rc::Rc};
6
7use async_channel::{Receiver, Sender, unbounded};
8use futures_timer::Delay;
9use parking_lot::{Mutex, RwLock};
10
11use crate::arbiter::Arbiter;
12#[cfg(target_os = "linux")]
13use crate::capture::CAPTURE;
14use crate::pool::ThreadPool;
15use crate::{BlockingResult, Builder, Handle, HashMap, HashSet, Runner, SystemRunner};
16
17static SYSTEM_COUNT: AtomicUsize = AtomicUsize::new(0);
18
19thread_local!(
20    static PINGS: RefCell<HashMap<Id, VecDeque<PingRecord>>> = RefCell::new(HashMap::default());
21);
22
23#[derive(Default)]
24struct Arbiters {
25    all: HashMap<Id, Arbiter>,
26    list: Vec<Arbiter>,
27}
28
29/// Identifier assigned to a running [`System`].
30#[derive(Copy, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Debug)]
31pub struct Id(pub(crate) usize);
32
33/// Runtime manager for a group of arbiter threads.
34///
35/// A system stores runtime configuration, manages arbiters, dispatches process
36/// signals, and owns the blocking thread pool.
37pub struct System(Arc<SystemInner>);
38
39struct SystemInner {
40    id: usize,
41    arbiter: Arbiter,
42    config: SystemConfig,
43    sender: Sender<SystemCommand>,
44    receiver: Receiver<SystemCommand>,
45    storage: RwLock<HashMap<TypeId, Box<dyn Any + Sync + Send>>>,
46    arbiters: Mutex<Arbiters>,
47    signals: AtomicBool,
48    pool: ThreadPool,
49}
50
51/// Configuration shared by a running [`System`].
52#[derive(Clone)]
53pub struct SystemConfig {
54    pub(super) name: String,
55    pub(super) stack_size: usize,
56    pub(super) ping_interval: usize,
57    #[allow(dead_code)]
58    pub(super) ping_threshold: usize,
59    pub(super) pool_limit: usize,
60    pub(super) pool_recv_timeout: Duration,
61    pub(super) testing: bool,
62    pub(super) runner: Arc<dyn Runner>,
63}
64
65thread_local!(
66    static CURRENT: RefCell<Option<System>> = const { RefCell::new(None) };
67);
68
69impl Clone for System {
70    fn clone(&self) -> Self {
71        Self(self.0.clone())
72    }
73}
74
75impl System {
76    /// Constructs new system and sets it as current
77    pub(super) fn start(config: SystemConfig) -> (Self, oneshot::Receiver<i32>) {
78        let id = SYSTEM_COUNT.fetch_add(1, Ordering::SeqCst);
79        let (sender, receiver) = unbounded();
80
81        let pool = ThreadPool::new(&config.name, config.pool_limit, config.pool_recv_timeout);
82        let (arbiter, controller) = Arbiter::new_system(id, config.name.clone());
83
84        let mut arbiters = Arbiters::default();
85        arbiters.all.insert(arbiter.id(), arbiter.clone());
86        arbiters.list.push(arbiter.clone());
87
88        let sys = System(Arc::new(SystemInner {
89            id,
90            config,
91            arbiter,
92            sender,
93            receiver,
94            pool,
95            arbiters: Mutex::new(arbiters),
96            storage: RwLock::new(HashMap::default()),
97            signals: AtomicBool::new(false),
98        }));
99        System::set_current(sys.clone());
100
101        let (stop_tx, stop) = oneshot::channel();
102
103        // system support tasks
104        crate::spawn(SystemSupport::new(&sys, stop_tx).run());
105        crate::spawn(controller.run(sys.clone()));
106
107        (sys, stop)
108    }
109
110    /// Creates a builder for a system with a customized runtime.
111    ///
112    /// See [`Builder`] for the available configuration options.
113    pub fn build() -> Builder {
114        Builder::new()
115    }
116
117    #[allow(clippy::new_ret_no_self)]
118    /// Creates a system runner with the specified name and runtime runner.
119    ///
120    /// # Panics
121    ///
122    /// Panics if the runtime cannot be created.
123    pub fn new<R: Runner>(name: &str, runner: R) -> SystemRunner {
124        Self::build().name(name).build(runner)
125    }
126
127    #[allow(clippy::new_ret_no_self)]
128    /// Creates a system runner from an existing configuration.
129    ///
130    /// # Panics
131    ///
132    /// Panics if the runtime cannot be created.
133    pub fn with_config(name: &str, mut config: SystemConfig) -> SystemRunner {
134        name.clone_into(&mut config.name);
135        Self::build().name(name).build_with(config)
136    }
137
138    /// Returns the system running on the current thread.
139    ///
140    /// # Panics
141    ///
142    /// Panics if no system is running on the current thread.
143    pub fn current() -> System {
144        CURRENT.with(|cell| match *cell.borrow() {
145            Some(ref sys) => sys.clone(),
146            None => panic!("System is not running"),
147        })
148    }
149
150    /// Returns the system running on the current thread, if one exists.
151    pub fn try_current() -> Option<System> {
152        CURRENT.with(|cell| cell.borrow().as_ref().map(Clone::clone))
153    }
154
155    /// Set current running system
156    #[doc(hidden)]
157    pub fn set_current(sys: System) {
158        CURRENT.with(|s| {
159            *s.borrow_mut() = Some(sys);
160        });
161    }
162
163    pub(crate) fn register_arbiter(&self, arb: Arbiter) {
164        CURRENT.with(|s| {
165            *s.borrow_mut() = Some(self.clone());
166        });
167        let mut arbiters = self.0.arbiters.lock();
168        arbiters.all.insert(arb.id(), arb.clone());
169        arbiters.list.push(arb);
170    }
171
172    pub(crate) fn unregister_arbiter(&self, id: Id) {
173        CURRENT.with(|s| {
174            *s.borrow_mut() = None;
175        });
176        let mut arbiters = self.0.arbiters.lock();
177        if let Some(hnd) = arbiters.all.remove(&id) {
178            for (idx, arb) in arbiters.list.iter().enumerate() {
179                if &hnd == arb {
180                    arbiters.list.remove(idx);
181                    break;
182                }
183            }
184        }
185    }
186
187    pub(super) fn remove_current() {
188        CURRENT.with(|cell| {
189            cell.borrow_mut().take();
190        });
191    }
192
193    /// Returns the system identifier.
194    pub fn id(&self) -> Id {
195        Id(self.0.id)
196    }
197
198    /// Returns the system name.
199    pub fn name(&self) -> &str {
200        &self.0.config.name
201    }
202
203    /// Stops the system with exit code `0`.
204    pub fn stop(&self) {
205        self.stop_with_code(0);
206    }
207
208    /// Stops the system with the specified exit code.
209    pub fn stop_with_code(&self, code: i32) {
210        let _ = self.0.sender.try_send(SystemCommand::Exit(code));
211    }
212
213    /// Returns whether process signal handling is enabled.
214    pub fn signals(&self) -> bool {
215        self.0.signals.load(Ordering::Relaxed)
216    }
217
218    /// Enables process signal handling.
219    ///
220    /// Signals are handled by one system at a time, this has no effect while
221    /// another system handles signals.
222    pub fn enable_signals(&self) {
223        if !self.signals() && crate::signals::start(self) {
224            self.0.signals.store(true, Ordering::Relaxed);
225        }
226    }
227
228    /// Disables process signal handling.
229    pub fn disable_signals(&self) {
230        if self.signals() {
231            crate::signals::stop(self);
232            self.0.signals.store(false, Ordering::Relaxed);
233        }
234    }
235
236    /// Returns the system's primary arbiter.
237    ///
238    /// # Panics
239    ///
240    /// Panics if the system has not been started.
241    pub fn arbiter(&self) -> Arbiter {
242        self.0.arbiter.clone()
243    }
244
245    /// Provides access to all arbiters registered with this system.
246    ///
247    /// This method should be called from the thread where the system has been initialized,
248    /// typically the "main" thread.
249    pub fn list_arbiters<F, R>(&self, f: F) -> R
250    where
251        F: FnOnce(&[Arbiter]) -> R,
252    {
253        f(&self.0.arbiters.lock().list)
254    }
255
256    /// Visits the latest ping records for each registered arbiter.
257    ///
258    /// This method should be called from the thread where the system has been initialized,
259    /// typically the "main" thread.
260    pub fn list_arbiter_pings<F>(mut f: F)
261    where
262        F: FnMut(&Arbiter, &mut VecDeque<PingRecord>),
263    {
264        PINGS.with(|pings| {
265            let mut p = pings.borrow_mut();
266            let sys = System::current();
267            let arbiters = sys.0.arbiters.lock();
268
269            for (id, recs) in &mut *p {
270                if let Some(arb) = arbiters.all.get(id) {
271                    f(arb, recs);
272                }
273            }
274        });
275    }
276
277    #[cfg(target_os = "linux")]
278    #[doc(hidden)]
279    /// Set arbiter latency callback.
280    ///
281    /// This callback is called when the arbiter response latency exceeds the
282    /// configured threshold. The provided backtrace is not resolved.
283    pub fn set_latency_callback<F>(f: F)
284    where
285        F: Fn(ntex_error::Backtrace) + Send + Sync + 'static,
286    {
287        *ARB_CB.lock() = Some(Arc::new(f));
288    }
289
290    /// Returns a clone of the system configuration.
291    pub fn config(&self) -> SystemConfig {
292        self.0.config.clone()
293    }
294
295    #[inline]
296    /// Returns a runtime handle for the primary arbiter.
297    pub fn handle(&self) -> Handle {
298        self.arbiter().handle().clone()
299    }
300
301    /// Returns whether the system is configured for testing.
302    pub fn testing(&self) -> bool {
303        self.0.config.testing()
304    }
305
306    /// Spawns a blocking task on a new thread and waits for it to complete.
307    ///
308    /// Dropping the returned future prevents queued work from starting, but
309    /// cannot interrupt work that is already running. Call
310    /// [`BlockingResult::detach`] to let queued work continue even if its
311    /// result is no longer needed.
312    pub fn spawn_blocking<F, R>(&self, f: F) -> BlockingResult<R>
313    where
314        F: FnOnce() -> R + Send + 'static,
315        R: Send + 'static,
316    {
317        self.0.pool.execute(f)
318    }
319
320    /// Returns a previously registered type, or inserts and returns a new one.
321    ///
322    /// This method acquires a lock on the internal data structure.
323    /// To avoid repeated locking, prefer storing a cloned value in the arbiter's storage.
324    pub fn get_value<T>(&self, f: impl FnOnce() -> T) -> T
325    where
326        T: Clone + Send + Sync + 'static,
327    {
328        if let Some(boxed) = self.0.storage.read().get(&TypeId::of::<T>())
329            && let Some(val) = (&**boxed as &(dyn Any + 'static)).downcast_ref::<T>()
330        {
331            val.clone()
332        } else {
333            let val = f();
334            self.0
335                .storage
336                .write()
337                .insert(TypeId::of::<T>(), Box::new(val.clone()));
338            val
339        }
340    }
341}
342
343impl SystemConfig {
344    #[inline]
345    /// Returns whether the system is configured for testing.
346    pub fn testing(&self) -> bool {
347        self.testing
348    }
349}
350
351impl fmt::Debug for System {
352    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
353        f.debug_struct("System")
354            .field("id", &self.0.id)
355            .field("config", &self.0.config)
356            .field("signals", &self.signals())
357            .field("pool", &self.0.pool)
358            .finish()
359    }
360}
361
362impl fmt::Debug for SystemConfig {
363    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
364        f.debug_struct("SystemConfig")
365            .field("name", &self.name)
366            .field("testing", &self.testing)
367            .field("stack_size", &self.stack_size)
368            .finish()
369    }
370}
371
372#[derive(Debug)]
373pub(super) enum SystemCommand {
374    Exit(i32),
375}
376
377#[derive(Debug)]
378struct SystemSupport {
379    sys: System,
380    stop: Option<oneshot::Sender<i32>>,
381    commands: Receiver<SystemCommand>,
382}
383
384impl SystemSupport {
385    fn new(sys: &System, stop: oneshot::Sender<i32>) -> Self {
386        Self {
387            sys: sys.clone(),
388            stop: Some(stop),
389            commands: sys.0.receiver.clone(),
390        }
391    }
392
393    async fn run(mut self) {
394        if self.sys.0.config.ping_interval != 0 {
395            crate::spawn(ping_arbiters(self.sys.clone()));
396        }
397
398        loop {
399            match self.commands.recv().await {
400                Ok(SystemCommand::Exit(code)) => {
401                    log::debug!("Stopping system with {code} code");
402
403                    // stop arbiters
404                    let mut arbiters = self.sys.0.arbiters.lock();
405                    for arb in arbiters.list.drain(..) {
406                        arb.stop();
407                    }
408                    arbiters.all.clear();
409                    drop(arbiters);
410
411                    crate::arbiter::run_shutdown_callbacks();
412
413                    // stop event loop
414                    if let Some(stop) = self.stop.take() {
415                        let _ = stop.send(code);
416                    }
417                }
418                Err(_) => {
419                    log::debug!("System stopped");
420                    return;
421                }
422            }
423        }
424    }
425}
426
427#[derive(Copy, Clone, Debug)]
428pub struct PingRecord {
429    /// Ping start time
430    pub start: Instant,
431    /// Round-trip time, if value is not set then ping is in process
432    pub rtt: Option<Duration>,
433}
434
435async fn ping_arbiters(sys: System) {
436    let arbs = Rc::new(RefCell::new(HashSet::default()));
437    let interval = Duration::from_millis(sys.0.config.ping_interval as u64);
438    #[cfg(target_os = "linux")]
439    let threshold = Duration::from_millis(sys.0.config.ping_threshold as u64);
440
441    loop {
442        // interval between pings
443        Delay::new(interval).await;
444
445        // send pings
446        {
447            arbs.borrow_mut().clear();
448
449            let start = Instant::now();
450            let arbiters = sys.0.arbiters.lock();
451
452            // drop records of stopped arbiters
453            PINGS.with(|pings| {
454                pings
455                    .borrow_mut()
456                    .retain(|id, _| arbiters.all.contains_key(id));
457            });
458
459            for arb in &arbiters.list {
460                let id = arb.id();
461                let arbs = arbs.clone();
462                let fut = arb.handle().spawn(async move {
463                    yield_to().await;
464                });
465
466                // calc ttl
467                PINGS.with(|pings| {
468                    let mut p = pings.borrow_mut();
469                    let recs = p.entry(arb.id()).or_default();
470                    recs.push_front(PingRecord { start, rtt: None });
471                    recs.truncate(10);
472                });
473
474                crate::spawn(async move {
475                    if fut.await.is_ok() {
476                        arbs.borrow_mut().insert(id);
477
478                        PINGS.with(|pings| {
479                            // a late pong belongs to its own round
480                            if let Some(recs) = pings.borrow_mut().get_mut(&id)
481                                && let Some(rec) = recs.iter_mut().find(|r| r.start == start)
482                            {
483                                rec.rtt = Some(start.elapsed());
484                            }
485                        });
486                    }
487                });
488            }
489        }
490
491        // check pings
492        #[cfg(target_os = "linux")]
493        {
494            const SPIN: Duration = Duration::from_micros(100);
495
496            // threshold
497            Delay::new(threshold).await;
498
499            let mut no_pongs = Vec::new();
500            {
501                for arb in &sys.0.arbiters.lock().list {
502                    let pong = arbs.borrow_mut().remove(&arb.id());
503                    if !pong {
504                        no_pongs.push(arb.clone());
505                    }
506                }
507            }
508
509            if !crate::signals::is_enabled() {
510                continue;
511            }
512
513            for arb in no_pongs {
514                // no response from arbiter
515                log::error!("Arbiter {}({:?}) did not return pong", arb.name(), arb.id());
516
517                // send tgkill to thread id to capture backtrace
518                let tid = arb.tid();
519                if !CAPTURE.arm(tid) {
520                    continue;
521                }
522                let result = unsafe {
523                    libc::syscall(libc::SYS_tgkill, libc::getpid(), arb.tid(), libc::SIGUSR2)
524                };
525
526                if result == -1 {
527                    log::error!(
528                        "Unsable to send SIGUSR2 to arbiter {}({:?}): {}",
529                        arb.name(),
530                        arb.id(),
531                        std::io::Error::last_os_error()
532                    );
533                } else {
534                    // Spin
535                    for _ in 0..1000 {
536                        Delay::new(SPIN).await;
537                        if let Some(bt) = CAPTURE.take(tid) {
538                            let bt = ntex_error::Backtrace::from(bt);
539                            if let Some(f) = latency_callback() {
540                                f(bt);
541                            } else {
542                                bt.resolver().resolve();
543                                log::error!(
544                                    "Worker does not returned pong within {interval:?} time.\n{bt:?}"
545                                );
546                            }
547                            break;
548                        }
549                    }
550                }
551                CAPTURE.disarm(tid);
552            }
553        }
554    }
555}
556
557async fn yield_to() {
558    use std::task::{Context, Poll};
559
560    struct Yield {
561        completed: bool,
562    }
563
564    impl Future for Yield {
565        type Output = ();
566
567        fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<()> {
568            if self.completed {
569                return Poll::Ready(());
570            }
571            self.completed = true;
572            cx.waker().wake_by_ref();
573            Poll::Pending
574        }
575    }
576
577    Yield { completed: false }.await;
578}
579
580#[cfg(target_os = "linux")]
581#[allow(clippy::type_complexity)]
582static ARB_CB: Mutex<Option<Arc<dyn Fn(ntex_error::Backtrace) + Send + Sync>>> = Mutex::new(None);
583
584#[cfg(target_os = "linux")]
585#[allow(clippy::type_complexity)]
586/// Clone the callback so it runs without the lock and survives replacement.
587fn latency_callback() -> Option<Arc<dyn Fn(ntex_error::Backtrace) + Send + Sync>> {
588    ARB_CB.lock().clone()
589}
590
591#[track_caller]
592#[cfg(target_family = "unix")]
593pub(crate) fn sig_usr2() {
594    #[cfg(target_os = "linux")]
595    #[allow(clippy::cast_possible_truncation)]
596    {
597        let tid = unsafe { libc::syscall(libc::SYS_gettid) } as i32;
598        CAPTURE.capture(tid, panic::Location::caller().file());
599    }
600}
601
602#[cfg(test)]
603mod ping_tests {
604    use super::*;
605    use crate::testing::TestRunner;
606
607    #[test]
608    fn ping_records_pruned() {
609        System::build()
610            .name("test")
611            .ping_interval(5)
612            .build(TestRunner)
613            .block_on(async {
614                let mut arb = Arbiter::new();
615                let id = arb.id();
616                let answered = || {
617                    PINGS.with(|p| {
618                        p.borrow()
619                            .get(&id)
620                            .is_some_and(|recs| recs.iter().any(|r| r.rtt.is_some()))
621                    })
622                };
623                let has = || PINGS.with(|p| p.borrow().contains_key(&id));
624
625                for _ in 0..400 {
626                    if answered() {
627                        break;
628                    }
629                    Delay::new(Duration::from_millis(5)).await;
630                }
631                assert!(answered());
632
633                arb.stop();
634                arb.join().unwrap();
635                for _ in 0..400 {
636                    if !has() {
637                        break;
638                    }
639                    Delay::new(Duration::from_millis(5)).await;
640                }
641                assert!(!has());
642            });
643    }
644}
645
646#[cfg(test)]
647mod api_tests {
648    use std::sync::atomic::{AtomicUsize, Ordering};
649
650    use super::*;
651    use crate::testing::TestRunner;
652
653    #[test]
654    #[should_panic(expected = "System is not running")]
655    fn current_without_system() {
656        assert!(System::try_current().is_none());
657        let _ = System::current();
658    }
659
660    #[test]
661    fn system_api() {
662        let config = System::new("base", TestRunner).block_on(async {
663            let sys = System::current();
664            assert_eq!(sys.name(), "base");
665            assert!(!sys.testing());
666            assert!(!sys.signals());
667            assert_eq!(System::try_current().unwrap().id(), sys.id());
668            assert_eq!(sys.arbiter(), Arbiter::current());
669            assert!(format!("{sys:?}").contains("base"));
670            assert!(format!("{:?}", sys.config()).contains("base"));
671
672            // values are stored once per system
673            assert_eq!(sys.get_value(|| 1usize), 1);
674            assert_eq!(sys.get_value(|| 2usize), 1);
675
676            assert_eq!(sys.spawn_blocking(|| 10).await, Ok(10));
677
678            let cnt = Arc::new(AtomicUsize::new(0));
679            let cnt2 = cnt.clone();
680            sys.handle()
681                .spawn(async move {
682                    cnt2.fetch_add(1, Ordering::Relaxed);
683                })
684                .await
685                .unwrap();
686            assert_eq!(cnt.load(Ordering::Relaxed), 1);
687            sys.list_arbiters(|arbs| assert_eq!(arbs, [sys.arbiter()]));
688            sys.config()
689        });
690        assert!(System::try_current().is_none());
691
692        System::with_config("copy", config).block_on(async {
693            let sys = System::current();
694            assert_eq!(sys.name(), "copy");
695            assert_eq!(Arbiter::current().name(), "copy");
696        });
697    }
698
699    #[test]
700    fn list_arbiter_pings() {
701        System::build()
702            .name("test")
703            .ping_interval(5)
704            .build(TestRunner)
705            .block_on(async {
706                let mut arb = Arbiter::new();
707                let visited = || {
708                    let mut found = false;
709                    System::list_arbiter_pings(|a, recs| {
710                        if a == &arb && recs.iter().any(|r| r.rtt.is_some()) {
711                            found = true;
712                        }
713                    });
714                    found
715                };
716                for _ in 0..400 {
717                    if visited() {
718                        break;
719                    }
720                    Delay::new(Duration::from_millis(5)).await;
721                }
722                assert!(visited());
723                arb.stop();
724                arb.join().unwrap();
725            });
726    }
727}
728
729#[cfg(all(test, target_os = "linux"))]
730mod tests {
731    use std::sync::atomic::{AtomicUsize, Ordering};
732    use std::{panic::Location, sync::Arc};
733
734    use super::*;
735
736    #[test]
737    fn latency_callback_replaced_while_running() {
738        let calls = Arc::new(AtomicUsize::new(0));
739        let data = Arc::new(vec![1u8; 64]);
740        let calls2 = calls.clone();
741        System::set_latency_callback(move |_| {
742            // replace itself, then use captured state
743            System::set_latency_callback(|_| {});
744            assert_eq!(data.len(), 64);
745            calls2.fetch_add(1, Ordering::Relaxed);
746        });
747
748        latency_callback().unwrap()(ntex_error::Backtrace::new(Location::caller()));
749        latency_callback().unwrap()(ntex_error::Backtrace::new(Location::caller()));
750        assert_eq!(calls.load(Ordering::Relaxed), 1);
751        *ARB_CB.lock() = None;
752    }
753}