diff --git a/crates/tracedecay-runtime-core/src/shard_runtime/registry/close.rs b/crates/tracedecay-runtime-core/src/shard_runtime/registry/close.rs index 81e4e108f0..29533ff949 100644 --- a/crates/tracedecay-runtime-core/src/shard_runtime/registry/close.rs +++ b/crates/tracedecay-runtime-core/src/shard_runtime/registry/close.rs @@ -55,7 +55,10 @@ impl StoreRuntimeRegistry { /// no destructor closes these attachments otherwise. In-flight opens are /// cancelled and joined first, so none publishes or keeps a physical /// handle after the scan. Held runtimes stay mounted and are logged with - /// their blockers. Returns the number of runtimes closed. + /// their blockers. Each reserved runtime owns its own writer and readers, + /// so they close concurrently and the shutdown waits for the slowest + /// checkpoint rather than their sum. Returns the number of runtimes + /// closed. #[hotpath::measure(label = "runtime_core.registry.close_idle_for_shutdown", future = true)] pub async fn close_idle_for_shutdown(&self) -> Result { self.cancel_and_join_opens_for_shutdown().await; @@ -104,22 +107,26 @@ impl StoreRuntimeRegistry { } (reservations, reserve_failure) }; - let registry = self.clone(); - tokio::task::spawn_blocking(move || { - let closed = reservations.len(); - let mut first_failure = reserve_failure; - for reservation in reservations { - if let Err(failure) = registry.complete_eviction(reservation) { - first_failure.get_or_insert(failure); - } + let closes = reservations + .into_iter() + .map(|reservation| { + let registry = self.clone(); + tokio::task::spawn_blocking(move || registry.complete_eviction(reservation)) + }) + .collect::>(); + let closed = closes.len(); + let mut first_failure = reserve_failure; + for close in closes { + if let Err(failure) = close.await.unwrap_or_else(|error| { + Err(StoreRuntimeRegistryFailure::PhysicalRuntimeFailed { + operation: "join shutdown close of idle registered runtimes", + message: error.to_string(), + }) + }) { + first_failure.get_or_insert(failure); } - first_failure.map_or(Ok(closed), Err) - }) - .await - .map_err(|error| StoreRuntimeRegistryFailure::PhysicalRuntimeFailed { - operation: "join shutdown close of idle registered runtimes", - message: error.to_string(), - })? + } + first_failure.map_or(Ok(closed), Err) } #[hotpath::measure(label = "runtime_core.registry.close_path")] diff --git a/crates/tracedecay-runtime-core/src/shard_runtime/registry/tests/attachment.rs b/crates/tracedecay-runtime-core/src/shard_runtime/registry/tests/attachment.rs index a48691dedb..f4e163d963 100644 --- a/crates/tracedecay-runtime-core/src/shard_runtime/registry/tests/attachment.rs +++ b/crates/tracedecay-runtime-core/src/shard_runtime/registry/tests/attachment.rs @@ -451,3 +451,57 @@ async fn failed_blocking_drain_restores_fault_and_wakes_reserved_open_joiners() StoreRuntimeLookup::Evicting { .. } )); } + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn shutdown_close_drains_every_idle_runtime_concurrently() { + let publisher = Arc::new(AttachmentPublisher::default()); + let registry = StoreRuntimeRegistry::with_config( + Arc::new(TestResolver::default()), + publisher.clone(), + StoreRuntimeRegistryConfig::new(2).unwrap(), + ) + .unwrap(); + let pin = profile_pin(®istry).await; + let first = open_published(®istry, code_request("worktree.shutdown-first", &pin)).await; + let second = open_published(®istry, code_request("worktree.shutdown-second", &pin)).await; + let bindings = [first.binding().clone(), second.binding().clone()]; + let controls = [publisher.control(1), publisher.control(2)]; + drop((first, second)); + for control in &controls { + control.drain_gate.block(); + } + + let closing = registry.clone(); + let close = tokio::spawn(async move { closing.close_idle_for_shutdown().await }); + // A close that ran one runtime at a time would hold the second runtime's + // drain behind the first, which stays blocked until both have entered. + let entered = || { + controls + .iter() + .filter(|control| control.drain_gate.entered.load(Ordering::SeqCst)) + .count() + }; + let _ = tokio::time::timeout(Duration::from_secs(2), async { + while entered() < controls.len() { + tokio::task::yield_now().await; + } + }) + .await; + let concurrent_drains = entered(); + for control in &controls { + control.drain_gate.release(); + } + + assert_eq!(close.await.unwrap().unwrap(), 2); + assert_eq!( + concurrent_drains, 2, + "every idle runtime drains while the others are still draining" + ); + for (binding, control) in bindings.iter().zip(&controls) { + assert_eq!(control.close_calls.load(Ordering::SeqCst), 1); + assert!(matches!( + registry.lookup(binding), + StoreRuntimeLookup::Missing { .. } + )); + } +}