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#[derive(Copy, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Debug)]
31pub struct Id(pub(crate) usize);
32
33pub 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#[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 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 crate::spawn(SystemSupport::new(&sys, stop_tx).run());
105 crate::spawn(controller.run(sys.clone()));
106
107 (sys, stop)
108 }
109
110 pub fn build() -> Builder {
114 Builder::new()
115 }
116
117 #[allow(clippy::new_ret_no_self)]
118 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 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 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 pub fn try_current() -> Option<System> {
152 CURRENT.with(|cell| cell.borrow().as_ref().map(Clone::clone))
153 }
154
155 #[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 pub fn id(&self) -> Id {
195 Id(self.0.id)
196 }
197
198 pub fn name(&self) -> &str {
200 &self.0.config.name
201 }
202
203 pub fn stop(&self) {
205 self.stop_with_code(0);
206 }
207
208 pub fn stop_with_code(&self, code: i32) {
210 let _ = self.0.sender.try_send(SystemCommand::Exit(code));
211 }
212
213 pub fn signals(&self) -> bool {
215 self.0.signals.load(Ordering::Relaxed)
216 }
217
218 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 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 pub fn arbiter(&self) -> Arbiter {
242 self.0.arbiter.clone()
243 }
244
245 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 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 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 pub fn config(&self) -> SystemConfig {
292 self.0.config.clone()
293 }
294
295 #[inline]
296 pub fn handle(&self) -> Handle {
298 self.arbiter().handle().clone()
299 }
300
301 pub fn testing(&self) -> bool {
303 self.0.config.testing()
304 }
305
306 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 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 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 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 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 pub start: Instant,
431 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 Delay::new(interval).await;
444
445 {
447 arbs.borrow_mut().clear();
448
449 let start = Instant::now();
450 let arbiters = sys.0.arbiters.lock();
451
452 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 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 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 #[cfg(target_os = "linux")]
493 {
494 const SPIN: Duration = Duration::from_micros(100);
495
496 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 log::error!("Arbiter {}({:?}) did not return pong", arb.name(), arb.id());
516
517 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 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)]
586fn 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 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 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}