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
17 changes: 17 additions & 0 deletions crates/tracedecay-domain/src/errors.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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<String>,
retryable: bool,
Expand Down
38 changes: 24 additions & 14 deletions crates/tracedecay-global-db/src/registered.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -119,17 +120,20 @@ impl RegisteredGlobalDbOwnerV1 {
) -> tracedecay_domain::errors::Result<Self> {
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(&registered.database)
.await?
{
RegisteredSchemaAttachmentV1::Admitted(_) => {
super::schema_stages::converge_attached_registered_schema(&registered.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(
&registered.database,
&CancellationToken::new(),
)
.await?
{
RegisteredSchemaAttachmentV1::Admitted(_) => {
super::schema_stages::converge_attached_registered_schema(&registered.database)
.await?;
None
}
RegisteredSchemaAttachmentV1::SessionsRefused(refused) => Some(refused),
};
drop(registered);
Ok(Self {
database,
Expand All @@ -141,19 +145,25 @@ 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<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) =
match super::schema_stages::ensure_attached_registered_schema(&registered.database)
.await?
match super::schema_stages::ensure_attached_registered_schema(
&registered.database,
cancellation,
)
.await?
{
RegisteredSchemaAttachmentV1::Admitted(convergence) => (Some(convergence), None),
RegisteredSchemaAttachmentV1::SessionsRefused(refused) => (None, Some(refused)),
Expand Down
101 changes: 92 additions & 9 deletions crates/tracedecay-global-db/src/schema_stages.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -659,7 +664,7 @@ async fn install_registered_schema_stages(
}

async fn install_and_commit_registered_schema<T>(
transaction: T,
install: CancellableSchemaTransaction<'_, T>,
configuration_fresh: Option<&configuration::FreshConfigurationStoreEvidence>,
temporal_admission: session_temporal_schema::SessionTemporalSchemaAdmission,
workflow_admission: WorkflowSchemaAdmission,
Expand All @@ -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| {
Expand All @@ -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,
Expand All @@ -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<T> 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<T: QueryExecutor + Sync> QueryExecutor for CancellableSchemaTransaction<'_, T> {
async fn query<P>(&self, sql: &str, params: P) -> engine::Result<Rows>
where
P: IntoParams,
{
self.checkpoint()?;
self.transaction.query(sql, params).await
}
}

impl<T: Executor + Sync> Executor for CancellableSchemaTransaction<'_, T> {
async fn execute<P>(&self, sql: &str, params: P) -> engine::Result<u64>
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<Output = Result<(), String>> + Send;
fn rollback(self) -> impl std::future::Future<Output = Result<(), String>> + Send;
Expand Down Expand Up @@ -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<RegisteredSchemaAttachmentV1> {
const OPERATION: &str = "install attached registered global database schema";
schema_install_checkpoint(cancellation, OPERATION)?;
let read_connection = database.read_connection();
let RegisteredSchemaAdmissionClassification {
configuration_fresh,
Expand All @@ -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,
Expand Down
2 changes: 2 additions & 0 deletions crates/tracedecay-global-db/src/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand Down
12 changes: 7 additions & 5 deletions crates/tracedecay-global-db/src/tests/harness.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<tracedecay_runtime_core::resident_memory::ProcessResidentMemoryV1>,
Expand Down Expand Up @@ -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 {
Expand Down
Loading
Loading