Skip to main content

rightkit_process/
stall.rs

1//! Stall watchdog: detects that something that should keep answering has
2//! stopped. This is the shared form of HeardRight's `hang_watchdog` (main
3//! thread, 5 s) and `ui_stall_watchdog` (UI thread, 200 ms, one outstanding
4//! ping, observer-gap guard, capture cooldown).
5//!
6//! Mechanism: a dedicated OS thread sends one *ping* at a time to the watched
7//! party (a closure the app supplies: for the UI thread it is
8//! `run_on_main_thread(move || ack.done())`; for a worker, a message on its
9//! queue). The watched party calls [`Ack::done`] when it gets round to it. The
10//! gap between sending and acknowledging is the stall. The watchdog thread
11//! takes no application lock, so it can observe a deadlock without joining it.
12//!
13//! Semantics, all inherited from the donor and tested end to end:
14//! - **Edge-triggered**: one [`StallEvent::Stalled`] per episode, then
15//!   [`StallEvent::Recovered`] when the ack lands; the next stall reports again.
16//! - **One outstanding ping**: a frozen target never accumulates queued pings.
17//! - **Observer gap**: if the watchdog thread itself was not scheduled for
18//!   longer than `observer_gap` (laptop sleep, SIGSTOP, heavy load) the interval
19//!   proves nothing about the target, so its timing is discarded
20//!   ([`StallEvent::ObserverGap`]) rather than reported as a stall.
21//! - [`Cooldown`] rate-limits expensive captures (stack samples) per process.
22
23use 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/// Handed to the watched party; calling [`Ack::done`] proves it is alive.
30#[derive(Clone, Debug)]
31pub struct Ack {
32    epoch: Instant,
33    done_ms: Arc<AtomicU64>,
34}
35
36impl Ack {
37    /// Mark the ping answered. Cheap, lock-free, callable from any thread.
38    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    /// How often the watchdog wakes (HeardRight UI: 50 ms; main-thread: 1 s).
49    pub poll: Duration,
50    /// Ack latency at or above which an episode is reported (UI: 200 ms; hang: 5 s).
51    pub threshold: Duration,
52    /// A watchdog wake later than this since the previous one is an observer gap.
53    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    /// The target has not acknowledged for `gap`. `recovered` is true when the
69    /// ack landed late (a stall seen only after the fact: nothing left to sample).
70    Stalled {
71        episode: u64,
72        gap: Duration,
73        recovered: bool,
74    },
75    /// The target is answering again after a reported stall.
76    Recovered { episode: u64, gap: Duration },
77    /// The watchdog thread was not scheduled for `gap`; that interval's timing was discarded.
78    ObserverGap { gap: Duration },
79    /// The ping closure returned an error; the watchdog stopped.
80    Stopped { reason: String },
81}
82
83/// Pure episode state: report once per stall, re-arm on recovery.
84#[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/// Minimum spacing between expensive actions (HeardRight samples stacks at
100/// most once a minute). Pure; the caller supplies `now`.
101#[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    /// True (and records `now`) when the action may run.
112    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
124/// A running watchdog. Dropping it stops and joins the thread.
125pub struct StallWatchdog {
126    stop: Arc<AtomicBool>,
127    handle: Option<JoinHandle<()>>,
128}
129
130impl StallWatchdog {
131    /// Start watching. `ping` must hand the [`Ack`] to the target without
132    /// blocking on it (enqueue and return); `on_event` runs on the watchdog
133    /// thread, so keep it short and lock-free (log, flip an atomic, spawn work).
134    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    /// Stop and join. Idempotent.
152    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}