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 292e8e032b..ffc5102959 100644 --- a/crates/tracedecay-global-db/src/registered.rs +++ b/crates/tracedecay-global-db/src/registered.rs @@ -5,6 +5,7 @@ use std::sync::{Arc, OnceLock, RwLock, Weak}; use crate::schema_stages::RegisteredSchemaAttachmentV1; use tracedecay_domain::errors::TraceDecayError; use tracedecay_runtime_core::{ + cancellation::CancellationToken, db::{ Database, DatabaseAuthority, DatabaseEngineReadConnection, DatabaseEngineReadSnapshot, DatabaseOwnerErrorV1, DatabaseOwnerRetirementReservationV1, DatabaseOwnerV1, @@ -119,17 +120,20 @@ 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 = - match super::schema_stages::ensure_attached_registered_schema(®istered.database) - .await? - { - RegisteredSchemaAttachmentV1::Admitted(_) => { - super::schema_stages::converge_attached_registered_schema(®istered.database) - .await?; - None - } - RegisteredSchemaAttachmentV1::SessionsRefused(refused) => Some(refused), - }; + // Short-lived attaches have no daemon shutdown to observe. + let refused_authority = match super::schema_stages::ensure_attached_registered_schema( + ®istered.database, + &CancellationToken::new(), + ) + .await? + { + RegisteredSchemaAttachmentV1::Admitted(_) => { + super::schema_stages::converge_attached_registered_schema(®istered.database) + .await?; + None + } + RegisteredSchemaAttachmentV1::SessionsRefused(refused) => Some(refused), + }; drop(registered); Ok(Self { database, @@ -141,10 +145,13 @@ impl RegisteredGlobalDbOwnerV1 { /// Returns the resumable convergence plan for an already admitted schema /// without retaining an unowned client lease. A store admitted in its - /// typed reset-required state has no plan: its reset deletes it. + /// typed reset-required state has no plan: its reset deletes it. 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, Option, @@ -152,8 +159,11 @@ impl RegisteredGlobalDbOwnerV1 { let temporary = database.issue_lease().map_err(registered_owner_error)?; let registered = RegisteredGlobalDb::from_owned_database(temporary); let (convergence, refused_authority) = - match super::schema_stages::ensure_attached_registered_schema(®istered.database) - .await? + match super::schema_stages::ensure_attached_registered_schema( + ®istered.database, + cancellation, + ) + .await? { RegisteredSchemaAttachmentV1::Admitted(convergence) => (Some(convergence), None), RegisteredSchemaAttachmentV1::SessionsRefused(refused) => (None, Some(refused)), diff --git a/crates/tracedecay-global-db/src/schema_stages.rs b/crates/tracedecay-global-db/src/schema_stages.rs index 908a26b684..56fc840af2 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}, }, ports::registered_schema::{ RegisteredSchemaInstallationTransactionV1, RegisteredSchemaInstallationV1, @@ -591,13 +592,17 @@ pub async fn ensure_registered_schema_for_admission( }; 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 .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, @@ -659,7 +664,7 @@ 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, @@ -671,13 +676,21 @@ where T: Executor + Sync + SchemaInstallTransaction, { let admission = install_registered_schema_stages( - &transaction, + &install, configuration_fresh, temporal_admission, workflow_admission, force_exhaustive, ) - .await; + .await + .and_then(|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| { @@ -686,6 +699,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, @@ -695,6 +713,60 @@ 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(()) +} + +/// 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 CancellableSchemaTransaction<'_, 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 CancellableSchemaTransaction<'_, T> { + async fn query

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

(&self, sql: &str, params: P) -> engine::Result + where + P: IntoParams, + { + self.checkpoint()?; + self.transaction.execute(sql, params).await + } + + async fn execute_batch(&self, sql: &str) -> engine::Result<()> { + self.checkpoint()?; + self.transaction.execute_batch(sql).await + } +} + trait SchemaInstallTransaction: Sized { fn commit(self) -> impl std::future::Future> + Send; fn rollback(self) -> impl std::future::Future> + Send; @@ -1047,10 +1119,19 @@ pub async fn converge_attached_registered_schema( /// is admitted untouched for its other authorities; one whose observation /// rows it refuses is installed without them. Either returns the refused /// authority instead of a 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 { + const OPERATION: &str = "install attached registered global database schema"; + schema_install_checkpoint(cancellation, OPERATION)?; let read_connection = database.read_connection(); let RegisteredSchemaAdmissionClassification { configuration_fresh, @@ -1063,11 +1144,13 @@ pub(crate) async fn ensure_attached_registered_schema( } }; 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, + CancellableSchemaTransaction { + transaction, + cancellation, + }, configuration_fresh.as_ref(), temporal_admission, workflow_admission, 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 74baddb5c8..2a803c10a2 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 90f6bdd794..43404709bd 100644 --- a/crates/tracedecay-store-runtime/src/session_registry/maintenance.rs +++ b/crates/tracedecay-store-runtime/src/session_registry/maintenance.rs @@ -639,7 +639,11 @@ impl DaemonSessionRuntimeRegistryV1 { // reset deletes it. let long_lived = self.long_lived_session_maintenance; let (database, convergence) = if long_lived { - RegisteredGlobalDbOwnerV1::admit_and_attach_for_daemon(database).await? + RegisteredGlobalDbOwnerV1::admit_and_attach_for_daemon( + database, + self.registry.open_cancellation(), + ) + .await? } else { ( RegisteredGlobalDbOwnerV1::admit_and_attach(database).await?, 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 {