Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 5 additions & 24 deletions crates/tracedecay-hooks/src/delivery_spool.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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"),
Expand Down Expand Up @@ -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
Expand Down
44 changes: 38 additions & 6 deletions crates/tracedecay-hooks/src/spool/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand All @@ -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(),
Expand Down
9 changes: 7 additions & 2 deletions crates/tracedecay-mcp/src/handlers/edit_check.rs
Original file line number Diff line number Diff line change
Expand Up @@ -335,8 +335,13 @@ mod tests {
file: "caller.rs".to_owned(),
line: 0,
};
let error = signature_edits(root.path(), &[target.clone()], &[target, caller], &[])
.unwrap_err();
let error = signature_edits(
root.path(),
std::slice::from_ref(&target),
&[target.clone(), caller],
&[],
)
.unwrap_err();
assert!(
matches!(error, TraceDecayError::Io(ref error) if error.kind() == std::io::ErrorKind::NotFound),
"{missing}: {error}"
Expand Down
96 changes: 83 additions & 13 deletions crates/tracedecay-private-fs/src/lock_admission.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(()),
Comment thread
devin-ai-integration[bot] marked this conversation as resolved.
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)
));
}
}
Loading