rightkit_process/
stall.rs1use std::io;
24use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
25use std::sync::Arc;
26use std::thread::{self, JoinHandle};
27use std::time::{Duration, Instant};
28
29#[derive(Clone, Debug)]
31pub struct Ack {
32 epoch: Instant,
33 done_ms: Arc<AtomicU64>,
34}
35
36impl Ack {
37 pub fn done(&self) {
39 self.done_ms.store(
40 self.epoch.elapsed().as_millis() as u64 + 1,
41 Ordering::Release,
42 );
43 }
44}
45
46#[derive(Clone, Copy, Debug)]
47pub struct StallConfig {
48 pub poll: Duration,
50 pub threshold: Duration,
52 pub observer_gap: Duration,
54}
55
56impl Default for StallConfig {
57 fn default() -> Self {
58 Self {
59 poll: Duration::from_millis(50),
60 threshold: Duration::from_millis(200),
61 observer_gap: Duration::from_secs(1),
62 }
63 }
64}
65
66#[derive(Clone, Debug, PartialEq, Eq)]
67pub enum StallEvent {
68 Stalled {
71 episode: u64,
72 gap: Duration,
73 recovered: bool,
74 },
75 Recovered { episode: u64, gap: Duration },
77 ObserverGap { gap: Duration },
79 Stopped { reason: String },
81}
82
83#[derive(Debug, Default)]
85struct Episode {
86 reported: bool,
87}
88
89impl Episode {
90 fn detect(&mut self, gap: Duration, threshold: Duration) -> bool {
91 if gap < threshold || self.reported {
92 return false;
93 }
94 self.reported = true;
95 true
96 }
97}
98
99#[derive(Debug, Clone, Copy)]
102pub struct Cooldown {
103 every: Duration,
104 last: Option<Duration>,
105}
106
107impl Cooldown {
108 pub fn new(every: Duration) -> Self {
109 Self { every, last: None }
110 }
111 pub fn reserve(&mut self, now: Duration) -> bool {
113 if self
114 .last
115 .is_some_and(|l| now.saturating_sub(l) < self.every)
116 {
117 return false;
118 }
119 self.last = Some(now);
120 true
121 }
122}
123
124pub struct StallWatchdog {
126 stop: Arc<AtomicBool>,
127 handle: Option<JoinHandle<()>>,
128}
129
130impl StallWatchdog {
131 pub fn spawn(
135 name: &str,
136 config: StallConfig,
137 ping: impl Fn(Ack) -> Result<(), String> + Send + 'static,
138 on_event: impl Fn(StallEvent) + Send + 'static,
139 ) -> io::Result<Self> {
140 let stop = Arc::new(AtomicBool::new(false));
141 let flag = Arc::clone(&stop);
142 let handle = thread::Builder::new()
143 .name(name.to_string())
144 .spawn(move || run(config, flag, ping, on_event))?;
145 Ok(Self {
146 stop,
147 handle: Some(handle),
148 })
149 }
150
151 pub fn stop(&mut self) {
153 self.stop.store(true, Ordering::SeqCst);
154 if let Some(h) = self.handle.take() {
155 let _ = h.join();
156 }
157 }
158}
159
160impl Drop for StallWatchdog {
161 fn drop(&mut self) {
162 self.stop();
163 }
164}
165
166fn run(
167 cfg: StallConfig,
168 stop: Arc<AtomicBool>,
169 ping: impl Fn(Ack) -> Result<(), String>,
170 on_event: impl Fn(StallEvent),
171) {
172 let epoch = Instant::now();
173 let ms = |at: Instant| at.saturating_duration_since(epoch).as_millis() as u64 + 1;
174 let mut done = Arc::new(AtomicU64::new(0));
175 let mut pending: Option<u64> = None;
176 let mut episode = Episode::default();
177 let mut episode_no = 0u64;
178 let mut previous_poll = ms(Instant::now());
179 while !stop.load(Ordering::SeqCst) {
180 thread::sleep(cfg.poll);
181 if stop.load(Ordering::SeqCst) {
182 return;
183 }
184 let now = ms(Instant::now());
185 let observer = Duration::from_millis(now.saturating_sub(previous_poll));
186 previous_poll = now;
187 if observer > cfg.observer_gap {
188 on_event(StallEvent::ObserverGap { gap: observer });
189 if pending.is_some() {
190 pending = Some(now);
191 }
192 episode.reported = false;
193 }
194 if let Some(sent) = pending {
195 let acked = done.load(Ordering::Acquire);
196 let gap = Duration::from_millis(if acked != 0 {
197 acked.saturating_sub(sent)
198 } else {
199 now.saturating_sub(sent)
200 });
201 if episode.detect(gap, cfg.threshold) {
202 episode_no += 1;
203 on_event(StallEvent::Stalled {
204 episode: episode_no,
205 gap,
206 recovered: acked != 0,
207 });
208 }
209 if acked == 0 {
210 continue;
211 }
212 if episode.reported {
213 on_event(StallEvent::Recovered {
214 episode: episode_no,
215 gap,
216 });
217 }
218 episode.reported = false;
219 pending = None;
220 }
221 if pending.is_none() {
222 done = Arc::new(AtomicU64::new(0));
223 let ack = Ack {
224 epoch,
225 done_ms: Arc::clone(&done),
226 };
227 match ping(ack) {
228 Ok(()) => pending = Some(ms(Instant::now())),
229 Err(reason) => {
230 on_event(StallEvent::Stopped { reason });
231 return;
232 }
233 }
234 }
235 }
236}