diff --git a/crates/tracedecay-hooks/src/delivery_spool.rs b/crates/tracedecay-hooks/src/delivery_spool.rs index 1f767094d6..f5a0db60ff 100644 --- a/crates/tracedecay-hooks/src/delivery_spool.rs +++ b/crates/tracedecay-hooks/src/delivery_spool.rs @@ -276,23 +276,10 @@ impl HookDeliveryReceiptSpoolV1 { let root = root.into(); ensure_root(&root)?; let lock = open_lock_file(&root, LOCK_FILE)?; - let try_lock_result = { - let _span = tracing::trace_span!("hooks.delivery.lock.try_lock").entered(); - lock.try_lock() - }; - match try_lock_result { - Ok(()) => {} - Err(std::fs::TryLockError::WouldBlock) => { - if wait_budget.is_zero() { - return Err(HookDeliverySpoolError::Busy); - } - lock_until(&lock, Instant::now() + wait_budget).map_err(|error| match error { - LockAdmissionError::TimedOut => HookDeliverySpoolError::Busy, - LockAdmissionError::Io(_) => HookDeliverySpoolError::Io, - })?; - } - Err(std::fs::TryLockError::Error(_)) => return Err(HookDeliverySpoolError::Io), - } + lock_until(&lock, Instant::now() + wait_budget).map_err(|error| match error { + LockAdmissionError::TimedOut => HookDeliverySpoolError::Busy, + LockAdmissionError::Io(_) => HookDeliverySpoolError::Io, + })?; let spool = Self { root, _lock: FileLease::held(lock, "hooks.delivery.writer"), @@ -834,13 +821,7 @@ mod tests { HookDeliveryReceiptSpoolV1::open(&root.0, Duration::ZERO).unwrap_err(), HookDeliverySpoolError::Busy ); - // An exhausted budget never admits, even when the lease is free. - assert_eq!( - HookDeliveryReceiptWriterV1::open_within(&root.0, Duration::ZERO).unwrap_err(), - HookDeliverySpoolError::AdmissionTimedOut - ); - let writer = - HookDeliveryReceiptWriterV1::open_within(&root.0, Duration::from_millis(20)).unwrap(); + let writer = HookDeliveryReceiptWriterV1::open_within(&root.0, Duration::ZERO).unwrap(); assert_eq!( writer.retain(&receipt()).unwrap(), HookDeliveryRetentionV1::Staged diff --git a/crates/tracedecay-hooks/src/spool/tests.rs b/crates/tracedecay-hooks/src/spool/tests.rs index 4334e7bba2..efdbe49b21 100644 --- a/crates/tracedecay-hooks/src/spool/tests.rs +++ b/crates/tracedecay-hooks/src/spool/tests.rs @@ -1699,6 +1699,44 @@ fn settling_and_reclaiming_hold_the_writer_lease_across_no_durability_barrier() assert_eq!(fs::metadata(records_path(&root.0)).unwrap().len(), 0); } +/// Bounded admission waits only on contention. A writer whose budget is +/// already spent, like a callback or a drain descheduled past its deadline, +/// still takes a free lease, including the lease a settlement reacquires +/// after its barriers; a held lease refuses it without waiting. +#[test] +fn a_writer_whose_budget_is_spent_takes_a_free_lease_but_never_a_held_one() { + let root = TestDir::new("spent-budget"); + let (owner, _) = HookSpoolV1::open(&root.0, config(), UtcMicros(10)).unwrap(); + assert_eq!( + HookSpoolV1::open_within(&root.0, config(), UtcMicros(11), Duration::ZERO).unwrap_err(), + HookSpoolError::AdmissionTimedOut + ); + drop(owner); + + let (mut admitted, _) = + HookSpoolV1::open_within(&root.0, config(), UtcMicros(11), Duration::ZERO) + .expect("a free lease admits a writer with no budget left"); + let record = admitted + .append(envelope(1, 9), &binding(), UtcMicros(11)) + .unwrap(); + assert!( + admitted + .acknowledge( + HookSpoolAckV1 { + sequence: record.sequence, + receipt_id: [41; 16], + disposition: HookSpoolAckDispositionV1::Committed, + }, + UtcMicros(11), + ) + .unwrap() + ); + drop(admitted); + let (_, report) = HookSpoolV1::open(&root.0, config(), UtcMicros(12)).unwrap(); + assert_eq!(report.pending_records, 0); + assert_eq!(report.next_sequence, 2); +} + #[test] fn bounded_writer_admission_preserves_failfast_and_times_out_without_mutation() { let root = TestDir::new("bounded-admission"); @@ -1720,12 +1758,6 @@ fn bounded_writer_admission_preserves_failfast_and_times_out_without_mutation() ); assert_eq!(fs::read(meta_path(&root.0)).unwrap(), before); drop(owner); - // An exhausted budget never admits, even when the lease is free. - assert_eq!( - HookSpoolV1::open_within(&root.0, config(), UtcMicros(11), std::time::Duration::ZERO) - .unwrap_err(), - HookSpoolError::AdmissionTimedOut - ); let (mut admitted, _) = HookSpoolV1::open_within( &root.0, config(), diff --git a/crates/tracedecay-private-fs/src/lock_admission.rs b/crates/tracedecay-private-fs/src/lock_admission.rs index 84400bede9..eeeff20b0a 100644 --- a/crates/tracedecay-private-fs/src/lock_admission.rs +++ b/crates/tracedecay-private-fs/src/lock_admission.rs @@ -12,45 +12,115 @@ pub enum LockAdmissionError { // Avoid spinning while the admitted writer completes durable filesystem work. const LOCK_POLL_INTERVAL: Duration = Duration::from_millis(1); -/// Exclusive-locks `file`, waiting until `deadline`. -/// A lock taken after the deadline is released and the call times out. +/// Exclusive-locks `file`, waiting until `deadline` for any other holder. #[tracing::instrument(name = "private_fs.lock.admission", level = "trace", skip_all)] pub fn lock_until(file: &File, deadline: Instant) -> Result<(), LockAdmissionError> { admit_until(file, deadline, File::try_lock) } /// Shared-locks `file`, waiting until `deadline` for any exclusive holder. -/// A lock taken after the deadline is released and the call times out. #[tracing::instrument(name = "private_fs.lock.shared_admission", level = "trace", skip_all)] pub fn lock_shared_until(file: &File, deadline: Instant) -> Result<(), LockAdmissionError> { admit_until(file, deadline, File::try_lock_shared) } +/// The deadline bounds waiting on contention, not the caller's own latency: +/// the lock is always attempted once, so a caller descheduled past its +/// deadline still takes a free lock. Once contended, every retry happens +/// before the deadline, so a holder that outlasts the deadline times out. fn admit_until( file: &File, deadline: Instant, try_lock: impl Fn(&File) -> Result<(), std::fs::TryLockError>, ) -> Result<(), LockAdmissionError> { loop { - if Instant::now() >= deadline { - return Err(LockAdmissionError::TimedOut); - } match try_lock(file) { - Ok(()) => { - if Instant::now() >= deadline { - file.unlock().map_err(LockAdmissionError::Io)?; - return Err(LockAdmissionError::TimedOut); - } - return Ok(()); - } + Ok(()) => return Ok(()), Err(std::fs::TryLockError::WouldBlock) => { let remaining = deadline.saturating_duration_since(Instant::now()); if remaining.is_zero() { return Err(LockAdmissionError::TimedOut); } std::thread::sleep(remaining.min(LOCK_POLL_INTERVAL)); + if Instant::now() >= deadline { + return Err(LockAdmissionError::TimedOut); + } } Err(std::fs::TryLockError::Error(error)) => return Err(LockAdmissionError::Io(error)), } } } + +#[cfg(test)] +mod tests { + use std::fs::OpenOptions; + use std::path::Path; + + use super::*; + + fn open(path: &Path) -> File { + OpenOptions::new() + .read(true) + .write(true) + .create(true) + .truncate(false) + .open(path) + .unwrap() + } + + #[test] + fn a_free_lock_is_admitted_after_the_deadline_has_passed() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("admission.lock"); + let elapsed = Instant::now(); + + let exclusive = open(&path); + assert!(lock_until(&exclusive, elapsed).is_ok()); + exclusive.unlock().unwrap(); + + let shared = open(&path); + assert!(lock_shared_until(&shared, elapsed).is_ok()); + let second_reader = open(&path); + assert!(lock_shared_until(&second_reader, elapsed).is_ok()); + } + + #[test] + fn a_held_lock_times_out_once_the_deadline_passes() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("admission.lock"); + let holder = open(&path); + holder.lock().unwrap(); + let waiter = open(&path); + + assert!(matches!( + lock_until(&waiter, Instant::now()), + Err(LockAdmissionError::TimedOut) + )); + assert!(matches!( + lock_shared_until(&waiter, Instant::now() + Duration::from_millis(20)), + Err(LockAdmissionError::TimedOut) + )); + + holder.unlock().unwrap(); + assert!(lock_until(&waiter, Instant::now()).is_ok()); + } + + #[test] + fn a_holder_releasing_after_the_deadline_never_admits_a_contended_waiter() { + let dir = tempfile::tempdir().unwrap(); + let file = open(&dir.path().join("admission.lock")); + let deadline = Instant::now() + Duration::from_millis(5); + let released_after_deadline = |_: &File| { + if Instant::now() < deadline { + Err(std::fs::TryLockError::WouldBlock) + } else { + Ok(()) + } + }; + + assert!(matches!( + admit_until(&file, deadline, released_after_deadline), + Err(LockAdmissionError::TimedOut) + )); + } +}