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..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 @@ -137,60 +137,44 @@ 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(); + 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; + } }