Skip to main content

trusty_memory/
worker_liveness.rs

1//! Live worker-occupancy gauge for the palace open/write path (issue #4001).
2//!
3//! Why: during the #3992 incident six daemon threads were parked in
4//! `concurrent_open::backoff_sleep_ms` with a `memory_remember` hung ~1800 s,
5//! and BOTH `tm doctor` and `trusty-memory doctor` reported HEALTHY the whole
6//! time. Doctor observed only *process* liveness — an HTTP listener that still
7//! answers, and lock files that still look clean — neither of which can see a
8//! wedged worker pool. This module supplies the missing observation: how long
9//! the oldest in-flight palace operation has been running. That is the cheapest
10//! signal that actually distinguishes "the process is up" from "work is moving".
11//!
12//! What: a fixed-size, lock-free slot table of operation start timestamps. An
13//! operation claims a slot on entry and releases it on drop; the probe reads
14//! the table and reports the age of the oldest occupied slot. No allocation, no
15//! mutex, and no syscall on the hot path — one CAS in and one store out — so
16//! the gauge can never itself become the load problem it exists to detect.
17//! Test: see `worker_liveness_tests.rs`.
18
19use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
20use std::time::{Duration, Instant};
21
22/// Number of concurrently-trackable operations.
23///
24/// Why: the table is scanned linearly on both claim and probe, so it must stay
25/// small enough that a scan is trivial (64 relaxed atomic loads is ~nothing)
26/// while comfortably exceeding realistic in-flight concurrency for a
27/// single-user memory daemon. Operations beyond this count are still *counted*
28/// via the overflow gauge — they simply do not contribute an age sample, which
29/// degrades the signal gracefully instead of blocking or allocating.
30/// Test: `overflow_is_counted_when_slots_exhausted`.
31const SLOTS: usize = 64;
32
33/// Sentinel meaning "this slot is free". Real timestamps are offsets from the
34/// tracker's epoch and are stored as `millis + 1`, so 0 is never a valid entry.
35const FREE: u64 = 0;
36
37/// How long the oldest in-flight operation may run before the pool is called
38/// wedged.
39///
40/// Why: this must sit ABOVE every legitimately-bounded wait, or healthy load
41/// would read as a wedge. The longest such bound on the palace open path is
42/// `memory_core::timeouts::open_queue_timeout()` (default 60 s, issue #3992),
43/// after which `open_palace` gives up and returns an error. Doubling it means
44/// any operation still outstanding has already blown through the bound that
45/// was supposed to release it — which is precisely the #3992 signature, where
46/// a `memory_remember` ran ~1800 s. Env-overridable so an operator running a
47/// deliberately long `open_queue_timeout` can move the wedge line with it.
48/// What: `TRUSTY_WEDGE_THRESHOLD_SECS` if set and parseable, else
49/// `2 × open_queue_timeout()`.
50/// Test: `wedge_threshold_exceeds_the_open_queue_bound`.
51pub fn wedge_threshold() -> Duration {
52    if let Some(secs) = std::env::var("TRUSTY_WEDGE_THRESHOLD_SECS")
53        .ok()
54        .and_then(|v| v.parse::<u64>().ok())
55    {
56        return Duration::from_secs(secs);
57    }
58    trusty_common::memory_core::timeouts::open_queue_timeout() * 2
59}
60
61/// Tracks how long the oldest in-flight palace operation has been running.
62///
63/// Why (issue #4001): see the module docs — this is the signal that would have
64/// revealed the #3992 wedge. Alternatives considered and rejected: sampling
65/// thread state (needs platform-specific debugging APIs and is expensive),
66/// mutex wait-time histograms (needs instrumenting `parking_lot` internals),
67/// and a plain in-flight *count* (a count alone cannot distinguish healthy
68/// concurrency from a wedge — only the *age* of outstanding work can).
69/// What: a slot table of start timestamps plus an overflow counter. Cloneable
70/// and `Send`/`Sync` via the caller's `Arc`.
71/// Test: `worker_liveness_tests.rs`.
72#[derive(Debug)]
73pub struct WorkerLiveness {
74    /// Start timestamps as `epoch.elapsed().as_millis() + 1`; [`FREE`] when
75    /// the slot is unoccupied.
76    slots: [AtomicU64; SLOTS],
77    /// Operations in flight that could not claim a slot. Contributes to the
78    /// in-flight count but not to the oldest-age sample.
79    overflow: AtomicUsize,
80    /// Reference point for every stored timestamp. Using a monotonic `Instant`
81    /// rather than wall-clock time keeps the gauge immune to clock steps.
82    epoch: Instant,
83}
84
85impl Default for WorkerLiveness {
86    fn default() -> Self {
87        Self::new()
88    }
89}
90
91impl WorkerLiveness {
92    /// Create an empty tracker.
93    ///
94    /// Why: `AtomicU64` is not `Copy`, so the slot array cannot be built with
95    /// `[AtomicU64::new(0); SLOTS]`; `from_fn` is the idiomatic construction.
96    /// What: all slots [`FREE`], overflow zero, epoch = now.
97    /// Test: `idle_tracker_reports_no_work`.
98    pub fn new() -> Self {
99        Self {
100            slots: std::array::from_fn(|_| AtomicU64::new(FREE)),
101            overflow: AtomicUsize::new(0),
102            epoch: Instant::now(),
103        }
104    }
105
106    /// Register the start of an operation, returning a guard that unregisters
107    /// it on drop.
108    ///
109    /// Why: the release MUST be drop-driven rather than an explicit call. The
110    /// operations being tracked are exactly the ones that fail, time out, and
111    /// `?`-propagate mid-function; an explicit `finish()` would be skipped on
112    /// every error path and leak slots, which would then be misread as a
113    /// permanent wedge — turning this fix into a new false positive.
114    /// What: claims the first [`FREE`] slot via a relaxed CAS, storing the
115    /// current offset from `epoch`. Falls back to the overflow counter when
116    /// every slot is taken.
117    /// Test: `guard_releases_slot_on_drop`, `guard_releases_slot_on_panic`.
118    pub fn track(&self) -> WorkGuard<'_> {
119        // `+ 1` keeps 0 reserved as the FREE sentinel, so an operation
120        // starting in the tracker's first millisecond is still distinguishable
121        // from an empty slot.
122        let now = self.epoch.elapsed().as_millis() as u64 + 1;
123        for (idx, slot) in self.slots.iter().enumerate() {
124            if slot
125                .compare_exchange(FREE, now, Ordering::AcqRel, Ordering::Relaxed)
126                .is_ok()
127            {
128                return WorkGuard {
129                    tracker: self,
130                    slot: Some(idx),
131                };
132            }
133        }
134        self.overflow.fetch_add(1, Ordering::AcqRel);
135        WorkGuard {
136            tracker: self,
137            slot: None,
138        }
139    }
140
141    /// Number of operations currently in flight.
142    ///
143    /// What: occupied slots plus the overflow count.
144    /// Test: `in_flight_counts_active_operations`.
145    pub fn in_flight(&self) -> usize {
146        let occupied = self
147            .slots
148            .iter()
149            .filter(|s| s.load(Ordering::Acquire) != FREE)
150            .count();
151        occupied + self.overflow.load(Ordering::Acquire)
152    }
153
154    /// Age of the oldest in-flight operation, or `None` when idle.
155    ///
156    /// Why: this is the actual health signal. A daemon with zero in-flight work
157    /// is trivially not wedged; a daemon whose oldest operation has been
158    /// running for minutes is wedged regardless of what the HTTP listener says.
159    /// What: scans for the smallest occupied timestamp and returns the elapsed
160    /// duration since it. Returns `None` when every slot is free — note that
161    /// overflow-only operations produce `None` here by design, since no age
162    /// sample exists for them.
163    /// Test: `oldest_age_tracks_the_earliest_operation`.
164    pub fn oldest_age(&self) -> Option<Duration> {
165        let oldest = self
166            .slots
167            .iter()
168            .map(|s| s.load(Ordering::Acquire))
169            .filter(|&v| v != FREE)
170            .min()?;
171        let now = self.epoch.elapsed().as_millis() as u64 + 1;
172        Some(Duration::from_millis(now.saturating_sub(oldest)))
173    }
174
175    /// True when the oldest in-flight operation has exceeded `threshold`.
176    ///
177    /// Why: turns the raw age into the verdict doctor reports. Keeping the
178    /// threshold a parameter (rather than a constant baked in here) lets the
179    /// caller pick a bound derived from the real operation timeout, and lets
180    /// tests drive the wedge condition deterministically without sleeping.
181    /// What: `oldest_age() > threshold`; `false` when idle.
182    /// Test: `wedged_when_oldest_exceeds_threshold`.
183    pub fn is_wedged(&self, threshold: Duration) -> bool {
184        self.oldest_age().is_some_and(|age| age > threshold)
185    }
186}
187
188/// RAII registration for one in-flight operation.
189///
190/// Why: see [`WorkerLiveness::track`] — drop-driven release is what makes the
191/// gauge correct across the `?` and panic paths that a wedge actually travels.
192/// What: releases its slot (or decrements overflow) on drop.
193/// Test: `guard_releases_slot_on_drop`, `guard_releases_slot_on_panic`.
194#[derive(Debug)]
195pub struct WorkGuard<'a> {
196    tracker: &'a WorkerLiveness,
197    /// `Some(idx)` when a slot was claimed, `None` when the operation landed
198    /// in the overflow bucket.
199    slot: Option<usize>,
200}
201
202impl Drop for WorkGuard<'_> {
203    fn drop(&mut self) {
204        match self.slot {
205            Some(idx) => {
206                self.tracker.slots[idx].store(FREE, Ordering::Release);
207            }
208            None => {
209                self.tracker.overflow.fetch_sub(1, Ordering::AcqRel);
210            }
211        }
212    }
213}
214
215#[cfg(test)]
216#[path = "worker_liveness_tests.rs"]
217mod tests;