trusty_memory/lock_stall.rs
1//! Holder-agnostic stall detector for each open palace's handle locks (#4001).
2//!
3//! Why: on 2026-09-13 a dream cycle held `PalaceHandle::write_mutex`, and
4//! `memory.health` stayed `ok`. [`crate::worker_liveness`] sees only operations
5//! that register with it. None of the inner lock's acquirers in `trusty-common`
6//! registers: remember, forget, orphan compaction, the dream cycle's compact,
7//! rebuild and KG passes, and share import. Writers queued behind the holder
8//! give up at their bound, which sits below the wedge threshold, so no tracked
9//! age ever crossed it.
10//! What: a sweep `try_lock`s `write_mutex` and `commit_mutex` on every open
11//! palace. A lock found held with no probe outstanding gets a stamp and a probe
12//! task. The probe queues on the lock like a writer but never gives up, and
13//! only its own acquisition clears the stamp. The stamp's age is therefore how
14//! long a writer arriving at the stamp has waited, whoever holds the lock.
15//! Each stamp also carries the identity of the mutex it was taken against, so
16//! a stamp left by a handle the registry has since replaced cannot be read as
17//! a stall of the live handle.
18//!
19//! A probe is a real writer, not a passive observer: it queues on the lock with
20//! `lock().await`, and tokio assigns the permit at the holder's release. The
21//! probe therefore owns the lock from that release until its task is next
22//! polled, and any writer that arrives inside that window queues behind it. The
23//! window is one scheduler hop, and it buys the guarantee the detector exists
24//! for — the stamp clears only when a queued writer would really have been let
25//! through. A `try_lock` poll loop would avoid the window at the cost of that
26//! guarantee, since `try_lock` can succeed on a lock a queued writer is still
27//! waiting for.
28//! Test: `lock_stall_tests.rs`,
29//! `tools::tests::write_liveness_tests::a_dream_cycle_holding_the_handle_write_mutex_reads_as_wedged`.
30
31use std::collections::HashMap;
32use std::sync::atomic::{AtomicBool, Ordering};
33use std::sync::{Arc, Mutex, MutexGuard, Weak};
34use std::time::{Duration, Instant};
35
36use trusty_common::memory_core::PalaceRegistry;
37
38/// A ticker heartbeat older than this many intervals reads as stopped.
39const TICKER_GRACE_INTERVALS: u32 = 3;
40
41/// Which handle-level lock a stamp describes.
42#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize)]
43#[serde(rename_all = "snake_case")]
44pub enum PalaceLock {
45 /// `PalaceHandle::write_mutex`: the per-palace write critical section.
46 Write,
47 /// `PalaceHandle::commit_mutex`: the durable-commit tail (#6366), which can
48 /// outlive a write that gave up.
49 Commit,
50}
51
52/// One held lock, as first sighted.
53#[derive(Debug, Clone)]
54struct Stall {
55 /// When a sweep first found the lock held with no probe outstanding.
56 since: Instant,
57 /// True while a probe task is queued on the lock.
58 probing: bool,
59 /// The mutex this stamp was taken against (#4001).
60 ///
61 /// Why: the key is `(palace id, lock)`, which a reopened palace reuses. A
62 /// probe still queued on a superseded handle's mutex would otherwise keep
63 /// the stamp alive against a healthy live handle, and would clear a live
64 /// handle's stamp when its own lock finally freed.
65 /// What: a `Weak` rather than a raw address — it keeps the allocation from
66 /// being reused, so pointer identity cannot alias a later mutex.
67 mutex: Weak<tokio::sync::Mutex<()>>,
68}
69
70/// Whether `stamped` was taken against `mutex`.
71fn stamped_against(
72 stamped: &Weak<tokio::sync::Mutex<()>>,
73 mutex: &Arc<tokio::sync::Mutex<()>>,
74) -> bool {
75 std::ptr::eq(stamped.as_ptr(), Arc::as_ptr(mutex))
76}
77
78/// The longest-standing stall, as reported to `memory.health`.
79#[derive(Debug, Clone, PartialEq, Eq)]
80pub struct StalledLock {
81 /// Palace id whose lock is held.
82 pub palace: String,
83 /// Which of the palace's locks.
84 pub lock: PalaceLock,
85 /// How long a writer arriving at the first sighting has waited.
86 pub age: Duration,
87}
88
89type StallKey = (String, PalaceLock);
90
91/// Stamps for handle locks found held, plus the ticker's heartbeat.
92///
93/// Why: see the module docs. A `std` mutex guards the table because every
94/// critical section is a map operation with no `.await`.
95/// What: the stamp table, a poisoned flag, the health-path rate limit, and the
96/// ticker heartbeat. A poisoned lock is recovered with `into_inner` so stamps
97/// stay readable, and [`Self::degraded_at`] reports it from then on.
98/// Test: `lock_stall_tests.rs`.
99#[derive(Debug, Default)]
100pub struct LockStallTracker {
101 stalls: Mutex<HashMap<StallKey, Stall>>,
102 poisoned: AtomicBool,
103 last_sweep: Mutex<Option<Instant>>,
104 /// `(last beat, interval)` once a ticker has started.
105 ticker: Mutex<Option<(Instant, Duration)>>,
106}
107
108/// How often to sweep for a given wedge threshold.
109///
110/// Why: a stall is first stamped up to one interval after the holder took the
111/// lock, so detection lags by at most this much. A quarter of the threshold
112/// keeps that lag small next to the threshold itself.
113/// What: `threshold / 4`, clamped to 10 ms ..= 30 s.
114/// Test: `probe_interval_is_a_quarter_of_the_threshold_within_bounds`.
115pub fn probe_interval(threshold: Duration) -> Duration {
116 (threshold / 4).clamp(Duration::from_millis(10), Duration::from_secs(30))
117}
118
119impl LockStallTracker {
120 /// Lock one of the tracker's own mutexes, surviving poison.
121 ///
122 /// Why: a panic inside a critical section must not erase the stamps. That
123 /// would forget a lock that is still held.
124 /// What: returns the guard; on poison sets the flag and recovers the data.
125 /// Test: `a_poisoned_tracker_keeps_its_stamps_and_reports_degraded`.
126 fn guard<'a, T>(&self, m: &'a Mutex<T>) -> MutexGuard<'a, T> {
127 m.lock().unwrap_or_else(|poisoned| {
128 self.poisoned.store(true, Ordering::Release);
129 poisoned.into_inner()
130 })
131 }
132
133 /// Sweep unless a sweep ran within `interval`.
134 ///
135 /// Why: `memory.health` is polled about once a second. Sweeping from it
136 /// keeps stamps fresh even where no ticker runs, and the rate limit keeps
137 /// the cheap path cheap. The ticker sweeps through here too, so the two
138 /// paths share one rate limit instead of doubling each other's work.
139 /// What: records the sweep time under the rate-limit lock, then sweeps.
140 /// Test: `health_sweeps_are_rate_limited_to_the_interval`,
141 /// `a_ticker_sweep_records_itself_against_the_health_rate_limit`.
142 pub fn sweep_if_due(self: &Arc<Self>, registry: &PalaceRegistry, interval: Duration) {
143 let now = Instant::now();
144 {
145 let mut last = self.guard(&self.last_sweep);
146 if last.is_some_and(|t| now.saturating_duration_since(t) < interval) {
147 return;
148 }
149 *last = Some(now);
150 }
151 self.sweep_at(registry, now);
152 }
153
154 /// Observe both handle locks of every open palace at `now`.
155 ///
156 /// What: `peek`s each handle so the LRU order is untouched, then
157 /// [`Self::observe`]s `write_mutex` and `commit_mutex`. Probes keep only the
158 /// mutex `Arc`, never the handle, so idle eviction is unaffected.
159 /// Test: `a_sweep_stamps_both_handle_locks_of_an_open_palace`.
160 pub(crate) fn sweep_at(self: &Arc<Self>, registry: &PalaceRegistry, now: Instant) {
161 for id in registry.list() {
162 let Some(handle) = registry.peek(&id) else {
163 continue;
164 };
165 self.observe(id.as_str(), PalaceLock::Write, &handle.write_mutex, now);
166 self.observe(id.as_str(), PalaceLock::Commit, &handle.commit_mutex, now);
167 }
168 }
169
170 /// Observe one lock: clear an abandoned stamp if it is free, otherwise
171 /// ensure a stamp and a live probe exist.
172 ///
173 /// What: a free lock removes only a stamp with no live probe, or one taken
174 /// against a mutex this sighting has superseded — that probe can never
175 /// report on the live handle. A live probe's stamp on this same mutex is
176 /// left for the probe, which may not yet have queued. A held lock gets a
177 /// stamp (an existing one for this same mutex keeps its older `since`) and
178 /// a spawned probe unless one is already queued. Without a Tokio runtime
179 /// the stamp is kept unprobed, and the next sweep retries.
180 /// Test: `a_free_sighting_never_clears_a_live_probe_stamp`,
181 /// `an_abandoned_probe_keeps_the_stamp_until_the_lock_is_seen_free`,
182 /// `a_reopened_palace_clears_a_stamp_left_by_a_superseded_handle`.
183 pub(crate) fn observe(
184 self: &Arc<Self>,
185 palace: &str,
186 lock: PalaceLock,
187 mutex: &Arc<tokio::sync::Mutex<()>>,
188 now: Instant,
189 ) {
190 let key = (palace.to_string(), lock);
191 if let Ok(free) = mutex.try_lock() {
192 drop(free);
193 let mut stalls = self.guard(&self.stalls);
194 if stalls
195 .get(&key)
196 .is_some_and(|s| !s.probing || !stamped_against(&s.mutex, mutex))
197 {
198 stalls.remove(&key);
199 }
200 return;
201 }
202 let Some(token) = self.claim(key, now, mutex) else {
203 return;
204 };
205 match tokio::runtime::Handle::try_current() {
206 Ok(rt) => {
207 rt.spawn(token.wait(Arc::clone(mutex)));
208 }
209 // Dropping the token marks the stamp unprobed; it is not removed.
210 Err(_) => drop(token),
211 }
212 }
213
214 /// Stamp `key` as probing against `mutex`, returning the probe's token, or
215 /// `None` when a probe on that same mutex is already queued.
216 ///
217 /// What: a stamp taken against a different mutex is replaced outright — the
218 /// palace was reopened, so the older `since` measures a wait on a lock no
219 /// writer can queue on any more.
220 /// Test: `a_reopened_palace_restamps_instead_of_inheriting_the_old_age`.
221 pub(crate) fn claim(
222 self: &Arc<Self>,
223 key: StallKey,
224 now: Instant,
225 mutex: &Arc<tokio::sync::Mutex<()>>,
226 ) -> Option<ProbeToken> {
227 let mut stalls = self.guard(&self.stalls);
228 match stalls.get_mut(&key) {
229 Some(s) if stamped_against(&s.mutex, mutex) => {
230 if s.probing {
231 return None;
232 }
233 s.probing = true;
234 }
235 _ => {
236 stalls.insert(
237 key.clone(),
238 Stall {
239 since: now,
240 probing: true,
241 mutex: Arc::downgrade(mutex),
242 },
243 );
244 }
245 }
246 Some(ProbeToken {
247 tracker: Arc::clone(self),
248 key,
249 mutex: Arc::downgrade(mutex),
250 acquired: false,
251 })
252 }
253
254 /// The oldest stamp's age at `now`, or `None` when nothing is stamped.
255 ///
256 /// Why: an injected `now` lets tests age a stamp without sleeping.
257 /// Test: `a_lock_held_past_the_threshold_ages_past_it_and_clears_on_release`.
258 pub fn oldest_stall_at(&self, now: Instant) -> Option<StalledLock> {
259 let stalls = self.guard(&self.stalls);
260 stalls
261 .iter()
262 .min_by_key(|(_, s)| s.since)
263 .map(|((palace, lock), s)| StalledLock {
264 palace: palace.clone(),
265 lock: *lock,
266 age: now.saturating_duration_since(s.since),
267 })
268 }
269
270 /// Record a ticker heartbeat at `now` for a ticker running every `interval`.
271 pub(crate) fn beat(&self, now: Instant, interval: Duration) {
272 *self.guard(&self.ticker) = Some((now, interval));
273 }
274
275 /// Why the tracker cannot vouch for its stamps at `now`, if it cannot.
276 ///
277 /// Why: a stopped ticker or poisoned state leaves stall detection partial.
278 /// Reporting nothing would make health read `ok` on a signal that is gone.
279 /// What: `Some(reason)` when a started ticker has missed
280 /// three beats (`TICKER_GRACE_INTERVALS`), or when tracking state was
281 /// poisoned. The stamp table is touched first: the poison flag is set by
282 /// [`Self::guard`], so a caller that has not read the stamps yet would
283 /// otherwise see a stale `false` and report healthy.
284 /// Test: `a_ticker_that_stops_beating_reports_degraded`,
285 /// `a_poisoned_tracker_keeps_its_stamps_and_reports_degraded`,
286 /// `degraded_at_sees_poison_without_a_prior_stamp_read`.
287 pub fn degraded_at(&self, now: Instant) -> Option<String> {
288 drop(self.guard(&self.stalls));
289 let beat = *self.guard(&self.ticker);
290 if self.poisoned.load(Ordering::Acquire) {
291 return Some(
292 "palace lock stall tracking was poisoned by a panic; stall ages may be \
293 incomplete (#4001)"
294 .to_string(),
295 );
296 }
297 let (last, interval) = beat?;
298 let silent = now.saturating_duration_since(last);
299 (silent > interval * TICKER_GRACE_INTERVALS).then(|| {
300 format!(
301 "palace lock stall ticker has not run for {}s (interval {}ms); a held \
302 lock may go unnoticed between health polls (#4001)",
303 silent.as_secs(),
304 interval.as_millis()
305 )
306 })
307 }
308}
309
310/// A probe's claim on one stamp.
311///
312/// Why: a probe dropped before it acquires has not seen the lock free, so its
313/// stamp must survive. Removing it would forget a lock that may still be held.
314/// What: [`Self::wait`] queues on the lock and removes the stamp once
315/// acquired. Dropping the token without acquiring only marks the stamp
316/// unprobed, so the next sweep re-probes it and keeps the original `since`.
317/// Both touch the stamp only while it is still the one this token was issued
318/// for: a stamp restamped against a reopened palace's mutex belongs to a
319/// different lock, and clearing it here would erase a live stall.
320/// Test: `an_abandoned_probe_keeps_the_stamp_until_the_lock_is_seen_free`,
321/// `a_reopened_palace_restamps_instead_of_inheriting_the_old_age`.
322#[derive(Debug)]
323pub(crate) struct ProbeToken {
324 tracker: Arc<LockStallTracker>,
325 key: StallKey,
326 /// The mutex this token queues on; see [`Stall::mutex`].
327 mutex: Weak<tokio::sync::Mutex<()>>,
328 acquired: bool,
329}
330
331impl ProbeToken {
332 /// Queue on `mutex` until it is acquired, then clear this token's stamp.
333 ///
334 /// The acquisition is a real one — see the module docs on the ownership
335 /// window it opens between the holder's release and this task's next poll.
336 pub(crate) async fn wait(mut self, mutex: Arc<tokio::sync::Mutex<()>>) {
337 drop(mutex.lock().await);
338 let mut stalls = self.tracker.guard(&self.tracker.stalls);
339 if stalls
340 .get(&self.key)
341 .is_some_and(|s| Weak::ptr_eq(&s.mutex, &self.mutex))
342 {
343 stalls.remove(&self.key);
344 }
345 drop(stalls);
346 self.acquired = true;
347 }
348}
349
350impl Drop for ProbeToken {
351 fn drop(&mut self) {
352 if self.acquired {
353 return;
354 }
355 let mut stalls = self.tracker.guard(&self.tracker.stalls);
356 if let Some(s) = stalls
357 .get_mut(&self.key)
358 .filter(|s| Weak::ptr_eq(&s.mutex, &self.mutex))
359 {
360 s.probing = false;
361 }
362 }
363}
364
365/// Spawn the background sweep for the daemon.
366///
367/// Why: `tm doctor` reads health once. Without a ticker the first sweep would
368/// happen on that read, and a held lock would show age zero.
369/// What: beats, sleeps `interval`, beats, sweeps, forever. The sweep goes
370/// through the same rate limit the health path uses, so a ticker round records
371/// its work and a health poll arriving just after it does not sweep again. A
372/// panicking sweep ends the task, and the stale heartbeat then reads as
373/// degraded.
374/// Test: `the_ticker_stamps_a_held_lock_on_an_open_palace`,
375/// `a_ticker_sweep_records_itself_against_the_health_rate_limit`.
376pub fn spawn_lock_stall_ticker(
377 tracker: Arc<LockStallTracker>,
378 registry: Arc<PalaceRegistry>,
379 interval: Duration,
380) -> tokio::task::JoinHandle<()> {
381 tracker.beat(Instant::now(), interval);
382 tokio::spawn(async move {
383 loop {
384 tokio::time::sleep(interval).await;
385 tracker.beat(Instant::now(), interval);
386 tracker.sweep_if_due(®istry, interval);
387 }
388 })
389}
390
391#[cfg(test)]
392#[path = "lock_stall_tests.rs"]
393mod tests;