Skip to main content

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(&registry, interval);
387        }
388    })
389}
390
391#[cfg(test)]
392#[path = "lock_stall_tests.rs"]
393mod tests;