use super::*;
use std::sync::Arc;
use std::time::{Duration, Instant};
use trusty_common::memory_core::palace::{Palace, PalaceId};
const SETTLE: Duration = Duration::from_secs(10);
async fn wait_for(what: &str, mut cond: impl FnMut() -> bool) {
let started = Instant::now();
while !cond() {
assert!(started.elapsed() < SETTLE, "timed out waiting for {what}");
tokio::time::sleep(Duration::from_millis(5)).await;
}
}
fn registry_with(
id: &str,
) -> (
Arc<PalaceRegistry>,
Arc<trusty_common::memory_core::PalaceHandle>,
tempfile::TempDir,
) {
let tmp = tempfile::tempdir().expect("tempdir");
let registry = Arc::new(PalaceRegistry::new());
let pid = PalaceId::new(id);
let palace = Palace {
id: pid.clone(),
name: id.to_string(),
description: None,
created_at: chrono::Utc::now(),
data_dir: tmp.path().join(pid.as_str()),
};
let handle = registry
.create_palace(tmp.path(), palace)
.expect("create palace");
(registry, handle, tmp)
}
fn key(lock: PalaceLock) -> StallKey {
("p".to_string(), lock)
}
#[test]
fn probe_interval_is_a_quarter_of_the_threshold_within_bounds() {
assert_eq!(
probe_interval(Duration::from_secs(120)),
Duration::from_secs(30)
);
assert_eq!(
probe_interval(Duration::from_secs(40)),
Duration::from_secs(10)
);
assert_eq!(
probe_interval(Duration::from_secs(3600)),
Duration::from_secs(30)
);
assert_eq!(probe_interval(Duration::ZERO), Duration::from_millis(10));
}
#[tokio::test]
async fn a_lock_held_past_the_threshold_ages_past_it_and_clears_on_release() {
let tracker = Arc::new(LockStallTracker::default());
let mutex = Arc::new(tokio::sync::Mutex::new(()));
let held = mutex.clone().lock_owned().await;
let threshold = Duration::from_secs(120);
let t0 = Instant::now();
tracker.observe("p", PalaceLock::Write, &mutex, t0);
let under = tracker
.oldest_stall_at(t0 + Duration::from_secs(53))
.expect("stamped");
assert!(
under.age <= threshold,
"a 53 s hold is not past 120 s: {under:?}"
);
let over = tracker
.oldest_stall_at(t0 + Duration::from_secs(121))
.expect("stamped");
assert!(over.age > threshold, "{over:?}");
assert_eq!((over.palace.as_str(), over.lock), ("p", PalaceLock::Write));
drop(held);
wait_for("the probe to clear the stamp", || {
tracker.oldest_stall_at(Instant::now()).is_none()
})
.await;
}
#[tokio::test]
async fn a_holder_that_panics_releases_and_clears_the_stall() {
let tracker = Arc::new(LockStallTracker::default());
let mutex = Arc::new(tokio::sync::Mutex::new(()));
let (taken_tx, taken_rx) = tokio::sync::oneshot::channel::<()>();
let (go_tx, go_rx) = tokio::sync::oneshot::channel::<()>();
let holder_mutex = mutex.clone();
let holder = tokio::spawn(async move {
let _guard = holder_mutex.lock().await;
let _ = taken_tx.send(());
let _ = go_rx.await;
panic!("holder panics while holding the palace lock");
});
taken_rx.await.expect("holder took the lock");
tracker.observe("p", PalaceLock::Write, &mutex, Instant::now());
assert!(tracker.oldest_stall_at(Instant::now()).is_some());
let _ = go_tx.send(());
assert!(holder.await.expect_err("holder panicked").is_panic());
wait_for("the stamp to clear after the panic", || {
tracker.oldest_stall_at(Instant::now()).is_none()
})
.await;
assert!(
mutex.try_lock().is_ok(),
"the panicking holder released the lock"
);
}
#[tokio::test]
async fn an_abandoned_probe_keeps_the_stamp_until_the_lock_is_seen_free() {
let tracker = Arc::new(LockStallTracker::default());
let mutex = Arc::new(tokio::sync::Mutex::new(()));
let held = mutex.clone().lock_owned().await;
let t0 = Instant::now();
drop(
tracker
.claim(key(PalaceLock::Write), t0, &mutex)
.expect("first claim"),
);
let later = t0 + Duration::from_secs(200);
tracker.observe("p", PalaceLock::Write, &mutex, later);
let stall = tracker
.oldest_stall_at(later)
.expect("the stamp survives the drop");
assert_eq!(
stall.age,
Duration::from_secs(200),
"`since` must stay at t0"
);
assert!(
tracker
.claim(key(PalaceLock::Write), later, &mutex)
.is_none(),
"the re-observe must have queued a live probe"
);
drop(held);
wait_for("the re-spawned probe to clear the stamp", || {
tracker.oldest_stall_at(Instant::now()).is_none()
})
.await;
}
#[tokio::test]
async fn a_free_sighting_never_clears_a_live_probe_stamp() {
let tracker = Arc::new(LockStallTracker::default());
let mutex = Arc::new(tokio::sync::Mutex::new(()));
let now = Instant::now();
let token = tracker
.claim(key(PalaceLock::Commit), now, &mutex)
.expect("claim");
tracker.observe("p", PalaceLock::Commit, &mutex, now);
assert!(
tracker.oldest_stall_at(now).is_some(),
"live probe stamp kept"
);
drop(token);
tracker.observe("p", PalaceLock::Commit, &mutex, now);
assert!(
tracker.oldest_stall_at(now).is_none(),
"unprobed stamp cleared"
);
}
#[tokio::test]
async fn a_reopened_palace_clears_a_stamp_left_by_a_superseded_handle() {
let tracker = Arc::new(LockStallTracker::default());
let superseded = Arc::new(tokio::sync::Mutex::new(()));
let _held = superseded.clone().lock_owned().await;
let t0 = Instant::now();
tracker.observe("p", PalaceLock::Write, &superseded, t0);
assert!(
tracker.oldest_stall_at(t0).is_some(),
"the evicted handle's held lock is stamped"
);
let reopened = Arc::new(tokio::sync::Mutex::new(()));
let later = t0 + Duration::from_secs(200);
tracker.observe("p", PalaceLock::Write, &reopened, later);
assert_eq!(
tracker.oldest_stall_at(later),
None,
"a free lock on the live handle leaves no stamp"
);
}
#[tokio::test]
async fn a_reopened_palace_restamps_instead_of_inheriting_the_old_age() {
let tracker = Arc::new(LockStallTracker::default());
let superseded = Arc::new(tokio::sync::Mutex::new(()));
let _old = superseded.clone().lock_owned().await;
let t0 = Instant::now();
tracker.observe("p", PalaceLock::Write, &superseded, t0);
let reopened = Arc::new(tokio::sync::Mutex::new(()));
let _new = reopened.clone().lock_owned().await;
let later = t0 + Duration::from_secs(200);
tracker.observe("p", PalaceLock::Write, &reopened, later);
let stall = tracker
.oldest_stall_at(later + Duration::from_secs(1))
.expect("the live handle's held lock is stamped");
assert_eq!(
stall.age,
Duration::from_secs(1),
"the stamp must date from the live handle's sighting, not the old one"
);
}
#[tokio::test]
async fn a_poisoned_tracker_keeps_its_stamps_and_reports_degraded() {
let tracker = Arc::new(LockStallTracker::default());
let mutex = Arc::new(tokio::sync::Mutex::new(()));
let _held = mutex.clone().lock_owned().await;
let t0 = Instant::now();
tracker.observe("p", PalaceLock::Write, &mutex, t0);
assert!(tracker.degraded_at(t0).is_none());
let poisoner = Arc::clone(&tracker);
let joined = std::thread::spawn(move || {
let _g = poisoner.stalls.lock();
panic!("poison the stall table");
})
.join();
assert!(joined.is_err(), "the poisoning thread panicked");
let stall = tracker.oldest_stall_at(t0 + Duration::from_secs(1));
assert!(stall.is_some(), "stamps survive poisoning");
let reason = tracker.degraded_at(t0).expect("poisoning is reported");
assert!(reason.contains("poisoned"), "{reason}");
}
#[test]
fn degraded_at_sees_poison_without_a_prior_stamp_read() {
let tracker = Arc::new(LockStallTracker::default());
let poisoner = Arc::clone(&tracker);
let joined = std::thread::spawn(move || {
let _g = poisoner.stalls.lock();
panic!("poison the stall table");
})
.join();
assert!(joined.is_err(), "the poisoning thread panicked");
let reason = tracker
.degraded_at(Instant::now())
.expect("poisoning is reported on the first read");
assert!(reason.contains("poisoned"), "{reason}");
}
#[test]
fn a_ticker_that_stops_beating_reports_degraded() {
let tracker = LockStallTracker::default();
let t0 = Instant::now();
assert!(tracker.degraded_at(t0 + Duration::from_secs(60)).is_none());
tracker.beat(t0, Duration::from_millis(10));
assert!(tracker
.degraded_at(t0 + Duration::from_millis(30))
.is_none());
let reason = tracker
.degraded_at(t0 + Duration::from_millis(31))
.expect("a silent ticker is degraded");
assert!(reason.contains("ticker"), "{reason}");
}
#[tokio::test]
async fn health_sweeps_are_rate_limited_to_the_interval() {
let (registry, handle, _tmp) = registry_with("rate");
let tracker = Arc::new(LockStallTracker::default());
tracker.sweep_if_due(®istry, Duration::from_secs(3600));
let _held = handle.write_mutex.clone().lock_owned().await;
tracker.sweep_if_due(®istry, Duration::from_secs(3600));
assert!(tracker.oldest_stall_at(Instant::now()).is_none());
}
#[tokio::test]
async fn a_sweep_stamps_both_handle_locks_of_an_open_palace() {
let (registry, handle, _tmp) = registry_with("both");
let tracker = Arc::new(LockStallTracker::default());
let _w = handle.write_mutex.clone().lock_owned().await;
let _c = handle.commit_mutex.clone().lock_owned().await;
tracker.sweep_at(®istry, Instant::now());
let stalls = tracker.stalls.lock().expect("not poisoned");
assert!(stalls.contains_key(&("both".to_string(), PalaceLock::Write)));
assert!(stalls.contains_key(&("both".to_string(), PalaceLock::Commit)));
}
#[tokio::test]
async fn the_ticker_stamps_a_held_lock_on_an_open_palace() {
let (registry, handle, _tmp) = registry_with("ticked");
let tracker = Arc::new(LockStallTracker::default());
let _c = handle.commit_mutex.clone().lock_owned().await;
let ticker = spawn_lock_stall_ticker(Arc::clone(&tracker), registry, Duration::from_millis(10));
wait_for("the ticker to stamp the held commit lock", || {
tracker
.oldest_stall_at(Instant::now())
.is_some_and(|s| s.palace == "ticked" && s.lock == PalaceLock::Commit)
})
.await;
assert!(tracker.degraded_at(Instant::now()).is_none());
ticker.abort();
}
#[tokio::test]
async fn a_ticker_sweep_records_itself_against_the_health_rate_limit() {
let (registry, handle, _tmp) = registry_with("shared");
let tracker = Arc::new(LockStallTracker::default());
let _c = handle.commit_mutex.clone().lock_owned().await;
let ticker = spawn_lock_stall_ticker(
Arc::clone(&tracker),
Arc::clone(®istry),
Duration::from_millis(10),
);
wait_for("the ticker to sweep once", || {
tracker.oldest_stall_at(Instant::now()).is_some()
})
.await;
ticker.abort();
assert!(
ticker
.await
.expect_err("the ticker was aborted")
.is_cancelled(),
"the ticker must be stopped before the health sweep below"
);
let _w = handle.write_mutex.clone().lock_owned().await;
tracker.sweep_if_due(®istry, Duration::from_secs(3600));
let stalls = tracker.stalls.lock().expect("not poisoned");
assert!(
!stalls.contains_key(&("shared".to_string(), PalaceLock::Write)),
"the health sweep must be rate limited by the ticker's own sweep: {stalls:?}"
);
}