Skip to content
Merged
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
39 changes: 23 additions & 16 deletions crates/tracedecay-runtime-core/src/shard_runtime/registry/close.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<usize, StoreRuntimeRegistryFailure> {
self.cancel_and_join_opens_for_shutdown().await;
Expand Down Expand Up @@ -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::<Vec<_>>();
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")]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(&registry).await;
let first = open_published(&registry, code_request("worktree.shutdown-first", &pin)).await;
let second = open_published(&registry, 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 { .. }
));
}
}
Loading