From 526517654488c87371d421a2b7a81be59cf32aa4 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Fri, 2 Oct 2026 20:19:24 +0000 Subject: [PATCH 1/2] test(sessions): drive run_with_semaphore in the deadline detach test The test reimplemented the semaphore and poll loop inline, so it never exercised run_with_semaphore and passed with the permit dropped early. It now calls run_with_semaphore and retries when the wall-clock deadline expires before the closure is spawned, asserting that such an attempt leaks no permit. --- .../backfill/bounded/blocking/tests.rs | 79 +++++++------------ 1 file changed, 29 insertions(+), 50 deletions(-) diff --git a/crates/tracedecay-sessions/src/runtime/git_correlation/backfill/bounded/blocking/tests.rs b/crates/tracedecay-sessions/src/runtime/git_correlation/backfill/bounded/blocking/tests.rs index 1e0174b940..a8f0c0c602 100644 --- a/crates/tracedecay-sessions/src/runtime/git_correlation/backfill/bounded/blocking/tests.rs +++ b/crates/tracedecay-sessions/src/runtime/git_correlation/backfill/bounded/blocking/tests.rs @@ -137,60 +137,39 @@ async fn detached_task_holds_capacity_until_its_closure_finishes() { #[tokio::test(flavor = "current_thread")] async fn deadline_detaches_without_leaking_the_permit() { let permits = Arc::new(Semaphore::new(1)); - let (started_tx, started_rx) = oneshot::channel(); - let (closure_running_tx, closure_running_rx) = oneshot::channel(); - let (release_tx, release_rx) = mpsc::channel(); - let run_permits = Arc::clone(&permits); - let run = tokio::spawn(async move { - let permit = run_permits - .acquire_owned() - .await - .map_err(|_| BoundedBackfillInterruption::SourceUnavailable)?; - let join = tokio::task::spawn_blocking(move || { - let _permit = permit; - started_tx.send(()).unwrap(); - closure_running_tx.send(()).unwrap(); - release_rx.recv().unwrap(); - Ok::<_, BoundedBackfillInterruption>(()) - }); - tokio::pin!(join); - closure_running_rx.await.unwrap(); + // The deadline is wall-clock, so a loaded host can expire it before the + // closure is spawned. That attempt must leak nothing; retry until the + // closure is running when the deadline detaches it. + loop { + let (started_tx, started_rx) = oneshot::channel(); + let (release_tx, release_rx) = mpsc::channel(); let deadline_control = control( ObservationCancellation::default(), Duration::from_millis(25), ); - loop { - tokio::select! { - biased; - result = &mut join => { - return result - .map_err(|_| BoundedBackfillInterruption::SourceUnavailable)? - .and_then(|value| { - deadline_control.check()?; - Ok(value) - }); - } - () = tokio::time::sleep(Duration::from_millis(10)) => { - deadline_control.check()?; - } - } - } - }); - started_rx.await.unwrap(); - assert_eq!( - tokio::time::timeout(Duration::from_secs(1), run) + let run_permits = Arc::clone(&permits); + let run = tokio::spawn(async move { + run_with_semaphore(run_permits, &deadline_control, move || { + started_tx.send(()).unwrap(); + release_rx.recv().unwrap(); + Ok::<_, BoundedBackfillInterruption>(()) + }) .await - .unwrap() - .unwrap(), - Err(BoundedBackfillInterruption::CommandTimedOut) - ); - assert_eq!(permits.available_permits(), 0); + }); + let started = started_rx.await.is_ok(); + assert_eq!( + run.await.unwrap(), + Err(BoundedBackfillInterruption::CommandTimedOut) + ); + if !started { + assert_eq!(permits.available_permits(), 1); + continue; + } + assert_eq!(permits.available_permits(), 0); - release_tx.send(()).unwrap(); - let permit = tokio::time::timeout(Duration::from_secs(1), Arc::clone(&permits).acquire_owned()) - .await - .unwrap() - .unwrap(); - drop(permit); - assert_eq!(permits.available_permits(), 1); + release_tx.send(()).unwrap(); + drop(Arc::clone(&permits).acquire_owned().await.unwrap()); + assert_eq!(permits.available_permits(), 1); + return; + } } From e4dbb0c2631a970d6ed589e1f148b01487b3744b Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Fri, 2 Oct 2026 21:14:58 +0000 Subject: [PATCH 2/2] test(sessions): bound the permit reacquire so a leak fails --- .../git_correlation/backfill/bounded/blocking/tests.rs | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/crates/tracedecay-sessions/src/runtime/git_correlation/backfill/bounded/blocking/tests.rs b/crates/tracedecay-sessions/src/runtime/git_correlation/backfill/bounded/blocking/tests.rs index a8f0c0c602..b7e276cac6 100644 --- a/crates/tracedecay-sessions/src/runtime/git_correlation/backfill/bounded/blocking/tests.rs +++ b/crates/tracedecay-sessions/src/runtime/git_correlation/backfill/bounded/blocking/tests.rs @@ -168,7 +168,12 @@ async fn deadline_detaches_without_leaking_the_permit() { assert_eq!(permits.available_permits(), 0); release_tx.send(()).unwrap(); - drop(Arc::clone(&permits).acquire_owned().await.unwrap()); + let permit = + tokio::time::timeout(Duration::from_secs(5), Arc::clone(&permits).acquire_owned()) + .await + .expect("the detached closure releases its permit") + .unwrap(); + drop(permit); assert_eq!(permits.available_permits(), 1); return; }