From 814866e3022bd41682954ee6239d611904a34c24 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 29 Sep 2026 11:45:19 +0000 Subject: [PATCH 1/3] fix(store-runtime): cancel in-flight store opens at daemon shutdown A daemon stopped while a project open was inside its session-store mount waited for the mount to finish: the registry open ran detached, the schema install could not be interrupted, and the shutdown close skipped runtimes still opening. Under CPU starvation the mount outlived the 2 s cooperative drain, the open was aborted, and the process missed its exit. The store runtime registry now owns a shutdown cancellation. The project open shutdown owner cancels it for every mounted session registry. New opens are refused with a typed OpenCancelled; an in-flight open stops before resolve and before publish, and one that finishes publishing after the cancel closes its unpublished runtime instead of publishing it. The registered schema install, both at initialization and at daemon attach, checks the token before every statement of its admission transaction and rolls back with a typed store_open_cancelled, so a store is left either with its prior schema or fully installed. The shutdown close cancels and joins in-flight opens before it closes idle runtimes, so none is skipped. --- crates/tracedecay-domain/src/errors.rs | 17 ++ crates/tracedecay-global-db/src/registered.rs | 20 ++- .../tracedecay-global-db/src/schema_stages.rs | 98 +++++++++- crates/tracedecay-global-db/src/tests.rs | 2 + .../tracedecay-global-db/src/tests/harness.rs | 12 +- .../src/tests/schema_install_cancellation.rs | 167 ++++++++++++++++++ crates/tracedecay-runtime-core/src/ports.rs | 32 +++- .../src/shard_runtime/registry.rs | 23 +++ .../src/shard_runtime/registry/close.rs | 8 +- .../src/shard_runtime/registry/open.rs | 70 ++++++++ .../src/shard_runtime/registry/ports.rs | 27 ++- .../src/shard_runtime/registry/tests.rs | 63 +++++++ .../src/session_registry.rs | 3 + .../src/session_registry/maintenance.rs | 6 +- .../src/session_registry/mounts.rs | 7 + .../branch_admin/session_runtime_shutdown.rs | 11 ++ .../tracedecay/src/daemon/engine/shutdown.rs | 6 +- 17 files changed, 540 insertions(+), 32 deletions(-) create mode 100644 crates/tracedecay-global-db/src/tests/schema_install_cancellation.rs diff --git a/crates/tracedecay-domain/src/errors.rs b/crates/tracedecay-domain/src/errors.rs index 479efe2aa7..b853b26222 100644 --- a/crates/tracedecay-domain/src/errors.rs +++ b/crates/tracedecay-domain/src/errors.rs @@ -2,6 +2,8 @@ use thiserror::Error; use crate::ApplicationProblemDetailV1; +const STORE_OPEN_CANCELLED_REASON_CODE: &str = "store_open_cancelled"; + #[derive(Error, Debug)] #[error("{detail}")] struct HookRuntimeErrorContext { @@ -335,6 +337,21 @@ impl TraceDecayError { } } + /// Daemon shutdown stopped a store open or schema install at a safe point + /// and rolled back its uncommitted work. + pub fn store_open_cancelled(operation: impl std::fmt::Display) -> Self { + Self::project_route( + STORE_OPEN_CANCELLED_REASON_CODE, + true, + format!("{operation} was cancelled by daemon shutdown"), + ) + } + + pub fn is_store_open_cancelled(&self) -> bool { + self.project_route_context() + .is_some_and(|(reason_code, _, _)| reason_code == STORE_OPEN_CANCELLED_REASON_CODE) + } + pub fn project_route_with_detail( reason_code: impl Into, retryable: bool, diff --git a/crates/tracedecay-global-db/src/registered.rs b/crates/tracedecay-global-db/src/registered.rs index bd1d423823..caa1d4ea54 100644 --- a/crates/tracedecay-global-db/src/registered.rs +++ b/crates/tracedecay-global-db/src/registered.rs @@ -4,6 +4,7 @@ use std::sync::{Arc, OnceLock, RwLock, Weak}; use tracedecay_domain::errors::TraceDecayError; use tracedecay_runtime_core::{ + cancellation::CancellationToken, db::{ Database, DatabaseAuthority, DatabaseEngineReadConnection, DatabaseEngineReadSnapshot, DatabaseOwnerErrorV1, DatabaseOwnerRetirementReservationV1, DatabaseOwnerV1, @@ -98,8 +99,12 @@ impl RegisteredGlobalDbOwnerV1 { ) -> tracedecay_domain::errors::Result { let temporary = database.issue_lease().map_err(registered_owner_error)?; let registered = RegisteredGlobalDb::from_owned_database(temporary); - let (_, refused_authority) = - super::schema_stages::ensure_attached_registered_schema(®istered.database).await?; + // Short-lived attaches have no daemon shutdown to observe. + let (_, refused_authority) = super::schema_stages::ensure_attached_registered_schema( + ®istered.database, + &CancellationToken::new(), + ) + .await?; // A store refused for reset is never converged: its reset deletes it. if refused_authority.is_none() { super::schema_stages::converge_attached_registered_schema(®istered.database).await?; @@ -114,16 +119,23 @@ impl RegisteredGlobalDbOwnerV1 { } /// Returns the resumable convergence plan for an already admitted schema - /// without retaining an unowned client lease. + /// without retaining an unowned client lease. A `cancellation` observed + /// before the admission transaction commits rolls it back and fails with + /// [`TraceDecayError::store_open_cancelled`]. #[hotpath::measure(future = true, label = "global_db.registered.admit_daemon")] pub async fn admit_and_attach_for_daemon( database: DatabaseOwnerV1, + cancellation: &CancellationToken, ) -> tracedecay_domain::errors::Result<(Self, super::schema_stages::RegisteredSchemaConvergence)> { let temporary = database.issue_lease().map_err(registered_owner_error)?; let registered = RegisteredGlobalDb::from_owned_database(temporary); let (convergence, refused_authority) = - super::schema_stages::ensure_attached_registered_schema(®istered.database).await?; + super::schema_stages::ensure_attached_registered_schema( + ®istered.database, + cancellation, + ) + .await?; drop(registered); Ok(( Self { diff --git a/crates/tracedecay-global-db/src/schema_stages.rs b/crates/tracedecay-global-db/src/schema_stages.rs index b0b25f5561..0bc5b5eb76 100644 --- a/crates/tracedecay-global-db/src/schema_stages.rs +++ b/crates/tracedecay-global-db/src/schema_stages.rs @@ -13,9 +13,10 @@ use super::{ }; use crate::registered::RefusedAuthorityV1; use tracedecay_runtime_core::{ + cancellation::CancellationToken, db::{ Database, DatabaseWriteTransaction, - engine::{Executor, QueryExecutor}, + engine::{self, Executor, IntoParams, QueryExecutor, Rows, WriteStatement}, }, ports::registered_schema::{ RegisteredSchemaInstallationTransactionV1, RegisteredSchemaInstallationV1, @@ -545,6 +546,7 @@ pub async fn ensure_registered_schema_for_admission( } = classify_registered_schema_admission(installation).await?; let is_fresh = configuration_fresh.is_some(); let force_exhaustive = !authority_invariant_triggers_intact(installation).await?; + schema_install_checkpoint(installation.cancellation(), OPERATION)?; let transaction = installation .begin_atomic_schema_transaction() .await @@ -556,6 +558,7 @@ pub async fn ensure_registered_schema_for_admission( temporal_admission, workflow_admission, force_exhaustive, + installation.cancellation(), "commit registered global schema", "roll back registered global schema", ) @@ -618,6 +621,7 @@ async fn install_and_commit_registered_schema( temporal_admission: session_temporal_schema::SessionTemporalSchemaAdmission, workflow_admission: WorkflowSchemaAdmission, force_exhaustive: bool, + cancellation: &CancellationToken, commit_operation: &'static str, rollback_operation: &'static str, ) -> tracedecay_domain::errors::Result> @@ -625,13 +629,19 @@ where T: Executor + Sync + SchemaInstallTransaction, { let admission = install_registered_schema_stages( - &transaction, + &CancellableSchemaExecutor { + inner: &transaction, + cancellation, + }, configuration_fresh, temporal_admission, workflow_admission, force_exhaustive, ) - .await; + .await + .and_then(|refused_authority| { + schema_install_checkpoint(cancellation, commit_operation).map(|()| refused_authority) + }); match admission { Ok(refused_authority) => { transaction.commit().await.map_err(|error| { @@ -640,6 +650,11 @@ where Ok(refused_authority) } Err(error) => match transaction.rollback().await { + Ok(()) if cancellation.is_cancelled() => Err( + tracedecay_domain::errors::TraceDecayError::store_open_cancelled( + rollback_operation, + ), + ), Ok(()) => Err(error), Err(rollback_error) => Err(global_db_operation_error( rollback_operation, @@ -649,6 +664,68 @@ where } } +fn schema_install_checkpoint( + cancellation: &CancellationToken, + operation: &'static str, +) -> tracedecay_domain::errors::Result<()> { + if cancellation.is_cancelled() { + return Err(tracedecay_domain::errors::TraceDecayError::store_open_cancelled(operation)); + } + Ok(()) +} + +/// Refuses the next statement once shutdown cancels the Store open, so the +/// enclosing schema transaction rolls back at a statement boundary instead of +/// running the whole install first. +struct CancellableSchemaExecutor<'a, T> { + inner: &'a T, + cancellation: &'a CancellationToken, +} + +impl CancellableSchemaExecutor<'_, T> { + fn checkpoint(&self) -> engine::Result<()> { + if self.cancellation.is_cancelled() { + return Err(engine::Error::InvalidOperation( + "registered schema install cancelled by daemon shutdown".to_owned(), + )); + } + Ok(()) + } +} + +impl QueryExecutor for CancellableSchemaExecutor<'_, T> { + async fn query

(&self, sql: &str, params: P) -> engine::Result + where + P: IntoParams, + { + self.checkpoint()?; + self.inner.query(sql, params).await + } +} + +impl Executor for CancellableSchemaExecutor<'_, T> { + async fn execute

(&self, sql: &str, params: P) -> engine::Result + where + P: IntoParams, + { + self.checkpoint()?; + self.inner.execute(sql, params).await + } + + async fn execute_statements( + &self, + statements: Vec, + ) -> engine::Result> { + self.checkpoint()?; + self.inner.execute_statements(statements).await + } + + async fn execute_batch(&self, sql: &str) -> engine::Result<()> { + self.checkpoint()?; + self.inner.execute_batch(sql).await + } +} + trait SchemaInstallTransaction: Sized { fn commit(self) -> impl std::future::Future> + Send; fn rollback(self) -> impl std::future::Future> + Send; @@ -1004,10 +1081,19 @@ pub async fn converge_attached_registered_schema( /// same work synchronously through [`converge_attached_registered_schema`]. /// A store whose observation rows this binary refuses is still admitted for /// its other authorities and returns that refused authority beside the plan. +/// +/// `cancellation` is observed before classification and before every +/// statement of the admission transaction: a cancelled attach rolls that +/// transaction back and fails typed, so the store keeps exactly its prior +/// schema. Once committed, the idempotent index builds and validation run to +/// completion. #[hotpath::measure(future = true, label = "global_db.schema.persist.attach")] pub(crate) async fn ensure_attached_registered_schema( database: &Database, + cancellation: &CancellationToken, ) -> tracedecay_domain::errors::Result<(RegisteredSchemaConvergence, Option)> { + const OPERATION: &str = "install attached registered global database schema"; + schema_install_checkpoint(cancellation, OPERATION)?; let read_connection = database.read_connection(); let RegisteredSchemaAdmissionClassification { configuration_fresh, @@ -1015,15 +1101,15 @@ pub(crate) async fn ensure_attached_registered_schema( workflow_admission, } = classify_registered_schema_admission(&read_connection).await?; let force_exhaustive = !authority_invariant_triggers_intact(&read_connection).await?; - let transaction = database - .begin_bulk_write_transaction("install attached registered global database schema") - .await?; + schema_install_checkpoint(cancellation, OPERATION)?; + let transaction = database.begin_bulk_write_transaction(OPERATION).await?; let refused_authority = install_and_commit_registered_schema( transaction, configuration_fresh.as_ref(), temporal_admission, workflow_admission, force_exhaustive, + cancellation, "commit attached registered global schema", "roll back attached registered global schema", ) diff --git a/crates/tracedecay-global-db/src/tests.rs b/crates/tracedecay-global-db/src/tests.rs index 317b42fc17..e34b4e36b1 100644 --- a/crates/tracedecay-global-db/src/tests.rs +++ b/crates/tracedecay-global-db/src/tests.rs @@ -12,6 +12,8 @@ pub(crate) mod lcm_privacy_rescan; #[cfg(test)] mod lcm_schema; #[cfg(test)] +mod schema_install_cancellation; +#[cfg(test)] mod session_sync; #[cfg(test)] diff --git a/crates/tracedecay-global-db/src/tests/harness.rs b/crates/tracedecay-global-db/src/tests/harness.rs index cb1ad44b5c..88cc7951ab 100644 --- a/crates/tracedecay-global-db/src/tests/harness.rs +++ b/crates/tracedecay-global-db/src/tests/harness.rs @@ -11,7 +11,7 @@ use tracedecay_runtime_core::db::DaemonDatabaseScope; #[cfg(test)] use tracedecay_runtime_core::db::engine::{Executor, IntoParams, QueryExecutor, Rows}; -static TEST_RUNTIME_NONCE: AtomicU64 = AtomicU64::new(1); +pub(super) static TEST_RUNTIME_NONCE: AtomicU64 = AtomicU64::new(1); #[cfg(test)] static HOST_ADMISSION_TEST_RESIDENT_MEMORY: OnceLock< Arc, @@ -492,10 +492,12 @@ impl RegisteredGlobalDbHarness { .await .expect("publish daemon test runtime") .into_parts(); - let (database, convergence) = - RegisteredGlobalDbOwnerV1::admit_and_attach_for_daemon(database_owner) - .await - .expect("daemon admission"); + let (database, convergence) = RegisteredGlobalDbOwnerV1::admit_and_attach_for_daemon( + database_owner, + &tracedecay_runtime_core::cancellation::CancellationToken::new(), + ) + .await + .expect("daemon admission"); let registered = database.issue_lease().expect("issue daemon test lease"); ( Self { diff --git a/crates/tracedecay-global-db/src/tests/schema_install_cancellation.rs b/crates/tracedecay-global-db/src/tests/schema_install_cancellation.rs new file mode 100644 index 0000000000..22eea5ae5a --- /dev/null +++ b/crates/tracedecay-global-db/src/tests/schema_install_cancellation.rs @@ -0,0 +1,167 @@ +//! Daemon shutdown cancelling a registered store attach mid schema install. + +use std::future::Future; +use std::path::PathBuf; +use std::sync::atomic::Ordering; +use std::task::Poll; + +use tempfile::TempDir; +use tracedecay_domain::errors::TraceDecayError; +use tracedecay_runtime_core::cancellation::CancellationToken; +use tracedecay_runtime_core::db::{ + DaemonDatabaseScope, Database, DatabaseAuthority, RegisteredTestRuntimeFixtureV1, + TestDatabaseRuntimeMode, TestDatabaseRuntimeScope, +}; + +use super::harness::TEST_RUNTIME_NONCE; +use crate::RegisteredGlobalDbOwnerV1; + +/// An existing profile-sessions store file with no schema objects, the shape +/// daemon attach installs the full registered schema into. +struct EmptyStore { + path: PathBuf, + _scope: DaemonDatabaseScope, + _directory: TempDir, +} + +impl EmptyStore { + fn create() -> Self { + crate::register_registered_schema_installer(); + let directory = tempfile::tempdir().unwrap(); + let profile_root = directory.path().join("profile"); + tracedecay_runtime_core::storage::PrivateStoreIo::create_dir_all(&profile_root).unwrap(); + let scope = tracedecay_runtime_core::db::enter_daemon_database_scope( + &profile_root, + TEST_RUNTIME_NONCE.fetch_add(1, Ordering::Relaxed), + "schema install cancellation", + ) + .unwrap(); + let path = tracedecay_sessions::runtime::user_sessions_db_path(&profile_root); + std::fs::create_dir_all(path.parent().unwrap()).unwrap(); + std::fs::File::create(&path).unwrap(); + Self { + path, + _scope: scope, + _directory: directory, + } + } + + async fn publish(&self) -> RegisteredTestRuntimeFixtureV1 { + let authority = + DatabaseAuthority::for_owned_runtime(&self.path, "schema install cancellation test") + .unwrap(); + Database::publish_registered_daemon_test_runtime_with_retirement_control( + &self.path, + &authority, + TestDatabaseRuntimeMode::Existing, + TestDatabaseRuntimeScope::ProfileSessions, + ) + .await + .unwrap() + } + + async fn schema_objects(&self) -> Vec { + let (owner, _runtime, _retirement) = self.publish().await.into_parts(); + let database = owner.issue_lease().unwrap(); + let mut rows = database + .read_connection() + .query( + "SELECT type || ':' || name FROM sqlite_schema + WHERE name NOT LIKE 'sqlite_%' ORDER BY type, name", + (), + ) + .await + .unwrap(); + let mut objects = Vec::new(); + while let Some(row) = rows.next().await.unwrap() { + objects.push(row.get::(0).unwrap()); + } + objects + } + + /// Runs daemon admission, cancelling its token just before the + /// `cancel_at`-th poll. Every poll after the first follows one completed + /// statement round trip, so sweeping `cancel_at` lands the cancellation at + /// each statement boundary of the install. Returns the outcome and how + /// many polls admission took. + async fn attach_cancelled_at(&self, cancel_at: usize) -> (Result<(), TraceDecayError>, usize) { + let (owner, _runtime, _retirement) = self.publish().await.into_parts(); + let cancellation = CancellationToken::new(); + let mut admission = Box::pin(RegisteredGlobalDbOwnerV1::admit_and_attach_for_daemon( + owner, + &cancellation, + )); + let mut polls = 0; + let outcome = std::future::poll_fn(|context| { + polls += 1; + if polls == cancel_at { + cancellation.cancel(); + } + match admission.as_mut().poll(context) { + Poll::Ready(result) => Poll::Ready(result.map(drop)), + Poll::Pending => Poll::Pending, + } + }) + .await; + (outcome, polls) + } +} + +#[tokio::test(flavor = "multi_thread")] +async fn cancelled_daemon_attach_leaves_the_store_empty_or_fully_installed_and_reopens() { + let reference = EmptyStore::create(); + let (outcome, total_polls) = reference.attach_cancelled_at(usize::MAX).await; + outcome.unwrap(); + let installed = reference.schema_objects().await; + assert!( + installed.contains(&"table:code_projects".to_owned()), + "a completed attach installs the registered schema: {installed:?}" + ); + + let stride = (total_polls / 20).max(1); + let mut cancelled = Vec::new(); + let mut completed = Vec::new(); + for cancel_at in (1..total_polls).step_by(stride).chain([total_polls + 10]) { + let store = EmptyStore::create(); + let (outcome, _) = store.attach_cancelled_at(cancel_at).await; + let objects = store.schema_objects().await; + match outcome { + Err(error) => { + assert!( + error.is_store_open_cancelled(), + "cancel at poll {cancel_at}/{total_polls} must fail typed: {error:?}" + ); + assert_eq!( + objects, + Vec::::new(), + "cancel at poll {cancel_at}/{total_polls} must roll the install back" + ); + cancelled.push(cancel_at); + } + Ok(()) => { + assert_eq!( + objects, installed, + "an attach that commits before observing cancel at poll {cancel_at}/{total_polls} installs everything" + ); + completed.push(cancel_at); + } + } + let (reopened, _) = store.attach_cancelled_at(usize::MAX).await; + reopened.unwrap_or_else(|error| { + panic!("reopen after cancel at poll {cancel_at}/{total_polls} failed: {error:?}") + }); + assert_eq!(store.schema_objects().await, installed); + } + // Only the commit and the idempotent post-commit tail may run past a + // cancel: every sampled point up to the commit boundary rolls back. + let last_cancelled = cancelled.iter().max().copied().unwrap_or(0); + let first_completed = completed.iter().min().copied().unwrap_or(usize::MAX); + assert!( + cancelled.contains(&1) + && last_cancelled >= total_polls / 2 + && first_completed.saturating_sub(last_cancelled) <= 2 * stride, + "cancellation must stop the install at every statement before its commit; \ + cancelled at {cancelled:?}, completed at {completed:?} of {total_polls} polls" + ); + assert!(completed.contains(&(total_polls + 10))); +} diff --git a/crates/tracedecay-runtime-core/src/ports.rs b/crates/tracedecay-runtime-core/src/ports.rs index 351476a236..c58db21d34 100644 --- a/crates/tracedecay-runtime-core/src/ports.rs +++ b/crates/tracedecay-runtime-core/src/ports.rs @@ -23,6 +23,7 @@ pub mod registered_schema { use std::pin::Pin; use std::sync::OnceLock; + use crate::cancellation::CancellationToken; use crate::db::engine::{Connection, Executor, QueryExecutor, Transaction}; use tracedecay_domain::errors::{Result, TraceDecayError}; use tracedecay_store::StoreRuntimeBindingV1; @@ -36,6 +37,7 @@ pub mod registered_schema { /// been validated for a final-schema installation. pub struct RegisteredSchemaInstallationV1 { connection: Connection, + cancellation: CancellationToken, } /// An atomic schema-installation transaction tied to its initializing @@ -48,8 +50,20 @@ pub mod registered_schema { } impl RegisteredSchemaInstallationV1 { - fn from_authorized_connection(connection: Connection) -> Self { - Self { connection } + fn from_authorized_connection( + connection: Connection, + cancellation: CancellationToken, + ) -> Self { + Self { + connection, + cancellation, + } + } + + /// Cancelled when shutdown stops the Store open this installation + /// belongs to; the installer rolls back at its next statement. + pub fn cancellation(&self) -> &CancellationToken { + &self.cancellation } /// The exact Store scope being initialized. @@ -290,8 +304,12 @@ pub mod registered_schema { /// This is crate-private so no dependent crate can fabricate an /// installation capability before Store publication. #[hotpath::skip] - pub(crate) async fn install_from_authorized_connection(connection: Connection) -> Result<()> { - let installation = RegisteredSchemaInstallationV1::from_authorized_connection(connection); + pub(crate) async fn install_from_authorized_connection( + connection: Connection, + cancellation: CancellationToken, + ) -> Result<()> { + let installation = + RegisteredSchemaInstallationV1::from_authorized_connection(connection, cancellation); ensure_registered_schema(&installation).await } @@ -368,8 +386,10 @@ pub mod registered_schema { &directory.path().join("registered-schema.sqlite3"), authority.clone(), ); - let installation = - RegisteredSchemaInstallationV1::from_authorized_connection((*connection).clone()); + let installation = RegisteredSchemaInstallationV1::from_authorized_connection( + (*connection).clone(), + CancellationToken::new(), + ); let transaction = installation .begin_atomic_schema_transaction() .await diff --git a/crates/tracedecay-runtime-core/src/shard_runtime/registry.rs b/crates/tracedecay-runtime-core/src/shard_runtime/registry.rs index befd753478..e8b54b525c 100644 --- a/crates/tracedecay-runtime-core/src/shard_runtime/registry.rs +++ b/crates/tracedecay-runtime-core/src/shard_runtime/registry.rs @@ -38,6 +38,7 @@ use tracedecay_store::{ use super::shard::ShardRuntime; use super::telemetry::{RuntimeRegistryInventory, RuntimeRegistryInventoryEntry}; use super::utc_now; +use crate::cancellation::CancellationToken; use crate::profiled_lock::{ProfiledMutex, ProfiledMutexGuard}; #[cfg(test)] @@ -1310,6 +1311,11 @@ pub enum StoreRuntimeRegistryFailure { OpenTaskAbandoned { key: Box, }, + /// Shutdown cancelled this open at a safe point. Nothing it opened is + /// published or left open. + OpenCancelled { + key: Box, + }, } #[derive(Clone, Debug)] @@ -1453,6 +1459,7 @@ struct StoreRuntimeRegistryInner { /// reservation funnels through `lock_state`, so it is the coarsest lock in /// the store runtime and the first place cross-shard queueing shows up. state: ProfiledMutex, + open_cancellation: CancellationToken, } #[derive(Clone)] @@ -1474,6 +1481,7 @@ impl StoreRuntimeRegistry { Mutex::new(RegistryState::default()), label = "runtime_core.shard_runtime.registry_state" ), + open_cancellation: CancellationToken::new(), }), } } @@ -1505,10 +1513,25 @@ impl StoreRuntimeRegistry { Mutex::new(RegistryState::default()), label = "runtime_core.shard_runtime.registry_state" ), + open_cancellation: CancellationToken::new(), }), }) } + /// Refuses every new runtime open and stops in-flight opens and schema + /// installs at their next safe point, each with a typed + /// [`StoreRuntimeRegistryFailure::OpenCancelled`]. Shutdown-only: the + /// registry never admits opens again. + pub fn cancel_opens_for_shutdown(&self) { + self.inner.open_cancellation.cancel(); + } + + /// The token opens and the schema installs that follow them observe. + #[must_use] + pub fn open_cancellation(&self) -> &CancellationToken { + &self.inner.open_cancellation + } + pub fn lookup(&self, expected: &StoreRuntimeBindingV1) -> StoreRuntimeLookup { let key = StoreRuntimeKey::from_binding(expected); let state = self.lock_state(); 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 8d22e06052..81e4e108f0 100644 --- a/crates/tracedecay-runtime-core/src/shard_runtime/registry/close.rs +++ b/crates/tracedecay-runtime-core/src/shard_runtime/registry/close.rs @@ -52,11 +52,13 @@ impl StoreRuntimeRegistry { /// Closes every mounted runtime that no lease, queued work, profile pin, /// or graph lease still holds, so each writer runs its shutdown TRUNCATE /// checkpoint. The daemon process exits with this registry reachable, so - /// no destructor closes these attachments otherwise. Held runtimes stay - /// mounted and are logged with their blockers. Returns the number of - /// runtimes closed. + /// 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. #[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; let (reservations, reserve_failure) = { let mut state = self.lock_state(); let mut idle = Vec::new(); diff --git a/crates/tracedecay-runtime-core/src/shard_runtime/registry/open.rs b/crates/tracedecay-runtime-core/src/shard_runtime/registry/open.rs index 0f48bd1758..79b60fe1aa 100644 --- a/crates/tracedecay-runtime-core/src/shard_runtime/registry/open.rs +++ b/crates/tracedecay-runtime-core/src/shard_runtime/registry/open.rs @@ -252,6 +252,11 @@ impl StoreRuntimeRegistry { }; } + if self.inner.open_cancellation.is_cancelled() { + return StoreRuntimeOpenBegin::Rejected( + StoreRuntimeRegistryFailure::OpenCancelled { key: Box::new(key) }, + ); + } let authority_epoch = match state.graph_publications.get(&key) { Some(graph) => graph.binding.authority_epoch, None => match allocate_authority_epoch() { @@ -360,6 +365,12 @@ impl StoreRuntimeRegistry { } (outcome, _) => outcome, }; + let outcome = match outcome { + Ok((published, ..)) if registry.inner.open_cancellation.is_cancelled() => { + Err(close_cancelled_open(&key, published).await) + } + outcome => outcome, + }; guard.complete(outcome); }); StoreRuntimeOpenBegin::Started(join) @@ -396,6 +407,32 @@ impl StoreRuntimeRegistry { } } + /// Cancels every open and waits until each in-flight attempt has settled, + /// so none can publish or hold a physical handle afterwards. + pub(super) async fn cancel_and_join_opens_for_shutdown(&self) { + self.cancel_opens_for_shutdown(); + let openings = self + .lock_state() + .entries + .values() + .filter_map(|entry| match entry { + RegistryEntry::Opening(opening) => Some(opening.updates.subscribe()), + _ => None, + }) + .collect::>(); + for mut updates in openings { + // A closed channel means the attempt's guard already settled it. + loop { + if !matches!(*updates.borrow_and_update(), OpenState::Opening) { + break; + } + if updates.changed().await.is_err() { + break; + } + } + } + } + fn fail_reserved_open( &self, key: &StoreRuntimeKey, @@ -427,6 +464,7 @@ impl StoreRuntimeRegistry { Result, > { Box::pin(async move { + self.open_checkpoint(key)?; let resolved = self .inner .resolver @@ -457,6 +495,7 @@ impl StoreRuntimeRegistry { } } let locator = RuntimeLocatorRecord::new(key.clone(), resolved); + self.open_checkpoint(key)?; let published = self .inner .publisher @@ -466,6 +505,7 @@ impl StoreRuntimeRegistry { mode, access, database_authority.clone(), + self.inner.open_cancellation.clone(), )) .await?; if published.binding() != &binding { @@ -477,6 +517,36 @@ impl StoreRuntimeRegistry { Ok((published, locator, database_authority)) }) } + + fn open_checkpoint(&self, key: &StoreRuntimeKey) -> Result<(), StoreRuntimeRegistryFailure> { + if self.inner.open_cancellation.is_cancelled() { + return Err(StoreRuntimeRegistryFailure::OpenCancelled { + key: Box::new(key.clone()), + }); + } + Ok(()) + } +} + +/// Closes a runtime that finished opening after shutdown cancelled it, so the +/// open ends with nothing published and no physical handle left behind. +async fn close_cancelled_open( + key: &StoreRuntimeKey, + published: PublishedShardRuntime, +) -> StoreRuntimeRegistryFailure { + let close = tokio::task::spawn_blocking(move || close_unpublished_runtime(published)) + .await + .map_err(|error| error.to_string()) + .and_then(|result| result); + match close { + Ok(()) => StoreRuntimeRegistryFailure::OpenCancelled { + key: Box::new(key.clone()), + }, + Err(message) => StoreRuntimeRegistryFailure::PhysicalRuntimeFailed { + operation: "close cancelled registered runtime open", + message, + }, + } } fn open_access_compatible(request: &StoreRuntimeOpenRequest, opening: &OpeningRuntime) -> bool { diff --git a/crates/tracedecay-runtime-core/src/shard_runtime/registry/ports.rs b/crates/tracedecay-runtime-core/src/shard_runtime/registry/ports.rs index 52bd7e37da..09e6a50948 100644 --- a/crates/tracedecay-runtime-core/src/shard_runtime/registry/ports.rs +++ b/crates/tracedecay-runtime-core/src/shard_runtime/registry/ports.rs @@ -19,6 +19,7 @@ use super::{ PublishedShardRuntime, StoreRuntimeAccessMode, StoreRuntimeKey, StoreRuntimeOpenMode, StoreRuntimeRegistryFailure, }; +use crate::cancellation::CancellationToken; use crate::shard_runtime::shard::{ShardRuntime, ShardRuntimeError}; #[derive(Clone, Debug, PartialEq, Eq)] @@ -416,12 +417,23 @@ async fn install_final_schema_before_publication( StoreShardScopeV1::Profile | StoreShardScopeV1::ProfileSessions | StoreShardScopeV1::ProjectSessions { .. } => { - crate::ports::registered_schema::install_from_authorized_connection(connection) - .await - .map_err(|error| StoreRuntimeRegistryFailure::PhysicalRuntimeFailed { - operation: "create initialized global/session schema", - message: error.to_string(), - })?; + crate::ports::registered_schema::install_from_authorized_connection( + connection, + request.open_cancellation.clone(), + ) + .await + .map_err(|error| { + if error.is_store_open_cancelled() { + StoreRuntimeRegistryFailure::OpenCancelled { + key: Box::new(StoreRuntimeKey::from_binding(&request.binding)), + } + } else { + StoreRuntimeRegistryFailure::PhysicalRuntimeFailed { + operation: "create initialized global/session schema", + message: error.to_string(), + } + } + })?; } StoreShardScopeV1::RemoteNode { .. } => { connection @@ -598,6 +610,7 @@ pub struct ShardRuntimeBuildRequest { mode: StoreRuntimeOpenMode, access: StoreRuntimeAccessMode, database_authority: Option, + open_cancellation: CancellationToken, } impl ShardRuntimeBuildRequest { @@ -607,6 +620,7 @@ impl ShardRuntimeBuildRequest { mode: StoreRuntimeOpenMode, access: StoreRuntimeAccessMode, database_authority: Option, + open_cancellation: CancellationToken, ) -> Self { Self { binding, @@ -614,6 +628,7 @@ impl ShardRuntimeBuildRequest { mode, access, database_authority, + open_cancellation, } } diff --git a/crates/tracedecay-runtime-core/src/shard_runtime/registry/tests.rs b/crates/tracedecay-runtime-core/src/shard_runtime/registry/tests.rs index fe2e4a7c3a..3dc10bed13 100644 --- a/crates/tracedecay-runtime-core/src/shard_runtime/registry/tests.rs +++ b/crates/tracedecay-runtime-core/src/shard_runtime/registry/tests.rs @@ -245,6 +245,69 @@ fn cancelled_open_task_wakes_every_joiner_and_allows_retry() { } } +/// Shutdown during an in-flight open: the shutdown close cancels and joins the +/// open instead of skipping it, the open ends typed-cancelled with nothing +/// published, and every idle published runtime still closes. +#[tokio::test] +async fn shutdown_close_cancels_and_joins_an_in_flight_open() { + let (registry, _, publisher) = registry(StoreRuntimeRegistryConfig::default()); + let pin = profile_pin(®istry).await; + publisher.block.store(true, Ordering::SeqCst); + let request = project_sessions_request("project.shutdown", &pin); + let opening = registry.begin_or_join_open(&request); + wait_for_calls(&publisher.calls, 2).await; + let opening_key = request.key.clone(); + drop((request, pin)); + + let close = tokio::spawn({ + let registry = registry.clone(); + async move { registry.close_idle_for_shutdown().await } + }); + tokio::time::timeout(Duration::from_secs(2), async { + while !registry.open_cancellation().is_cancelled() { + tokio::task::yield_now().await; + } + }) + .await + .expect("shutdown close cancels in-flight opens"); + tokio::task::yield_now().await; + assert!( + !close.is_finished(), + "shutdown close must join the open it cancelled, not skip it" + ); + publisher.release.notify_one(); + + let opened = opening.wait().await; + assert!( + matches!( + &opened, + StoreRuntimeOpenResult::Failed(StoreRuntimeRegistryFailure::OpenCancelled { key }) + if **key == opening_key + ), + "a cancelled open publishes nothing: {opened:?}" + ); + assert_eq!(close.await.unwrap().unwrap(), 1, "only the profile runtime"); + let inventory = registry.inventory(AdmissionConfigV1::default(), None); + assert_eq!(inventory.opening_shards, 0); + assert_eq!(inventory.entries.len(), 0); + + let reopened = registry + .open(StoreRuntimeOpenRequest::new( + profile_shard(), + incarnation(), + None, + )) + .await; + assert!( + matches!( + &reopened, + StoreRuntimeOpenResult::Failed(StoreRuntimeRegistryFailure::OpenCancelled { .. }) + ), + "a shut-down registry admits no new open: {reopened:?}" + ); + assert_eq!(publisher.calls.load(Ordering::SeqCst), 2); +} + #[tokio::test] async fn profile_pin_budget_and_all_runtime_blockers_are_authoritative() { let config = StoreRuntimeRegistryConfig::new(2).unwrap(); diff --git a/crates/tracedecay-store-runtime/src/session_registry.rs b/crates/tracedecay-store-runtime/src/session_registry.rs index da65652a36..d89f7e7a79 100644 --- a/crates/tracedecay-store-runtime/src/session_registry.rs +++ b/crates/tracedecay-store-runtime/src/session_registry.rs @@ -2772,6 +2772,9 @@ pub fn registry_open_error( StoreRuntimeRegistryFailure::ResetRequired { authority, reason } => { TraceDecayError::reset_required(authority, reason) } + StoreRuntimeRegistryFailure::OpenCancelled { .. } => { + TraceDecayError::store_open_cancelled(operation) + } failure => session_registry_error(operation, format!("{failure:?}")), } } diff --git a/crates/tracedecay-store-runtime/src/session_registry/maintenance.rs b/crates/tracedecay-store-runtime/src/session_registry/maintenance.rs index 2340b8285b..22bc1dc0e9 100644 --- a/crates/tracedecay-store-runtime/src/session_registry/maintenance.rs +++ b/crates/tracedecay-store-runtime/src/session_registry/maintenance.rs @@ -640,7 +640,11 @@ impl DaemonSessionRuntimeRegistryV1 { let long_lived = self.long_lived_session_maintenance; let (database, convergence) = if long_lived { let (database, convergence) = - RegisteredGlobalDbOwnerV1::admit_and_attach_for_daemon(database).await?; + RegisteredGlobalDbOwnerV1::admit_and_attach_for_daemon( + database, + self.registry.open_cancellation(), + ) + .await?; (database, Some(convergence)) } else { ( diff --git a/crates/tracedecay-store-runtime/src/session_registry/mounts.rs b/crates/tracedecay-store-runtime/src/session_registry/mounts.rs index ee45f078f3..307513ce6d 100644 --- a/crates/tracedecay-store-runtime/src/session_registry/mounts.rs +++ b/crates/tracedecay-store-runtime/src/session_registry/mounts.rs @@ -991,6 +991,13 @@ impl DaemonSessionRuntimeRegistryV1 { databases } + /// Stops every in-flight store open and schema install at its next safe + /// point and refuses new ones, so a draining daemon never waits for a + /// mount to run to completion. + pub fn cancel_store_opens_for_shutdown(&self) { + self.registry.cancel_opens_for_shutdown(); + } + /// Releases exclusive Grafeo writers at daemon shutdown, after graph /// operation settlement and reconciliation workers have joined. /// diff --git a/crates/tracedecay/src/daemon/branch_admin/session_runtime_shutdown.rs b/crates/tracedecay/src/daemon/branch_admin/session_runtime_shutdown.rs index f7c8b4a553..0fdf4e36bd 100644 --- a/crates/tracedecay/src/daemon/branch_admin/session_runtime_shutdown.rs +++ b/crates/tracedecay/src/daemon/branch_admin/session_runtime_shutdown.rs @@ -157,6 +157,17 @@ impl StoreAdministration { } } + /// Shutdown-only: stops the store opens and schema installs of every + /// mounted session runtime registry at their next safe point, so an + /// admitted project open returns promptly instead of finishing its mount. + pub(in crate::daemon) async fn cancel_store_opens_for_shutdown(&self) { + for entry in self.session_runtime_registries.lock().await.values() { + if let Some(registry) = entry.registry.get() { + registry.cancel_store_opens_for_shutdown(); + } + } + } + #[hotpath::skip] pub(in crate::daemon) async fn prepare_memory_graph_reconciliation_shutdown( &self, diff --git a/crates/tracedecay/src/daemon/engine/shutdown.rs b/crates/tracedecay/src/daemon/engine/shutdown.rs index e66f2dc527..7523af8946 100644 --- a/crates/tracedecay/src/daemon/engine/shutdown.rs +++ b/crates/tracedecay/src/daemon/engine/shutdown.rs @@ -42,6 +42,7 @@ impl DaemonEngine { )] pub(in crate::daemon) async fn shutdown_owner_phases(&self) -> Vec> { let project_open = project_open_tasks(&self.project_open_gates).await; + let store_open_cancel = self.store_administration.clone(); let manual_branch_cancel = self.store_administration.clone(); let manual_branch_join = self.store_administration.clone(); @@ -88,7 +89,9 @@ impl DaemonEngine { // An admitted open registers its owners with the invocation // registry, so it must settle before that registry drains: it // either registers in time to be released or stops at a - // cancellation boundary before registering anything. + // cancellation boundary before registering anything. The store + // mount it may be inside stops at its next safe point too, so the + // open never outlasts the cooperative window by finishing a mount. vec![ShutdownOwner::with_deadline_status( "project_open", { @@ -96,6 +99,7 @@ impl DaemonEngine { move || project_open_cancel.cancel_all() }, move |_| async move { + store_open_cancel.cancel_store_opens_for_shutdown().await; if project_open.shutdown().await { ShutdownStatus::Clean } else { From 6b180f95cb4a9bae606d5c83c5a6486630e06f90 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 29 Sep 2026 12:17:39 +0000 Subject: [PATCH 2/3] simplify(global-db): drop the redundant batched cancellation override --- crates/tracedecay-global-db/src/schema_stages.rs | 10 +--------- 1 file changed, 1 insertion(+), 9 deletions(-) diff --git a/crates/tracedecay-global-db/src/schema_stages.rs b/crates/tracedecay-global-db/src/schema_stages.rs index 0bc5b5eb76..bc19ecfaac 100644 --- a/crates/tracedecay-global-db/src/schema_stages.rs +++ b/crates/tracedecay-global-db/src/schema_stages.rs @@ -16,7 +16,7 @@ use tracedecay_runtime_core::{ cancellation::CancellationToken, db::{ Database, DatabaseWriteTransaction, - engine::{self, Executor, IntoParams, QueryExecutor, Rows, WriteStatement}, + engine::{self, Executor, IntoParams, QueryExecutor, Rows}, }, ports::registered_schema::{ RegisteredSchemaInstallationTransactionV1, RegisteredSchemaInstallationV1, @@ -712,14 +712,6 @@ impl Executor for CancellableSchemaExecutor<'_, T> { self.inner.execute(sql, params).await } - async fn execute_statements( - &self, - statements: Vec, - ) -> engine::Result> { - self.checkpoint()?; - self.inner.execute_statements(statements).await - } - async fn execute_batch(&self, sql: &str) -> engine::Result<()> { self.checkpoint()?; self.inner.execute_batch(sql).await From f17d16e7ae4fff6b3d54a2fe713ee08406b3209f Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 29 Sep 2026 12:36:59 +0000 Subject: [PATCH 3/3] refactor(global-db): let the cancellable install own its transaction --- .../tracedecay-global-db/src/schema_stages.rs | 49 ++++++++++--------- 1 file changed, 27 insertions(+), 22 deletions(-) diff --git a/crates/tracedecay-global-db/src/schema_stages.rs b/crates/tracedecay-global-db/src/schema_stages.rs index bc19ecfaac..dbd9cb4656 100644 --- a/crates/tracedecay-global-db/src/schema_stages.rs +++ b/crates/tracedecay-global-db/src/schema_stages.rs @@ -553,12 +553,14 @@ pub async fn ensure_registered_schema_for_admission( .map_err(|error| global_db_operation_error(OPERATION, error))?; if let Some(refused) = install_and_commit_registered_schema( - transaction, + CancellableSchemaTransaction { + transaction, + cancellation: installation.cancellation(), + }, configuration_fresh.as_ref(), temporal_admission, workflow_admission, force_exhaustive, - installation.cancellation(), "commit registered global schema", "roll back registered global schema", ) @@ -616,12 +618,11 @@ async fn install_registered_schema_stages( } async fn install_and_commit_registered_schema( - transaction: T, + install: CancellableSchemaTransaction<'_, T>, configuration_fresh: Option<&configuration::FreshConfigurationStoreEvidence>, temporal_admission: session_temporal_schema::SessionTemporalSchemaAdmission, workflow_admission: WorkflowSchemaAdmission, force_exhaustive: bool, - cancellation: &CancellationToken, commit_operation: &'static str, rollback_operation: &'static str, ) -> tracedecay_domain::errors::Result> @@ -629,10 +630,7 @@ where T: Executor + Sync + SchemaInstallTransaction, { let admission = install_registered_schema_stages( - &CancellableSchemaExecutor { - inner: &transaction, - cancellation, - }, + &install, configuration_fresh, temporal_admission, workflow_admission, @@ -640,8 +638,13 @@ where ) .await .and_then(|refused_authority| { - schema_install_checkpoint(cancellation, commit_operation).map(|()| refused_authority) + schema_install_checkpoint(install.cancellation, commit_operation) + .map(|()| refused_authority) }); + let CancellableSchemaTransaction { + transaction, + cancellation, + } = install; match admission { Ok(refused_authority) => { transaction.commit().await.map_err(|error| { @@ -674,15 +677,15 @@ fn schema_install_checkpoint( Ok(()) } -/// Refuses the next statement once shutdown cancels the Store open, so the -/// enclosing schema transaction rolls back at a statement boundary instead of -/// running the whole install first. -struct CancellableSchemaExecutor<'a, T> { - inner: &'a T, +/// A schema transaction that refuses its next statement once shutdown cancels +/// the Store open, so it rolls back at a statement boundary instead of running +/// the whole install first. +struct CancellableSchemaTransaction<'a, T> { + transaction: T, cancellation: &'a CancellationToken, } -impl CancellableSchemaExecutor<'_, T> { +impl CancellableSchemaTransaction<'_, T> { fn checkpoint(&self) -> engine::Result<()> { if self.cancellation.is_cancelled() { return Err(engine::Error::InvalidOperation( @@ -693,28 +696,28 @@ impl CancellableSchemaExecutor<'_, T> { } } -impl QueryExecutor for CancellableSchemaExecutor<'_, T> { +impl QueryExecutor for CancellableSchemaTransaction<'_, T> { async fn query

(&self, sql: &str, params: P) -> engine::Result where P: IntoParams, { self.checkpoint()?; - self.inner.query(sql, params).await + self.transaction.query(sql, params).await } } -impl Executor for CancellableSchemaExecutor<'_, T> { +impl Executor for CancellableSchemaTransaction<'_, T> { async fn execute

(&self, sql: &str, params: P) -> engine::Result where P: IntoParams, { self.checkpoint()?; - self.inner.execute(sql, params).await + self.transaction.execute(sql, params).await } async fn execute_batch(&self, sql: &str) -> engine::Result<()> { self.checkpoint()?; - self.inner.execute_batch(sql).await + self.transaction.execute_batch(sql).await } } @@ -1096,12 +1099,14 @@ pub(crate) async fn ensure_attached_registered_schema( schema_install_checkpoint(cancellation, OPERATION)?; let transaction = database.begin_bulk_write_transaction(OPERATION).await?; let refused_authority = install_and_commit_registered_schema( - transaction, + CancellableSchemaTransaction { + transaction, + cancellation, + }, configuration_fresh.as_ref(), temporal_admission, workflow_admission, force_exhaustive, - cancellation, "commit attached registered global schema", "roll back attached registered global schema", )