rings_core/lifecycle/
stop.rs1use std::sync::atomic::AtomicBool;
2use std::sync::atomic::Ordering;
3use std::sync::Arc;
4
5use event_listener::Event;
6
7#[derive(Default)]
8struct StopState {
9 requested: AtomicBool,
10 changed: Event,
11}
12
13#[derive(Clone, Default)]
15pub struct StopSource {
16 state: Arc<StopState>,
17}
18
19impl StopSource {
20 pub fn new() -> Self {
22 Self::default()
23 }
24
25 pub fn token(&self) -> StopToken {
27 StopToken {
28 state: Arc::clone(&self.state),
29 }
30 }
31
32 pub fn request_stop(&self) {
36 if !self.state.requested.swap(true, Ordering::AcqRel) {
37 self.state.changed.notify(usize::MAX);
38 }
39 }
40
41 pub fn is_stop_requested(&self) -> bool {
43 self.state.requested.load(Ordering::Acquire)
44 }
45}
46
47impl std::fmt::Debug for StopSource {
48 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
49 formatter
50 .debug_struct("StopSource")
51 .field("requested", &self.is_stop_requested())
52 .finish()
53 }
54}
55
56#[derive(Clone, Default)]
58pub struct StopToken {
59 state: Arc<StopState>,
60}
61
62impl StopToken {
63 pub fn never() -> Self {
65 Self::default()
66 }
67
68 pub fn should_stop(&self) -> bool {
70 self.state.requested.load(Ordering::Acquire)
71 }
72
73 pub async fn stopped(&self) {
78 loop {
79 if self.should_stop() {
80 return;
81 }
82 let changed = self.state.changed.listen();
83 if self.should_stop() {
84 return;
85 }
86 changed.await;
87 }
88 }
89}
90
91impl std::fmt::Debug for StopToken {
92 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
93 formatter
94 .debug_struct("StopToken")
95 .field("requested", &self.should_stop())
96 .finish()
97 }
98}