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. Two such bounds are tracked: the palace open queue,
42/// `memory_core::timeouts::open_queue_timeout()` (default 60 s, issue #3992),
43/// and a writer queued on the palace write lock, `write_lock_timeout()`
44/// (default 60 s), tracked since #4001. Doubling the larger means any operation
45/// still outstanding has already blown through the bound that was supposed to
46/// release it — the #3992 signature, where a `memory_remember` ran ~1800 s.
47/// Env-overridable so an operator running a deliberately long bound can move
48/// the wedge line with it.
49/// What: `TRUSTY_WEDGE_THRESHOLD_SECS` if set and parseable, else
50/// `2 × max(open_queue_timeout(), write_lock_timeout())`.
51/// Test: `wedge_threshold_exceeds_the_open_queue_bound`,
52/// `wedge_threshold_doubles_the_larger_wait_bound`.
53pub fn wedge_threshold() -> Duration {
54 let override_secs = std::env::var("TRUSTY_WEDGE_THRESHOLD_SECS")
55 .ok()
56 .and_then(|v| v.parse::<u64>().ok());
57 derive_wedge_threshold(
58 override_secs,
59 trusty_common::memory_core::timeouts::open_queue_timeout(),
60 trusty_common::memory_core::timeouts::write_lock_timeout(),
61 )
62}
63
64/// [`wedge_threshold`] with its inputs supplied (#4001).
65///
66/// Why: both bounds are process-wide env reads, so a test cannot vary them
67/// without racing its siblings. A raised `TRUSTY_WRITE_LOCK_TIMEOUT_SECS` must
68/// move the line, or a writer legitimately queued past the old line would read
69/// as wedged.
70/// What: the override in seconds when present, else twice the larger bound.
71/// Test: `wedge_threshold_doubles_the_larger_wait_bound`.
72pub(crate) fn derive_wedge_threshold(
73 override_secs: Option<u64>,
74 open_queue: Duration,
75 write_lock: Duration,
76) -> Duration {
77 match override_secs {
78 Some(secs) => Duration::from_secs(secs),
79 None => open_queue.max(write_lock) * 2,
80 }
81}
82
83/// Tracks how long the oldest in-flight palace operation has been running.
84///
85/// Why (issue #4001): see the module docs — this is the signal that would have
86/// revealed the #3992 wedge. Alternatives considered and rejected: sampling
87/// thread state (needs platform-specific debugging APIs and is expensive),
88/// mutex wait-time histograms (needs instrumenting `parking_lot` internals),
89/// and a plain in-flight *count* (a count alone cannot distinguish healthy
90/// concurrency from a wedge — only the *age* of outstanding work can).
91/// What: a slot table of start timestamps plus an overflow counter. Cloneable
92/// and `Send`/`Sync` via the caller's `Arc`.
93/// Test: `worker_liveness_tests.rs`.
94#[derive(Debug)]
95pub struct WorkerLiveness {
96 /// Start timestamps as `epoch.elapsed().as_millis() + 1`; [`FREE`] when
97 /// the slot is unoccupied.
98 slots: [AtomicU64; SLOTS],
99 /// Operations in flight that could not claim a slot. Contributes to the
100 /// in-flight count but not to the oldest-age sample.
101 overflow: AtomicUsize,
102 /// Reference point for every stored timestamp. Using a monotonic `Instant`
103 /// rather than wall-clock time keeps the gauge immune to clock steps.
104 epoch: Instant,
105}
106
107impl Default for WorkerLiveness {
108 fn default() -> Self {
109 Self::new()
110 }
111}
112
113impl WorkerLiveness {
114 /// Create an empty tracker.
115 ///
116 /// Why: `AtomicU64` is not `Copy`, so the slot array cannot be built with
117 /// `[AtomicU64::new(0); SLOTS]`; `from_fn` is the idiomatic construction.
118 /// What: all slots [`FREE`], overflow zero, epoch = now.
119 /// Test: `idle_tracker_reports_no_work`.
120 pub fn new() -> Self {
121 Self {
122 slots: std::array::from_fn(|_| AtomicU64::new(FREE)),
123 overflow: AtomicUsize::new(0),
124 epoch: Instant::now(),
125 }
126 }
127
128 /// Register the start of an operation, returning a guard that unregisters
129 /// it on drop.
130 ///
131 /// Why: the release MUST be drop-driven rather than an explicit call. The
132 /// operations being tracked are exactly the ones that fail, time out, and
133 /// `?`-propagate mid-function; an explicit `finish()` would be skipped on
134 /// every error path and leak slots, which would then be misread as a
135 /// permanent wedge — turning this fix into a new false positive.
136 /// What: claims the first [`FREE`] slot via a relaxed CAS, storing the
137 /// current offset from `epoch`. Falls back to the overflow counter when
138 /// every slot is taken.
139 /// Test: `guard_releases_slot_on_drop`, `guard_releases_slot_on_panic`.
140 pub fn track(&self) -> WorkGuard<'_> {
141 // `+ 1` keeps 0 reserved as the FREE sentinel, so an operation
142 // starting in the tracker's first millisecond is still distinguishable
143 // from an empty slot.
144 let now = self.epoch.elapsed().as_millis() as u64 + 1;
145 for (idx, slot) in self.slots.iter().enumerate() {
146 if slot
147 .compare_exchange(FREE, now, Ordering::AcqRel, Ordering::Relaxed)
148 .is_ok()
149 {
150 return WorkGuard {
151 tracker: self,
152 slot: Some(idx),
153 };
154 }
155 }
156 self.overflow.fetch_add(1, Ordering::AcqRel);
157 WorkGuard {
158 tracker: self,
159 slot: None,
160 }
161 }
162
163 /// Number of operations currently in flight.
164 ///
165 /// What: occupied slots plus the overflow count.
166 /// Test: `in_flight_counts_active_operations`.
167 pub fn in_flight(&self) -> usize {
168 let occupied = self
169 .slots
170 .iter()
171 .filter(|s| s.load(Ordering::Acquire) != FREE)
172 .count();
173 occupied + self.overflow.load(Ordering::Acquire)
174 }
175
176 /// Age of the oldest in-flight operation, or `None` when idle.
177 ///
178 /// Why: this is the actual health signal. A daemon with zero in-flight work
179 /// is trivially not wedged; a daemon whose oldest operation has been
180 /// running for minutes is wedged regardless of what the HTTP listener says.
181 /// What: scans for the smallest occupied timestamp and returns the elapsed
182 /// duration since it. Returns `None` when every slot is free — note that
183 /// overflow-only operations produce `None` here by design, since no age
184 /// sample exists for them.
185 /// Test: `oldest_age_tracks_the_earliest_operation`.
186 pub fn oldest_age(&self) -> Option<Duration> {
187 let oldest = self
188 .slots
189 .iter()
190 .map(|s| s.load(Ordering::Acquire))
191 .filter(|&v| v != FREE)
192 .min()?;
193 let now = self.epoch.elapsed().as_millis() as u64 + 1;
194 Some(Duration::from_millis(now.saturating_sub(oldest)))
195 }
196
197 /// True when the oldest in-flight operation has exceeded `threshold`.
198 ///
199 /// Why: turns the raw age into the verdict doctor reports. Keeping the
200 /// threshold a parameter (rather than a constant baked in here) lets the
201 /// caller pick a bound derived from the real operation timeout, and lets
202 /// tests drive the wedge condition deterministically without sleeping.
203 /// What: `oldest_age() > threshold`; `false` when idle.
204 /// Test: `wedged_when_oldest_exceeds_threshold`.
205 pub fn is_wedged(&self, threshold: Duration) -> bool {
206 self.oldest_age().is_some_and(|age| age > threshold)
207 }
208}
209
210/// RAII registration for one in-flight operation.
211///
212/// Why: see [`WorkerLiveness::track`] — drop-driven release is what makes the
213/// gauge correct across the `?` and panic paths that a wedge actually travels.
214/// What: releases its slot (or decrements overflow) on drop.
215/// Test: `guard_releases_slot_on_drop`, `guard_releases_slot_on_panic`.
216#[derive(Debug)]
217pub struct WorkGuard<'a> {
218 tracker: &'a WorkerLiveness,
219 /// `Some(idx)` when a slot was claimed, `None` when the operation landed
220 /// in the overflow bucket.
221 slot: Option<usize>,
222}
223
224impl Drop for WorkGuard<'_> {
225 fn drop(&mut self) {
226 match self.slot {
227 Some(idx) => {
228 self.tracker.slots[idx].store(FREE, Ordering::Release);
229 }
230 None => {
231 self.tracker.overflow.fetch_sub(1, Ordering::AcqRel);
232 }
233 }
234 }
235}
236
237#[cfg(test)]
238#[path = "worker_liveness_tests.rs"]
239mod tests;