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;