use std::collections::BTreeSet;
use std::path::PathBuf;
use std::time::Instant;
use sha2::{Digest, Sha256};
use time::OffsetDateTime;
use crate::__bypass::guarded_batch_execute;
use crate::config::{MigrateConfig, PolicyConfig};
use crate::context::{DjogiContext, PinnedCtx};
use crate::error::{DbError, DjogiError};
use crate::types::HeerId;
use super::guard::WorkspaceGuard;
use super::ledger::{
self, ChecksumFormatError, ChecksumMismatch, ExecutionMode, LedgerRow, LedgerStatus,
VerifyError, compute_checksum, load_full_row_by_version,
};
use super::projection::BucketKey;
use super::schema::SNAPSHOT_FORMAT_VERSION;
use super::segment::{MigrationPlan, Segment, SegmentKind};
use super::snapshot::{SnapshotError, save_snapshot};
use super::sql::{LossyRollbackKind, OperationSql};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SegmentSqlExecutionModeProblem {
TransactionControl {
keyword: &'static str,
},
RequiresNonTransactional {
statement_shape: &'static str,
},
}
#[derive(Debug)]
#[non_exhaustive]
pub enum RunnerError {
LockTimeout {
path: PathBuf,
holder_pid: Option<i32>,
},
GuardError(super::guard::GuardError),
AdvisoryLockFailed {
bucket: BucketKey,
key: i64,
attempts: u32,
},
AdvisoryLockQueryFailed {
app_label: String,
source: DjogiError,
},
ChecksumMismatch(ChecksumMismatch),
ChecksumFormat(ChecksumFormatError),
LedgerWriteFailed { version: String, source: DjogiError },
LedgerQueryFailed {
query_label: &'static str,
source: DjogiError,
},
RunIdGenerationFailed { source: DjogiError },
NodeIdentityBindingFailed {
node_id: i32,
source: DjogiError,
},
SingleNodeDevProvisioningFailed {
node_id: i32,
step: &'static str,
source: DjogiError,
},
LedgerBootstrapFailed { source: DjogiError },
VersionAlreadyApplied {
version: String,
applied_at: Option<OffsetDateTime>,
},
VersionCollisionNonTerminal {
version: String,
status: LedgerStatus,
run_id: i64,
},
TargetTableNotFound {
bucket: BucketKey,
index_name: String,
target_table: String,
},
RelpagesThresholdExceeded {
bucket: BucketKey,
index_name: String,
target_table: String,
relpages: i32,
threshold: u32,
},
SegmentSqlExecutionModeConflict {
segment_index: usize,
segment_kind: SegmentKind,
statement_label: String,
problem: SegmentSqlExecutionModeProblem,
},
TransactionalSegmentFailed {
segment_index: usize,
statement_label: String,
source: DjogiError,
},
NonTransactionalSegmentFailed {
segment_index: usize,
step_index: usize,
statement_label: String,
applied_steps_count: i32,
source: DjogiError,
},
NonTransactionalProgressAckFailed {
segment_index: usize,
step_index: usize,
statement_label: String,
applied_steps_count: i32,
source: DjogiError,
},
ConfigLoadFailed { source: figment::Error },
SnapshotPersistFailed {
path: PathBuf,
source: SnapshotError,
},
BaselineSnapshotShouldNotBeProvided,
StalePhaseZeroArtifact {
version: String,
refusal_reason: &'static str,
},
OutOfOrderRejected {
version: String,
conflicting_version: String,
conflicting_applied_at: Option<String>,
},
BaselineProjectionFailed {
source: Box<super::verify::VerifyRunError>,
},
DriftDetected {
bucket: BucketKey,
report: super::verify::VerifyReport,
},
DriftBaselineMissing { bucket: BucketKey },
DriftBaselineCorrupted { bucket: BucketKey, reason: String },
DriftPreflightFailed {
source: Box<super::verify::VerifyRunError>,
},
PkFlipHazardReplicaSessions {
walsenders: Vec<(String, String)>,
subscriptions: Vec<String>,
},
PkFlipHazardPreexistingZzzTrigger {
table: String,
trigger_names: Vec<String>,
},
PkFlipHazardDisabledTriggers {
table: String,
triggers: Vec<(String, char)>,
},
PkFlipHazardLongRunningTx {
offenders: Vec<(i32, i64)>,
threshold_secs: u32,
},
PkFlipVerificationFailed {
table: String,
count_violating: i64,
},
CatalogQueryFailed {
query_label: &'static str,
source: DjogiError,
},
PinnedSessionCheckoutFailed {
source: DjogiError,
},
AdvisoryUnlockReturnedFalse {
key: i64,
bucket: BucketKey,
},
PartitionExpansionNoLeaves {
parent: String,
statement_label: String,
},
MissingApplyIdentity {
version: String,
},
}
impl RunnerError {
#[must_use]
#[allow(clippy::match_like_matches_macro)]
pub fn is_operator_actionable(&self) -> bool {
match self {
Self::ChecksumMismatch(_)
| Self::ChecksumFormat(_)
| Self::VersionAlreadyApplied { .. }
| Self::VersionCollisionNonTerminal { .. }
| Self::TargetTableNotFound { .. }
| Self::RelpagesThresholdExceeded { .. }
| Self::SegmentSqlExecutionModeConflict { .. }
| Self::SnapshotPersistFailed { .. }
| Self::BaselineSnapshotShouldNotBeProvided
| Self::StalePhaseZeroArtifact { .. }
| Self::DriftDetected { .. }
| Self::DriftBaselineMissing { .. }
| Self::DriftBaselineCorrupted { .. }
| Self::OutOfOrderRejected { .. }
| Self::PkFlipHazardReplicaSessions { .. }
| Self::PkFlipHazardPreexistingZzzTrigger { .. }
| Self::PkFlipHazardDisabledTriggers { .. }
| Self::PkFlipHazardLongRunningTx { .. }
| Self::PkFlipVerificationFailed { .. }
| Self::AdvisoryUnlockReturnedFalse { .. }
| Self::PartitionExpansionNoLeaves { .. }
| Self::MissingApplyIdentity { .. } => true,
Self::LockTimeout { .. }
| Self::GuardError(_)
| Self::AdvisoryLockFailed { .. }
| Self::AdvisoryLockQueryFailed { .. }
| Self::LedgerWriteFailed { .. }
| Self::LedgerQueryFailed { .. }
| Self::RunIdGenerationFailed { .. }
| Self::NodeIdentityBindingFailed { .. }
| Self::SingleNodeDevProvisioningFailed { .. }
| Self::LedgerBootstrapFailed { .. }
| Self::TransactionalSegmentFailed { .. }
| Self::NonTransactionalSegmentFailed { .. }
| Self::NonTransactionalProgressAckFailed { .. }
| Self::ConfigLoadFailed { .. }
| Self::BaselineProjectionFailed { .. }
| Self::DriftPreflightFailed { .. }
| Self::CatalogQueryFailed { .. }
| Self::PinnedSessionCheckoutFailed { .. } => false,
}
}
}
impl std::fmt::Display for RunnerError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
RunnerError::LockTimeout { path, holder_pid } => match holder_pid {
Some(pid) => write!(
f,
"D025 lock held by another invocation (PID {pid}) at {}; \
refusing to apply",
path.display(),
),
None => write!(
f,
"D025 lock held by another invocation at {}; refusing to apply \
(PID unknown)",
path.display(),
),
},
RunnerError::GuardError(e) => write!(f, "workspace lock error: {e}"),
RunnerError::AdvisoryLockFailed {
bucket,
key,
attempts,
} => write!(
f,
"Postgres advisory lock for bucket database={db} app={app} \
(key=0x{key:016x}) could not be acquired after {attempts} attempts",
db = bucket.database,
app = bucket.app,
),
RunnerError::AdvisoryLockQueryFailed { app_label, source } => write!(
f,
"pg_try_advisory_lock query failed for app `{app_label}`: {source}",
),
RunnerError::ChecksumMismatch(m) => write!(f, "{m}"),
RunnerError::ChecksumFormat(e) => write!(f, "{e}"),
RunnerError::LedgerWriteFailed { version, source } => {
write!(f, "ledger write failed for version `{version}`: {source}")
}
RunnerError::LedgerQueryFailed {
query_label,
source,
} => write!(
f,
"ledger query `{query_label}` failed before the migration could proceed: {source}",
),
RunnerError::RunIdGenerationFailed { source } => write!(
f,
"run_id generation via `SELECT heerid_next()` failed before any \
migration ran: {source}",
),
RunnerError::NodeIdentityBindingFailed { node_id, source } => write!(
f,
"runner node identity binding failed for node {node_id} \
(unregistered/inactive or SQL error): {source}",
),
RunnerError::SingleNodeDevProvisioningFailed {
node_id,
step,
source,
} => write!(
f,
"single-node-dev provisioning failed for node {node_id} during {step}: {source}",
),
RunnerError::LedgerBootstrapFailed { source } => write!(
f,
"ledger bootstrap (CREATE TABLE IF NOT EXISTS djogi_schema_migrations) \
failed: {source}",
),
RunnerError::VersionAlreadyApplied {
version,
applied_at,
} => match applied_at {
Some(when) => write!(
f,
"migration version `{version}` was already applied at {when}; \
re-running is rejected — use `djogi migrations status` to confirm",
),
None => write!(
f,
"migration version `{version}` was already applied; \
re-running is rejected — use `djogi migrations status` to confirm",
),
},
RunnerError::VersionCollisionNonTerminal {
version,
status,
run_id,
} => {
let guidance = match status {
LedgerStatus::Pending => {
"use `djogi migrations status` to inspect it, then `repair_partial_apply` to resolve it in place"
}
LedgerStatus::Failed => {
"use `djogi migrations status` to inspect it, then `repair_resume_partial_apply` if it is still resumable or `repair_partial_apply` otherwise"
}
LedgerStatus::RolledBack => {
"use `djogi migrations status` to inspect it; the rolled_back row is retained as audit history, so use a new migration version or an explicit operator cleanup before re-applying"
}
LedgerStatus::Applied | LedgerStatus::Baseline | LedgerStatus::Faked => {
unreachable!(
"VersionCollisionNonTerminal only carries pending, failed, or rolled_back rows",
)
}
};
write!(
f,
"migration version `{version}` collided with an existing non-terminal \
ledger row (status `{status}`, run_id {run_id}); re-running is rejected \
until that run is reconciled — {guidance}",
status = status.as_db_str(),
)
}
RunnerError::TargetTableNotFound {
bucket,
index_name,
target_table,
} => write!(
f,
"relpages probe for `{index}` could not locate target table `{table}` \
(bucket database={db} app={app}); the index plan does not create \
this table either — check for a typo or a mis-quoted identifier",
index = index_name,
table = target_table,
db = bucket.database,
app = bucket.app,
),
RunnerError::RelpagesThresholdExceeded {
bucket,
index_name,
target_table,
relpages,
threshold,
} => write!(
f,
"relpages probe rejected `CREATE INDEX {index}` on `{table}` (bucket \
database={db} app={app}): {relpages} > {threshold}; \
set `requires_out_of_transaction = true` on the IndexSpec, or \
lower `migrate.strict_concurrent_warnings`",
index = index_name,
table = target_table,
db = bucket.database,
app = bucket.app,
),
RunnerError::SegmentSqlExecutionModeConflict {
segment_index,
segment_kind,
statement_label,
problem,
} => match problem {
SegmentSqlExecutionModeProblem::TransactionControl { keyword } => write!(
f,
"{} segment {segment_index} statement `{statement_label}` embeds top-level \
transaction control `{keyword}`; djogi owns migration transaction boundaries \
and refuses inline BEGIN/COMMIT/SAVEPOINT control",
segment_kind_name(*segment_kind),
),
SegmentSqlExecutionModeProblem::RequiresNonTransactional { statement_shape } => {
write!(
f,
"{} segment {segment_index} statement `{statement_label}` uses `{statement_shape}`, \
which must run in a non-transactional segment",
segment_kind_name(*segment_kind),
)
}
},
RunnerError::TransactionalSegmentFailed {
segment_index,
statement_label,
source,
} => write!(
f,
"transactional segment {segment_index} failed at `{statement_label}`: {source}",
),
RunnerError::NonTransactionalSegmentFailed {
segment_index,
step_index,
statement_label,
applied_steps_count,
source,
} => write!(
f,
"non-transactional segment {segment_index} step {step_index} `{statement_label}` \
failed after {applied_steps_count} successful step(s): {source}",
),
RunnerError::NonTransactionalProgressAckFailed {
segment_index,
step_index,
statement_label,
applied_steps_count,
source,
} => write!(
f,
"non-transactional segment {segment_index} step {} `{statement_label}` \
committed, but the runner failed to durably acknowledge \
applied_steps_count={applied_steps_count}; the row now carries a \
non-tx progress claim and must be reconciled before resume: {source}",
step_index + 1,
),
RunnerError::ConfigLoadFailed { source } => {
write!(f, "failed to load Djogi.toml: {source}")
}
RunnerError::SnapshotPersistFailed { path, source } => {
write!(f, "snapshot persist failed at {}: {source}", path.display(),)
}
RunnerError::BaselineSnapshotShouldNotBeProvided => f.write_str(
"baseline_plan rejects caller-supplied snapshots: baseline projects the \
live database itself; pass `runner_ctx.snapshot = None`",
),
RunnerError::StalePhaseZeroArtifact {
version,
refusal_reason,
} => write!(
f,
"Phase 0 artifact `{version}` refused before migration apply: {refusal_reason}",
),
RunnerError::BaselineProjectionFailed { source } => write!(
f,
"baseline live-DB projection failed before ledger insert: {source}",
),
RunnerError::DriftDetected { bucket, report } => {
let errors = report
.diagnostics
.iter()
.filter(|d| d.severity == super::verify::VerifySeverity::Error)
.count();
write!(
f,
"drift pre-flight refused apply for database={db} app={app}: \
{errors} error-severity drift diagnostic(s) between the \
recorded snapshot and the live database; no migration SQL \
was executed and no pending ledger row was written. Inspect \
with `djogi migrations verify`; reconcile intentional drift \
with `djogi migrations attune`.",
db = bucket.database,
app = bucket.app,
)
}
RunnerError::DriftBaselineMissing { bucket } => write!(
f,
"drift pre-flight refused apply for database={db} app={app}: \
the bucket has applied migration history but no recorded \
snapshot baseline (`schema_snapshot.json`) was found; no \
migration SQL was executed and no pending ledger row was \
written. Restore the snapshot from version control, or \
rebuild it from the live database with `djogi migrations \
repair snapshot-rebuild`.",
db = bucket.database,
app = bucket.app,
),
RunnerError::DriftBaselineCorrupted { bucket, reason } => write!(
f,
"drift pre-flight refused apply for database={db} app={app}: \
the recorded snapshot baseline (`schema_snapshot.json`) exists \
but could not be read: {reason}. No migration SQL was executed \
and no pending ledger row was written. Restore the snapshot \
from version control or rebuild it with \
`djogi migrations repair snapshot-rebuild`.",
db = bucket.database,
app = bucket.app,
),
RunnerError::DriftPreflightFailed { source } => {
write!(f, "drift pre-flight could not run: {source}")
}
RunnerError::PkFlipHazardReplicaSessions {
walsenders,
subscriptions,
} => write!(
f,
"D060 PK-flip cutover refused: logical-replication machinery is active and \
may be applying changes with session_replication_role = 'replica' (which \
suppresses BEFORE row triggers and would leave the autofill skipped). \
Pause the apply worker(s) or `ALTER TABLE ... ENABLE ALWAYS TRIGGER zzz_*` \
before retrying. Walsenders ({nw}): {walsenders:?}; \
enabled subscriptions ({ns}): {subscriptions:?}",
nw = walsenders.len(),
ns = subscriptions.len(),
),
RunnerError::PkFlipHazardPreexistingZzzTrigger {
table,
trigger_names,
} => write!(
f,
"D061 PK-flip cutover refused: pre-existing zzz_* trigger(s) on `{table}` \
collide with the autofill install: {trigger_names:?}. Rename or drop them \
before retrying.",
),
RunnerError::PkFlipHazardDisabledTriggers { table, triggers } => write!(
f,
"D062 PK-flip cutover refused: disabled trigger(s) on `{table}` would \
silently bypass the autofill: {triggers:?}. Re-enable them or pause \
whatever process disabled them before retrying.",
),
RunnerError::PkFlipHazardLongRunningTx {
offenders,
threshold_secs,
} => write!(
f,
"D063 PK-flip cutover refused: {n} transaction(s) have been open longer \
than {threshold_secs}s and would block AccessExclusiveLock or trigger \
lock_timeout. Cancel or terminate them via pg_cancel_backend / \
pg_terminate_backend, then retry. Offenders (pid, age_secs): {offenders:?}",
n = offenders.len(),
),
RunnerError::PkFlipVerificationFailed {
table,
count_violating,
} => write!(
f,
"D064 PK-flip verification halt: table `{table}` has {count_violating} row(s) \
with NULL or stale shadow values. Re-run the backfill (and audit any DISABLE \
TRIGGER / replica writes during the window) before retrying the cutover.",
),
RunnerError::OutOfOrderRejected {
version,
conflicting_version,
conflicting_applied_at,
} => match conflicting_applied_at {
Some(when) => write!(
f,
"version `{version}` would apply out-of-order: peer \
`{conflicting_version}` was already applied at {when}; \
the active OutOfOrderPolicy is Reject. Either rebase \
this migration to a later timestamp, supply \
OutOfOrderPolicy::AllowExplicit with an override reason, \
or run on a non-CI / non-production profile to inherit \
AllowWithDiagnostic."
),
None => write!(
f,
"version `{version}` would apply out-of-order: peer \
`{conflicting_version}` is already applied; \
the active OutOfOrderPolicy is Reject."
),
},
RunnerError::CatalogQueryFailed {
query_label,
source,
} => write!(f, "Postgres catalog query '{query_label}' failed: {source}",),
RunnerError::PinnedSessionCheckoutFailed { source } => write!(
f,
"failed to check out a pinned Postgres session from the pool before \
the migration operation began (GH #274): {source}",
),
RunnerError::AdvisoryUnlockReturnedFalse { key, bucket } => write!(
f,
"D274 pg_advisory_unlock returned false for bucket database={db} app={app} \
(key=0x{key:016x}); the advisory lock was not held on the session that \
called pg_advisory_unlock — this is a session-pinning correctness failure \
(GH #274/#280). The migration SQL and ledger writes may have succeeded; \
inspect the ledger row to determine the actual applied state.",
db = bucket.database,
app = bucket.app,
),
RunnerError::PartitionExpansionNoLeaves {
parent,
statement_label,
} => write!(
f,
"partition expansion for `{statement_label}` refused: \
partitioned parent `{parent}` has 0 leaves in replay-strict mode",
),
RunnerError::MissingApplyIdentity { version } => write!(
f,
"apply refused: version '{version}' requires a binding-capable runner identity \
(not Phase 0, no runner_identity set)",
),
}
}
}
impl std::error::Error for RunnerError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
RunnerError::ChecksumMismatch(e) => Some(e),
RunnerError::ChecksumFormat(e) => Some(e),
RunnerError::GuardError(e) => Some(e),
RunnerError::LedgerWriteFailed { source, .. } => Some(source),
RunnerError::AdvisoryLockQueryFailed { source, .. } => Some(source),
RunnerError::SegmentSqlExecutionModeConflict { .. } => None,
RunnerError::TransactionalSegmentFailed { source, .. } => Some(source),
RunnerError::NonTransactionalSegmentFailed { source, .. } => Some(source),
RunnerError::NonTransactionalProgressAckFailed { source, .. } => Some(source),
RunnerError::ConfigLoadFailed { source } => Some(source),
RunnerError::SnapshotPersistFailed { source, .. } => Some(source),
RunnerError::StalePhaseZeroArtifact { .. } => None,
RunnerError::BaselineProjectionFailed { source } => Some(source.as_ref()),
RunnerError::DriftDetected { .. } => None,
RunnerError::DriftBaselineMissing { .. } => None,
RunnerError::DriftBaselineCorrupted { .. } => None,
RunnerError::DriftPreflightFailed { source } => Some(source.as_ref()),
RunnerError::CatalogQueryFailed { source, .. } => Some(source),
RunnerError::LedgerQueryFailed { source, .. } => Some(source),
RunnerError::RunIdGenerationFailed { source } => Some(source),
RunnerError::NodeIdentityBindingFailed { source, .. } => Some(source),
RunnerError::SingleNodeDevProvisioningFailed { source, .. } => Some(source),
RunnerError::LedgerBootstrapFailed { source } => Some(source),
RunnerError::PinnedSessionCheckoutFailed { source } => Some(source),
RunnerError::VersionCollisionNonTerminal { .. } => None,
_ => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RunnerIdentity {
Selected { id: i32 },
SingleNodeDev,
IdentityFree,
}
impl RunnerIdentity {
pub(crate) fn node_id(self) -> Option<i32> {
match self {
Self::Selected { id } => Some(id),
Self::SingleNodeDev => Some(1),
Self::IdentityFree => None,
}
}
pub(crate) fn requires_binding(self) -> bool {
!matches!(self, Self::IdentityFree)
}
pub(crate) fn provisions_phase_zero_single_node_dev(self) -> bool {
matches!(self, Self::SingleNodeDev)
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub enum DriftBaseline {
Disabled,
Missing,
Snapshot(super::schema::AppliedSchema),
Corrupted(String),
}
pub struct RunnerCtx {
pub bucket: BucketKey,
pub version: String,
pub description: String,
pub checksum_up: String,
pub checksum_down: Option<String>,
pub snapshot: Option<super::schema::AppliedSchema>,
pub snapshot_path: Option<PathBuf>,
pub config: MigrateConfig,
pub out_of_order_policy: super::policy::OutOfOrderPolicy,
pub audit_pool: Option<deadpool_postgres::Pool>,
pub runner_identity: Option<RunnerIdentity>,
pub drift_baseline: DriftBaseline,
}
#[derive(Debug, Clone)]
pub struct RunReport {
pub ledger_id: i64,
pub run_id: i64,
pub transactional_segments: usize,
pub non_transactional_segments: usize,
pub metadata_segments: usize,
pub execution_time_ms: i64,
}
#[expect(
clippy::result_large_err,
reason = "runner guard helpers return the public RunnerError enum unchanged"
)]
fn preflight_phase_zero_apply_artifact(
plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
) -> Result<(), RunnerError> {
if runner_ctx.version != super::bootstrap::PHASE_ZERO_VERSION {
return Ok(());
}
let combined_up = plan
.segments
.iter()
.flat_map(|seg| seg.statements.iter())
.map(|stmt| stmt.up.as_str())
.collect::<Vec<_>>()
.join("\n");
if let Err(refusal_reason) =
super::phase_zero::require_identity_free_phase_zero_migration_artifact(
combined_up.as_bytes(),
)
{
return Err(RunnerError::StalePhaseZeroArtifact {
version: runner_ctx.version.clone(),
refusal_reason,
});
}
Ok(())
}
#[expect(
clippy::result_large_err,
reason = "runner guard helpers return the public RunnerError enum unchanged"
)]
fn verify_plan_checksum(plan: &MigrationPlan, runner_ctx: &RunnerCtx) -> Result<(), RunnerError> {
let computed_up = compute_checksum_for_plan_up(plan);
if let Err(e) =
ledger::verify_checksum(&runner_ctx.version, &runner_ctx.checksum_up, &computed_up)
{
return Err(match e {
VerifyError::Mismatch(m) => RunnerError::ChecksumMismatch(m),
VerifyError::Format(f) => RunnerError::ChecksumFormat(f),
});
}
Ok(())
}
#[expect(
clippy::result_large_err,
reason = "runner guard helpers return the public RunnerError enum unchanged"
)]
fn validate_apply_identity(runner_ctx: &RunnerCtx) -> Result<(), RunnerError> {
if runner_ctx.version == super::bootstrap::PHASE_ZERO_VERSION {
return Ok(());
}
match runner_ctx.runner_identity {
Some(identity) if identity.requires_binding() => Ok(()),
_ => Err(RunnerError::MissingApplyIdentity {
version: runner_ctx.version.clone(),
}),
}
}
#[expect(
clippy::result_large_err,
reason = "rollback guard helpers return the public RollbackError enum unchanged"
)]
fn validate_rollback_identity(runner_ctx: &RunnerCtx) -> Result<(), RollbackError> {
if runner_ctx.version == super::bootstrap::PHASE_ZERO_VERSION {
return Ok(());
}
match runner_ctx.runner_identity {
Some(identity) if identity.requires_binding() => Ok(()),
_ => Err(RollbackError::MissingRollbackIdentity {
version: runner_ctx.version.clone(),
}),
}
}
pub async fn apply_plan(
ctx: &mut DjogiContext,
plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
_guard: &WorkspaceGuard,
) -> Result<RunReport, RunnerError> {
preflight_phase_zero_apply_artifact(plan, runner_ctx)?;
verify_plan_checksum(plan, runner_ctx)?;
validate_apply_identity(runner_ctx)?;
let mut pinned = ctx
.pin_for_migration()
.await
.map_err(|e| RunnerError::PinnedSessionCheckoutFailed { source: e })?;
apply_plan_pinned(&mut pinned, plan, runner_ctx).await
}
async fn apply_plan_pinned(
ctx: &mut PinnedCtx<'_>,
plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
) -> Result<RunReport, RunnerError> {
ledger::bootstrap(ctx)
.await
.map_err(|e| RunnerError::LedgerBootstrapFailed { source: e })?;
let lock_key = advisory_lock_key(&plan.bucket);
acquire_advisory_lock(ctx, &plan.bucket, lock_key).await?;
let result = apply_plan_inner(ctx, plan, runner_ctx).await;
let released = release_advisory_lock(ctx, lock_key).await;
let result = handle_release_result(result, released, &plan.bucket, lock_key);
if result.is_ok() {
ctx.mark_clean();
}
result
}
async fn apply_plan_inner(
ctx: &mut DjogiContext,
plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
) -> Result<RunReport, RunnerError> {
let started = Instant::now();
if matches!(
plan.classification,
super::diff::Classification::PkTypeFlip { .. }
) {
pk_flip_preflight(ctx, runner_ctx, plan).await?;
}
let computed_up = compute_checksum_for_plan_up(plan);
if let Err(e) =
ledger::verify_checksum(&runner_ctx.version, &runner_ctx.checksum_up, &computed_up)
{
return Err(match e {
VerifyError::Mismatch(m) => RunnerError::ChecksumMismatch(m),
VerifyError::Format(f) => RunnerError::ChecksumFormat(f),
});
}
let conflicting_peer = find_higher_applied_version(ctx, &plan.bucket, &runner_ctx.version)
.await
.map_err(|e| RunnerError::LedgerQueryFailed {
query_label: "out_of_order_check",
source: e,
})?;
let is_out_of_order = conflicting_peer.is_some();
if is_out_of_order && !runner_ctx.out_of_order_policy.allows() {
let (conflicting_version, conflicting_applied_at) =
conflicting_peer.unwrap_or_else(|| (String::new(), None));
return Err(RunnerError::OutOfOrderRejected {
version: runner_ctx.version.clone(),
conflicting_version,
conflicting_applied_at,
});
}
if is_out_of_order {
let (conflicting_version, applied_at) = conflicting_peer
.as_ref()
.map(|(v, ts)| (v.as_str(), ts.as_deref()))
.unwrap_or(("", None));
tracing::warn!(
bucket_database = %plan.bucket.database,
bucket_app = %plan.bucket.app,
version = %runner_ctx.version,
conflicting_version,
conflicting_applied_at = applied_at.unwrap_or("<unknown>"),
policy = ?runner_ctx.out_of_order_policy,
"out-of-order migration apply allowed by policy",
);
}
match &runner_ctx.drift_baseline {
DriftBaseline::Disabled => {}
DriftBaseline::Missing | DriftBaseline::Corrupted(_) | DriftBaseline::Snapshot(_) => {
let has_applied_history = bucket_has_applied_history(ctx, &plan.bucket)
.await
.map_err(|e| RunnerError::LedgerQueryFailed {
query_label: "drift_bucket_history",
source: e,
})?;
if has_applied_history {
match &runner_ctx.drift_baseline {
DriftBaseline::Disabled => unreachable!("matched above"),
DriftBaseline::Missing => {
return Err(RunnerError::DriftBaselineMissing {
bucket: plan.bucket.clone(),
});
}
DriftBaseline::Corrupted(reason) => {
return Err(RunnerError::DriftBaselineCorrupted {
bucket: plan.bucket.clone(),
reason: reason.clone(),
});
}
DriftBaseline::Snapshot(snapshot) => {
let report = super::verify::verify_bucket(
ctx,
&plan.bucket,
snapshot,
&PolicyConfig::default(),
false,
false,
)
.await
.map_err(|source| {
RunnerError::DriftPreflightFailed {
source: Box::new(source),
}
})?;
if report.has_errors() {
return Err(RunnerError::DriftDetected {
bucket: plan.bucket.clone(),
report,
});
}
}
}
}
}
}
let initial_note = compose_initial_note(
is_out_of_order,
runner_ctx.out_of_order_policy.override_reason(),
conflicting_peer.as_ref(),
);
let durable_non_tx_note = initial_note.clone();
let (plan_owned, leaves_cache) =
materialize_execution_plan(ctx, plan, PartitionExpansionMode::ApplyLenient).await?;
let plan = &plan_owned;
preflight_segment_sql_execution_compatibility(plan)?;
let mut total_non_tx_steps: i32 = 0;
let mut has_non_tx = false;
for seg in &plan.segments {
if seg.kind == SegmentKind::NonTransactional {
has_non_tx = true;
total_non_tx_steps = total_non_tx_steps.saturating_add(seg.statements.len() as i32);
}
}
let execution_mode = if has_non_tx {
ExecutionMode::NonTransactional
} else {
ExecutionMode::Transactional
};
if runner_ctx.version != super::bootstrap::PHASE_ZERO_VERSION {
let identity = match runner_ctx.runner_identity {
Some(id) => id,
None => {
return Err(RunnerError::MissingApplyIdentity {
version: runner_ctx.version.clone(),
});
}
};
if !identity.requires_binding() {
return Err(RunnerError::MissingApplyIdentity {
version: runner_ctx.version.clone(),
});
}
bind_runner_node_identity(
ctx,
identity.node_id().expect(
"INVARIANT: node_id is Some when requires_binding is true; \
IdentityFree ruled out by refusal gate above",
),
)
.await?;
}
let run_id = generate_run_id(ctx, &runner_ctx.version).await?;
let ledger_row = LedgerRow {
version: runner_ctx.version.clone(),
description: runner_ctx.description.clone(),
checksum_up: runner_ctx.checksum_up.clone(),
checksum_down: runner_ctx.checksum_down.clone(),
execution_mode,
status: LedgerStatus::Pending,
execution_time_ms: 0,
out_of_order_flag: is_out_of_order,
applied_steps_count: 0,
total_steps: if has_non_tx {
Some(total_non_tx_steps)
} else {
None
},
partial_apply_note: initial_note,
run_id,
snapshot_version: SNAPSHOT_FORMAT_VERSION.to_string(),
app_label: plan.bucket.app.clone(),
leaf_identity: serialize_leaf_identity(&leaves_cache),
};
let ledger_id = match ledger::insert_pending(ctx, &ledger_row).await {
Ok(id) => id,
Err(e) => {
if is_unique_violation(&e) {
return Err(classify_duplicate_version_collision(
ctx,
&runner_ctx.version,
&runner_ctx.bucket.app,
)
.await);
}
return Err(RunnerError::LedgerWriteFailed {
version: runner_ctx.version.clone(),
source: e,
});
}
};
let mut transactional_segments = 0usize;
let mut non_transactional_segments = 0usize;
let mut metadata_segments = 0usize;
let mut applied_non_tx_steps: i32 = 0;
let add_table_set = collect_add_table_targets(plan);
for (seg_idx, segment) in plan.segments.iter().enumerate() {
match segment.kind {
SegmentKind::Transactional => {
if let Err(e) =
run_transactional_segment(ctx, segment, seg_idx, runner_ctx, &add_table_set)
.await
{
let note = note_for_failed_transactional_segment(seg_idx, &e);
let _ = ledger::mark_failed(ctx, ledger_id, ¬e).await;
return Err(e);
}
transactional_segments += 1;
}
SegmentKind::NonTransactional => {
match run_non_transactional_segment(
ctx,
segment,
NonTransactionalSegmentRun {
segment_index: seg_idx,
version: &runner_ctx.version,
ledger_id,
prior_steps_completed: applied_non_tx_steps,
total_non_tx_steps,
stable_note: durable_non_tx_note.as_deref(),
runner_ctx,
},
)
.await
{
Ok(steps_completed) => {
applied_non_tx_steps = applied_non_tx_steps.saturating_add(steps_completed);
non_transactional_segments += 1;
}
Err(e) => {
return Err(e);
}
}
}
SegmentKind::MetadataOnly => {
metadata_segments += 1;
}
}
}
if runner_ctx.version == super::bootstrap::PHASE_ZERO_VERSION
&& runner_ctx
.runner_identity
.is_some_and(RunnerIdentity::provisions_phase_zero_single_node_dev)
{
let node_id = super::bootstrap::DEFAULT_NODE_ID;
if let Err(e) = provision_phase_zero_single_node_dev(ctx, node_id).await {
let note = format!("single-node-dev provisioning failed before mark_applied: {e}");
let _ = ledger::mark_failed(ctx, ledger_id, ¬e).await;
return Err(e);
}
}
let elapsed_ms: i64 = elapsed_ms(started);
if is_out_of_order {
ledger::mark_applied_keep_note(ctx, ledger_id, elapsed_ms, applied_non_tx_steps)
.await
.map_err(|e| RunnerError::LedgerWriteFailed {
version: runner_ctx.version.clone(),
source: e,
})?;
} else {
ledger::mark_applied(ctx, ledger_id, elapsed_ms, applied_non_tx_steps)
.await
.map_err(|e| RunnerError::LedgerWriteFailed {
version: runner_ctx.version.clone(),
source: e,
})?;
}
if let (Some(snapshot), Some(path)) = (&runner_ctx.snapshot, &runner_ctx.snapshot_path) {
save_snapshot(snapshot, path).map_err(|e| RunnerError::SnapshotPersistFailed {
path: path.clone(),
source: e,
})?;
}
record_ddl_audit_for_plan(plan, runner_ctx, runner_ctx.snapshot.as_ref()).await;
Ok(RunReport {
ledger_id,
run_id,
transactional_segments,
non_transactional_segments,
metadata_segments,
execution_time_ms: elapsed_ms,
})
}
fn audit_signing_key_from_loaded(
loaded: Result<Option<[u8; 32]>, crate::snapshot::sign::SnapshotKeyError>,
) -> [u8; 32] {
loaded.ok().flatten().unwrap_or([0u8; 32])
}
fn audit_signature_hex_for_snapshot(
snapshot: &super::schema::AppliedSchema,
key: [u8; 32],
) -> Result<String, SnapshotError> {
let snapshot_bytes = super::snapshot::serialize_snapshot(snapshot)?;
let sig = crate::snapshot::sign::sign_snapshot(&snapshot_bytes, &key);
Ok(super::audit::signature_to_hex(&sig))
}
async fn record_ddl_audit_for_plan(
plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
snapshot: Option<&super::schema::AppliedSchema>,
) {
let Some(audit_pool) = runner_ctx.audit_pool.as_ref() else {
return;
};
let key = audit_signing_key_from_loaded(crate::snapshot::sign::load_signing_key_from_env());
let sig_hex_opt: Option<String> = match snapshot {
Some(s) => match audit_signature_hex_for_snapshot(s, key) {
Ok(sig_hex) => Some(sig_hex),
Err(e) => {
tracing::warn!(
target: "djogi::migrate::audit",
error = ?e,
"snapshot re-serialisation for audit signature failed; \
proceeding with NULL signature so the DDL audit row still records the apply",
);
None
}
},
None => None,
};
let audit_djogi_pool = crate::pg::pool::DjogiPool {
inner: audit_pool.clone(),
url: None,
pool_id: crate::pg::pool::next_pool_id(),
};
let mut audit_ctx = DjogiContext::from_pool(audit_djogi_pool);
if let Err(e) = super::audit::bootstrap_ddl_audit(&mut audit_ctx).await {
tracing::warn!(
target: "djogi::migrate::audit",
bucket_database = %plan.bucket.database,
bucket_app = %plan.bucket.app,
error = ?e,
"djogi_ddl_audit bootstrap failed; skipping audit rows for this apply",
);
return;
}
for (seg_idx, segment) in plan.segments.iter().enumerate() {
if segment.kind == SegmentKind::MetadataOnly {
continue;
}
let ddl_sql: String = segment
.statements
.iter()
.map(|s| s.up.as_str())
.collect::<Vec<_>>()
.join(";\n");
if let Err(e) = super::record_ddl_audit(
&mut audit_ctx,
&plan.bucket.database,
&plan.bucket.app,
&ddl_sql,
sig_hex_opt.as_deref(),
)
.await
{
tracing::warn!(
target: "djogi::migrate::audit",
bucket_database = %plan.bucket.database,
bucket_app = %plan.bucket.app,
segment_index = seg_idx,
error = ?e,
"djogi_ddl_audit insert failed; continuing with remaining segments",
);
}
}
}
#[derive(Debug, Clone)]
pub enum LossyRollbackPolicy {
Refuse,
Allow {
reason: String,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RollbackChecksumSide {
Up,
Down,
}
impl RollbackChecksumSide {
const fn as_str(self) -> &'static str {
match self {
Self::Up => "up",
Self::Down => "down",
}
}
}
#[derive(Debug)]
pub enum RollbackError {
Runner {
source: RunnerError,
live_db_committed: bool,
},
LossyRollbackRefused {
offending_labels: Vec<String>,
kinds: Vec<LossyRollbackKind>,
},
VersionNotRollbackable {
version: String,
current_status: LedgerStatus,
},
VersionNotFound { version: String },
BucketAppMismatch {
version: String,
row_app_label: String,
supplied_app: String,
},
DownStatementFailed {
segment_index: usize,
statement_label: String,
live_db_committed: bool,
source: DjogiError,
},
ChecksumDrift {
version: String,
side: RollbackChecksumSide,
ledger: Option<String>,
on_disk: Option<String>,
},
PriorSnapshotMissing,
LeafIdentityMismatch {
version: String,
stored_leaf_identity: String,
current_leaf_identity: String,
},
SnapshotPersistFailed {
path: PathBuf,
source: SnapshotError,
},
StalePhaseZeroDown {
version: String,
refusal_reason: &'static str,
},
MissingRollbackIdentity {
version: String,
},
}
impl RollbackError {
fn runner(source: RunnerError) -> Self {
Self::Runner {
source,
live_db_committed: false,
}
}
fn runner_committed(source: RunnerError) -> Self {
Self::Runner {
source,
live_db_committed: true,
}
}
#[must_use]
pub fn live_db_committed(&self) -> bool {
match self {
Self::Runner {
live_db_committed, ..
}
| Self::DownStatementFailed {
live_db_committed, ..
} => *live_db_committed,
Self::SnapshotPersistFailed { .. } => true,
Self::LossyRollbackRefused { .. }
| Self::VersionNotRollbackable { .. }
| Self::VersionNotFound { .. }
| Self::BucketAppMismatch { .. }
| Self::ChecksumDrift { .. }
| Self::PriorSnapshotMissing
| Self::LeafIdentityMismatch { .. }
| Self::StalePhaseZeroDown { .. }
| Self::MissingRollbackIdentity { .. } => false,
}
}
}
fn render_rollback_checksum_value(value: Option<&str>) -> &str {
value.unwrap_or("<none>")
}
impl std::fmt::Display for RollbackError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
RollbackError::Runner { source, .. } => {
write!(f, "rollback failed at runner level: {source}")
}
RollbackError::LossyRollbackRefused {
offending_labels, ..
} => write!(
f,
"rollback refused: {n} operation(s) carry a lossy down side; \
supply LossyRollbackPolicy::Allow {{ reason }} to proceed: {labels:?}",
n = offending_labels.len(),
labels = offending_labels,
),
RollbackError::VersionNotRollbackable {
version,
current_status,
} => write!(
f,
"version `{version}` is not in a rollbackable status (current: {current})",
current = current_status.as_db_str(),
),
RollbackError::VersionNotFound { version } => {
write!(f, "version `{version}` is not present in the ledger")
}
RollbackError::BucketAppMismatch {
version,
row_app_label,
supplied_app,
} => write!(
f,
"rollback rejected: version `{version}` belongs to app \
`{row_app_label}` but the supplied plan has bucket app \
`{supplied_app}`; the advisory lock would be held for \
the wrong logical bucket",
),
RollbackError::DownStatementFailed {
segment_index,
statement_label,
source,
..
} => write!(
f,
"rollback `down` segment {segment_index} `{statement_label}` failed: {source}",
),
RollbackError::ChecksumDrift {
version,
side,
ledger,
on_disk,
} => {
write!(
f,
"rollback refused: {side} checksum drift for `{version}` \
(ledger={} on_disk={}); restore the exact committed migration \
files or run `djogi migrations repair checksum-drift` before rollback",
render_rollback_checksum_value(ledger.as_deref()),
render_rollback_checksum_value(on_disk.as_deref()),
side = side.as_str(),
)
}
RollbackError::PriorSnapshotMissing => f.write_str(
"rollback requires a prior_snapshot to revert to but the caller passed None",
),
RollbackError::SnapshotPersistFailed { path, source } => {
write!(
f,
"rollback snapshot persist at {} failed: {source}",
path.display()
)
}
RollbackError::LeafIdentityMismatch {
version,
stored_leaf_identity: _,
current_leaf_identity: _,
} => write!(
f,
"[D624] rollback refused: partition leaf identity mismatch for `{version}`",
),
RollbackError::StalePhaseZeroDown {
version,
refusal_reason,
} => write!(
f,
"rollback refused: Phase 0 down SQL for `{version}` is {refusal_reason}; \
refusing before execution to prevent stale artifact mutation",
),
RollbackError::MissingRollbackIdentity { version } => write!(
f,
"rollback refused: non-Phase-0 rollback of `{version}` requires a \
binding-capable runner identity (Selected or SingleNodeDev); \
IdentityFree and missing identity are not allowed for down SQL execution",
),
}
}
}
impl std::error::Error for RollbackError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
RollbackError::Runner { source, .. } => Some(source),
RollbackError::DownStatementFailed { source, .. } => Some(source),
RollbackError::SnapshotPersistFailed { source, .. } => Some(source),
_ => None,
}
}
}
pub async fn rollback_plan(
ctx: &mut DjogiContext,
plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
_guard: &WorkspaceGuard,
lossy_policy: LossyRollbackPolicy,
prior_snapshot: Option<&super::schema::AppliedSchema>,
) -> Result<RollbackReport, RollbackError> {
if prior_snapshot.is_none() && runner_ctx.snapshot_path.is_some() {
return Err(RollbackError::PriorSnapshotMissing);
}
validate_rollback_identity(runner_ctx)?;
let mut pinned = ctx.pin_for_migration().await.map_err(|e| {
RollbackError::runner(RunnerError::PinnedSessionCheckoutFailed { source: e })
})?;
rollback_plan_pinned(&mut pinned, plan, runner_ctx, lossy_policy, prior_snapshot).await
}
#[allow(clippy::result_large_err)]
fn rollback_lossy_allow_reason(
plan: &MigrationPlan,
lossy_policy: &LossyRollbackPolicy,
) -> Result<Option<String>, RollbackError> {
let lossy_ops: Vec<(String, LossyRollbackKind)> = plan
.segments
.iter()
.flat_map(|s| s.statements.iter())
.filter_map(|stmt| stmt.lossy.as_ref().map(|w| (stmt.label.clone(), w.kind)))
.collect();
match (lossy_policy, lossy_ops.is_empty()) {
(LossyRollbackPolicy::Refuse, true) => Ok(None),
(LossyRollbackPolicy::Refuse, false) => Err(RollbackError::LossyRollbackRefused {
offending_labels: lossy_ops.iter().map(|(label, _)| label.clone()).collect(),
kinds: lossy_ops.iter().map(|(_, kind)| *kind).collect(),
}),
(LossyRollbackPolicy::Allow { reason }, _) => Ok(Some(reason.clone())),
}
}
#[expect(
clippy::result_large_err,
reason = "rollback guard helpers return the public RollbackError enum unchanged"
)]
fn preflight_phase_zero_down_payload(
replay_plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
) -> Result<(), RollbackError> {
if runner_ctx.version != super::bootstrap::PHASE_ZERO_VERSION {
return Ok(());
}
let down_sqls = replay_plan
.segments
.iter()
.flat_map(|seg| seg.statements.iter())
.map(|stmt| stmt.down.as_str());
if let Err(refusal_reason) =
super::phase_zero::require_identity_free_phase_zero_down_payload(down_sqls)
{
return Err(RollbackError::StalePhaseZeroDown {
version: runner_ctx.version.clone(),
refusal_reason,
});
}
Ok(())
}
#[expect(
clippy::result_large_err,
reason = "rollback guard helpers return the public RollbackError enum unchanged"
)]
fn validate_rollback_checksum_parity(
row: &LedgerRow,
runner_ctx: &RunnerCtx,
) -> Result<(), RollbackError> {
if row.checksum_up != runner_ctx.checksum_up {
return Err(RollbackError::ChecksumDrift {
version: runner_ctx.version.clone(),
side: RollbackChecksumSide::Up,
ledger: Some(row.checksum_up.clone()),
on_disk: Some(runner_ctx.checksum_up.clone()),
});
}
if row.checksum_down != runner_ctx.checksum_down {
return Err(RollbackError::ChecksumDrift {
version: runner_ctx.version.clone(),
side: RollbackChecksumSide::Down,
ledger: row.checksum_down.clone(),
on_disk: runner_ctx.checksum_down.clone(),
});
}
Ok(())
}
async fn rollback_plan_pinned(
ctx: &mut PinnedCtx<'_>,
plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
lossy_policy: LossyRollbackPolicy,
prior_snapshot: Option<&super::schema::AppliedSchema>,
) -> Result<RollbackReport, RollbackError> {
ledger::bootstrap(ctx)
.await
.map_err(|e| RollbackError::runner(RunnerError::LedgerBootstrapFailed { source: e }))?;
let lock_key = advisory_lock_key(&plan.bucket);
acquire_advisory_lock(ctx, &plan.bucket, lock_key)
.await
.map_err(RollbackError::runner)?;
let result = rollback_handle_lock(ctx, plan, runner_ctx, lossy_policy, prior_snapshot).await;
let released = release_advisory_lock(ctx, lock_key).await;
match (result, released) {
(Ok(r), true) => {
ctx.mark_clean();
Ok(r)
}
(Ok(_), false) => Err(RollbackError::runner_committed(
RunnerError::AdvisoryUnlockReturnedFalse {
key: lock_key,
bucket: plan.bucket.clone(),
},
)),
(Err(e), _) => Err(e),
}
}
async fn rollback_handle_lock(
ctx: &mut PinnedCtx<'_>,
plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
lossy_policy: LossyRollbackPolicy,
prior_snapshot: Option<&super::schema::AppliedSchema>,
) -> Result<RollbackReport, RollbackError> {
let row = load_ledger_row_for_version(ctx, &runner_ctx.version, &runner_ctx.bucket.app)
.await
.map_err(|e| {
RollbackError::runner(RunnerError::LedgerQueryFailed {
query_label: "load_row_for_version",
source: e,
})
})?
.ok_or_else(|| RollbackError::VersionNotFound {
version: runner_ctx.version.clone(),
})?;
if row.app_label != plan.bucket.app {
return Err(RollbackError::BucketAppMismatch {
version: runner_ctx.version.clone(),
row_app_label: row.app_label.clone(),
supplied_app: plan.bucket.app.clone(),
});
}
if !matches!(row.status, LedgerStatus::Applied | LedgerStatus::Faked) {
return Err(RollbackError::VersionNotRollbackable {
version: runner_ctx.version.clone(),
current_status: row.status,
});
}
validate_rollback_checksum_parity(&row, runner_ctx)?;
if let Some(ref stored_identity) = row.leaf_identity {
let pre_cache = compute_leaf_identity_cache(ctx, plan)
.await
.map_err(RollbackError::runner)?;
let pre_identity = serialize_leaf_identity(&pre_cache).unwrap_or_default();
if pre_identity != *stored_identity {
return Err(RollbackError::LeafIdentityMismatch {
version: runner_ctx.version.clone(),
stored_leaf_identity: stored_identity.clone(),
current_leaf_identity: pre_identity,
});
}
}
let (replay_plan, leaves_cache_rollback) =
materialize_execution_plan(ctx, plan, PartitionExpansionMode::ReplayStrict)
.await
.map_err(RollbackError::runner)?;
let current_rollback_identity =
serialize_leaf_identity(&leaves_cache_rollback).unwrap_or_default();
let rollback_leaf_mismatch = if let Some(ref stored_identity) = row.leaf_identity {
current_rollback_identity != *stored_identity
} else {
false
};
if rollback_leaf_mismatch {
return Err(RollbackError::LeafIdentityMismatch {
version: runner_ctx.version.clone(),
stored_leaf_identity: row.leaf_identity.clone().unwrap(),
current_leaf_identity: current_rollback_identity,
});
}
let allow_reason = rollback_lossy_allow_reason(&replay_plan, &lossy_policy)?;
preflight_phase_zero_down_payload(&replay_plan, runner_ctx)?;
if runner_ctx.version == super::bootstrap::PHASE_ZERO_VERSION {
} else {
let identity = match runner_ctx.runner_identity {
Some(id) => id,
None => {
return Err(RollbackError::MissingRollbackIdentity {
version: runner_ctx.version.clone(),
});
}
};
if !identity.requires_binding() {
return Err(RollbackError::MissingRollbackIdentity {
version: runner_ctx.version.clone(),
});
}
bind_runner_node_identity(
ctx,
identity.node_id().expect(
"INVARIANT: node_id is Some when requires_binding is true; \
IdentityFree ruled out by refusal gate above",
),
)
.await
.map_err(RollbackError::runner)?;
}
rollback_inner(ctx, &replay_plan, runner_ctx, prior_snapshot, allow_reason).await
}
async fn rollback_inner(
ctx: &mut DjogiContext,
plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
prior_snapshot: Option<&super::schema::AppliedSchema>,
allow_reason: Option<String>,
) -> Result<RollbackReport, RollbackError> {
let mut transactional_undone = 0usize;
let mut non_transactional_undone = 0usize;
for (rev_idx, segment) in plan.segments.iter().enumerate().rev() {
if segment.kind == SegmentKind::NonTransactional {
rollback_non_transactional_segment(ctx, segment, rev_idx, runner_ctx).await?;
non_transactional_undone += 1;
}
}
let has_transactional = plan
.segments
.iter()
.any(|s| s.kind == SegmentKind::Transactional);
if has_transactional {
ctx.batch_execute("BEGIN")
.await
.map_err(|e| RollbackError::DownStatementFailed {
segment_index: usize::MAX,
statement_label: "<BEGIN compound rollback tx>".to_string(),
live_db_committed: non_transactional_undone > 0,
source: e,
})?;
for (rev_idx, segment) in plan.segments.iter().enumerate().rev() {
if segment.kind != SegmentKind::Transactional {
continue;
}
for stmt in segment.statements.iter().rev() {
if stmt.down.is_empty() {
continue;
}
if let Err(e) = execute_runner_statement(ctx, &stmt.down, runner_ctx).await {
let _ = ctx.batch_execute("ROLLBACK").await;
return Err(RollbackError::DownStatementFailed {
segment_index: rev_idx,
statement_label: stmt.label.clone(),
live_db_committed: non_transactional_undone > 0,
source: e,
});
}
}
transactional_undone += 1;
}
ctx.batch_execute("COMMIT")
.await
.map_err(|e| RollbackError::DownStatementFailed {
segment_index: usize::MAX,
statement_label: "<COMMIT compound rollback tx>".to_string(),
live_db_committed: true,
source: e,
})?;
}
let timestamp = OffsetDateTime::now_utc()
.format(&time::format_description::well_known::Rfc3339)
.unwrap_or_else(|_| "<unknown timestamp>".to_string());
let note = match allow_reason.as_deref() {
Some(reason) => format!("rolled back at {timestamp}; lossy reason: {reason}"),
None => format!("rolled back at {timestamp}"),
};
ctx.execute(
"UPDATE djogi_schema_migrations \
SET status = 'rolled_back', \
applied_steps_count = 0, \
total_steps = NULL, \
partial_apply_note = $2 \
WHERE version = $1 AND app_label = $3",
&[&runner_ctx.version, ¬e, &runner_ctx.bucket.app],
)
.await
.map_err(|e| {
RollbackError::runner_committed(RunnerError::LedgerWriteFailed {
version: runner_ctx.version.clone(),
source: e,
})
})?;
let mut snapshot_reverted = false;
if let (Some(snap), Some(path)) = (prior_snapshot, &runner_ctx.snapshot_path) {
save_snapshot(snap, path).map_err(|e| RollbackError::SnapshotPersistFailed {
path: path.clone(),
source: e,
})?;
snapshot_reverted = true;
}
Ok(RollbackReport {
transactional_undone,
non_transactional_undone,
snapshot_reverted,
lossy_reason: allow_reason,
})
}
async fn rollback_non_transactional_segment(
ctx: &mut DjogiContext,
segment: &Segment,
segment_index: usize,
runner_ctx: &RunnerCtx,
) -> Result<(), RollbackError> {
for stmt in segment.statements.iter().rev() {
if stmt.down.is_empty() {
continue;
}
if let Err(e) = execute_runner_statement(ctx, &stmt.down, runner_ctx).await {
return Err(RollbackError::DownStatementFailed {
segment_index,
statement_label: stmt.label.clone(),
live_db_committed: true,
source: e,
});
}
}
Ok(())
}
#[derive(Debug, Clone)]
pub struct RollbackReport {
pub transactional_undone: usize,
pub non_transactional_undone: usize,
pub snapshot_reverted: bool,
pub lossy_reason: Option<String>,
}
pub async fn fake_apply_plan(
ctx: &mut DjogiContext,
plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
_guard: &WorkspaceGuard,
reason: &str,
) -> Result<RunReport, RunnerError> {
preflight_phase_zero_apply_artifact(plan, runner_ctx)?;
verify_plan_checksum(plan, runner_ctx)?;
validate_apply_identity(runner_ctx)?;
let mut pinned = ctx
.pin_for_migration()
.await
.map_err(|e| RunnerError::PinnedSessionCheckoutFailed { source: e })?;
fake_apply_pinned(&mut pinned, plan, runner_ctx, reason).await
}
async fn fake_apply_pinned(
ctx: &mut PinnedCtx<'_>,
plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
reason: &str,
) -> Result<RunReport, RunnerError> {
ledger::bootstrap(ctx)
.await
.map_err(|e| RunnerError::LedgerBootstrapFailed { source: e })?;
let lock_key = advisory_lock_key(&plan.bucket);
acquire_advisory_lock(ctx, &plan.bucket, lock_key).await?;
let result = fake_apply_inner(ctx, plan, runner_ctx, reason).await;
let released = release_advisory_lock(ctx, lock_key).await;
let result = handle_release_result(result, released, &plan.bucket, lock_key);
if result.is_ok() {
ctx.mark_clean();
}
result
}
async fn fake_apply_inner(
ctx: &mut DjogiContext,
plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
reason: &str,
) -> Result<RunReport, RunnerError> {
let started = Instant::now();
let computed_up = compute_checksum_for_plan_up(plan);
if let Err(e) =
ledger::verify_checksum(&runner_ctx.version, &runner_ctx.checksum_up, &computed_up)
{
return Err(match e {
VerifyError::Mismatch(m) => RunnerError::ChecksumMismatch(m),
VerifyError::Format(f) => RunnerError::ChecksumFormat(f),
});
}
let conflicting_peer = find_higher_applied_version(ctx, &plan.bucket, &runner_ctx.version)
.await
.map_err(|e| RunnerError::LedgerQueryFailed {
query_label: "out_of_order_check",
source: e,
})?;
let is_out_of_order = conflicting_peer.is_some();
if is_out_of_order && !runner_ctx.out_of_order_policy.allows() {
let (conflicting_version, conflicting_applied_at) =
conflicting_peer.unwrap_or_else(|| (String::new(), None));
return Err(RunnerError::OutOfOrderRejected {
version: runner_ctx.version.clone(),
conflicting_version,
conflicting_applied_at,
});
}
if is_out_of_order {
let (conflicting_version, applied_at) = conflicting_peer
.as_ref()
.map(|(v, ts)| (v.as_str(), ts.as_deref()))
.unwrap_or(("", None));
tracing::warn!(
bucket_database = %plan.bucket.database,
bucket_app = %plan.bucket.app,
version = %runner_ctx.version,
conflicting_version,
conflicting_applied_at = applied_at.unwrap_or("<unknown>"),
policy = ?runner_ctx.out_of_order_policy,
"out-of-order fake-apply allowed by policy",
);
}
let timestamp = OffsetDateTime::now_utc()
.format(&time::format_description::well_known::Rfc3339)
.unwrap_or_else(|_| "<unknown timestamp>".to_string());
let fake_note = format!("faked at {timestamp}; reason: {reason}");
let ooo_note = compose_initial_note(
is_out_of_order,
runner_ctx.out_of_order_policy.override_reason(),
conflicting_peer.as_ref(),
);
let note = match ooo_note {
Some(note_str) => format!("{fake_note}; {note_str}"),
None => fake_note,
};
if runner_ctx.version != super::bootstrap::PHASE_ZERO_VERSION {
let identity = match runner_ctx.runner_identity {
Some(id) => id,
None => {
return Err(RunnerError::MissingApplyIdentity {
version: runner_ctx.version.clone(),
});
}
};
if !identity.requires_binding() {
return Err(RunnerError::MissingApplyIdentity {
version: runner_ctx.version.clone(),
});
}
bind_runner_node_identity(
ctx,
identity.node_id().expect(
"INVARIANT: node_id is Some when requires_binding is true; \
IdentityFree ruled out by refusal gate above",
),
)
.await?;
}
let run_id = generate_run_id(ctx, &runner_ctx.version).await?;
let row = LedgerRow {
version: runner_ctx.version.clone(),
description: runner_ctx.description.clone(),
checksum_up: runner_ctx.checksum_up.clone(),
checksum_down: runner_ctx.checksum_down.clone(),
execution_mode: ExecutionMode::Transactional,
status: LedgerStatus::Faked,
execution_time_ms: 0,
out_of_order_flag: is_out_of_order,
applied_steps_count: 0,
total_steps: None,
partial_apply_note: Some(note.clone()),
run_id,
snapshot_version: SNAPSHOT_FORMAT_VERSION.to_string(),
app_label: plan.bucket.app.clone(),
leaf_identity: None,
};
let ledger_id = match ledger::insert_pending(ctx, &row).await {
Ok(id) => id,
Err(e) => {
if is_unique_violation(&e) {
return Err(classify_duplicate_version_collision(
ctx,
&runner_ctx.version,
&runner_ctx.bucket.app,
)
.await);
}
return Err(RunnerError::LedgerWriteFailed {
version: runner_ctx.version.clone(),
source: e,
});
}
};
if let (Some(snapshot), Some(path)) = (&runner_ctx.snapshot, &runner_ctx.snapshot_path) {
save_snapshot(snapshot, path).map_err(|e| RunnerError::SnapshotPersistFailed {
path: path.clone(),
source: e,
})?;
}
let elapsed = elapsed_ms(started);
Ok(RunReport {
ledger_id,
run_id,
transactional_segments: 0,
non_transactional_segments: 0,
metadata_segments: 0,
execution_time_ms: elapsed,
})
}
pub async fn baseline_plan(
ctx: &mut DjogiContext,
bucket: &BucketKey,
runner_ctx: &RunnerCtx,
_guard: &WorkspaceGuard,
reason: &str,
) -> Result<RunReport, RunnerError> {
if runner_ctx.snapshot.is_some() {
return Err(RunnerError::BaselineSnapshotShouldNotBeProvided);
}
validate_apply_identity(runner_ctx)?;
let mut pinned = ctx
.pin_for_migration()
.await
.map_err(|e| RunnerError::PinnedSessionCheckoutFailed { source: e })?;
baseline_pinned(&mut pinned, bucket, runner_ctx, reason).await
}
async fn baseline_pinned(
ctx: &mut PinnedCtx<'_>,
bucket: &BucketKey,
runner_ctx: &RunnerCtx,
reason: &str,
) -> Result<RunReport, RunnerError> {
ledger::bootstrap(ctx)
.await
.map_err(|e| RunnerError::LedgerBootstrapFailed { source: e })?;
let lock_key = advisory_lock_key(bucket);
acquire_advisory_lock(ctx, bucket, lock_key).await?;
let result = baseline_inner(ctx, bucket, runner_ctx, reason).await;
let released = release_advisory_lock(ctx, lock_key).await;
let result = handle_release_result(result, released, bucket, lock_key);
if result.is_ok() {
ctx.mark_clean();
}
result
}
async fn baseline_inner(
ctx: &mut DjogiContext,
bucket: &BucketKey,
runner_ctx: &RunnerCtx,
reason: &str,
) -> Result<RunReport, RunnerError> {
let started = Instant::now();
let projected = super::verify::live_schema_for_repair(ctx, bucket, None)
.await
.map_err(|e| RunnerError::BaselineProjectionFailed {
source: Box::new(e),
})?;
let checksum_up = checksum_for_baseline_snapshot(&projected);
let timestamp = OffsetDateTime::now_utc()
.format(&time::format_description::well_known::Rfc3339)
.unwrap_or_else(|_| "<unknown timestamp>".to_string());
let note = format!(
"baseline established at {timestamp} for bucket database={db} app={app}; reason: {reason}",
db = bucket.database,
app = bucket.app,
);
if runner_ctx.version != super::bootstrap::PHASE_ZERO_VERSION {
let identity = match runner_ctx.runner_identity {
Some(id) => id,
None => {
return Err(RunnerError::MissingApplyIdentity {
version: runner_ctx.version.clone(),
});
}
};
if !identity.requires_binding() {
return Err(RunnerError::MissingApplyIdentity {
version: runner_ctx.version.clone(),
});
}
bind_runner_node_identity(
ctx,
identity.node_id().expect(
"INVARIANT: node_id is Some when requires_binding is true; \
IdentityFree ruled out by refusal gate above",
),
)
.await?;
}
let run_id = generate_run_id(ctx, &runner_ctx.version).await?;
let row = LedgerRow {
version: runner_ctx.version.clone(),
description: format!("<baseline> {}", runner_ctx.description),
checksum_up: checksum_up.clone(),
checksum_down: None,
execution_mode: ExecutionMode::Transactional,
status: LedgerStatus::Baseline,
execution_time_ms: 0,
out_of_order_flag: false,
applied_steps_count: 0,
total_steps: None,
partial_apply_note: Some(note),
run_id,
snapshot_version: SNAPSHOT_FORMAT_VERSION.to_string(),
app_label: bucket.app.clone(),
leaf_identity: None,
};
let ledger_id = match ledger::insert_pending(ctx, &row).await {
Ok(id) => id,
Err(e) => {
if is_unique_violation(&e) {
return Err(classify_duplicate_version_collision(
ctx,
&runner_ctx.version,
&runner_ctx.bucket.app,
)
.await);
}
return Err(RunnerError::LedgerWriteFailed {
version: runner_ctx.version.clone(),
source: e,
});
}
};
if let Some(path) = &runner_ctx.snapshot_path {
save_snapshot(&projected, path).map_err(|e| RunnerError::SnapshotPersistFailed {
path: path.clone(),
source: e,
})?;
}
let elapsed = elapsed_ms(started);
Ok(RunReport {
ledger_id,
run_id,
transactional_segments: 0,
non_transactional_segments: 0,
metadata_segments: 0,
execution_time_ms: elapsed,
})
}
pub(crate) fn checksum_for_baseline_snapshot(schema: &super::schema::AppliedSchema) -> String {
let json = serde_json::to_string(schema).unwrap_or_default();
compute_checksum([json])
}
async fn load_ledger_row_for_version(
ctx: &mut DjogiContext,
version: &str,
app_label: &str,
) -> Result<Option<LedgerRow>, DjogiError> {
load_full_row_by_version(ctx, version, app_label).await
}
async fn execute_runner_statement(
ctx: &mut DjogiContext,
sql: &str,
runner_ctx: &RunnerCtx,
) -> Result<(), DjogiError> {
if runner_ctx.version == super::bootstrap::PHASE_ZERO_VERSION {
if let Err(reason) = super::phase_zero::require_current_phase_zero_statement(sql) {
return Err(DjogiError::StalePhaseZeroStatement {
refusal_reason: reason,
statement: truncate_for_log(sql),
});
}
ctx.batch_execute(sql).await
} else {
guarded_batch_execute(ctx, sql).await
}
}
fn truncate_for_log(sql: &str) -> String {
const MAX_LEN: usize = 256;
let trimmed = sql.trim();
if trimmed.len() <= MAX_LEN {
trimmed.to_string()
} else {
let mut end = 0;
for (idx, ch) in trimmed.char_indices() {
let next = idx + ch.len_utf8();
if next > MAX_LEN {
break;
}
end = next;
}
format!("{}...", &trimmed[..end])
}
}
async fn run_transactional_segment(
ctx: &mut DjogiContext,
segment: &Segment,
segment_index: usize,
runner_ctx: &RunnerCtx,
add_table_set: &BTreeSet<String>,
) -> Result<(), RunnerError> {
let all_verify = !segment.statements.is_empty()
&& segment
.statements
.iter()
.all(|s| s.label.starts_with("PkFlipVerify "));
if all_verify {
for stmt in &segment.statements {
let table = stmt
.label
.split_whitespace()
.nth(1)
.unwrap_or("")
.to_string();
let row = ctx.query_one(&stmt.up, &[]).await.map_err(|e| {
RunnerError::TransactionalSegmentFailed {
segment_index,
statement_label: stmt.label.clone(),
source: e,
}
})?;
let count: i64 = row.try_get(0).unwrap_or(0);
if count > 0 {
return Err(RunnerError::PkFlipVerificationFailed {
table,
count_violating: count,
});
}
}
return Ok(());
}
for stmt in &segment.statements {
if let Some((index_name, target_table)) = parse_create_index_statement(stmt) {
relpages_probe(ctx, runner_ctx, &index_name, &target_table, add_table_set).await?;
}
}
ctx.batch_execute("BEGIN")
.await
.map_err(|e| RunnerError::TransactionalSegmentFailed {
segment_index,
statement_label: "<BEGIN>".to_string(),
source: e,
})?;
for stmt in &segment.statements {
if let Err(e) = execute_runner_statement(ctx, &stmt.up, runner_ctx).await {
let _ = ctx.batch_execute("ROLLBACK").await;
return Err(RunnerError::TransactionalSegmentFailed {
segment_index,
statement_label: stmt.label.clone(),
source: e,
});
}
}
ctx.batch_execute("COMMIT")
.await
.map_err(|e| RunnerError::TransactionalSegmentFailed {
segment_index,
statement_label: "<COMMIT>".to_string(),
source: e,
})?;
Ok(())
}
struct NonTransactionalSegmentRun<'a> {
segment_index: usize,
version: &'a str,
ledger_id: i64,
prior_steps_completed: i32,
total_non_tx_steps: i32,
stable_note: Option<&'a str>,
runner_ctx: &'a RunnerCtx,
}
async fn run_non_transactional_segment(
ctx: &mut DjogiContext,
segment: &Segment,
run: NonTransactionalSegmentRun<'_>,
) -> Result<i32, RunnerError> {
let mut completed: i32 = 0;
for (step_idx, stmt) in segment.statements.iter().enumerate() {
let claimed_step = run
.prior_steps_completed
.saturating_add(completed)
.saturating_add(1);
let claim_note = ledger::format_non_tx_progress_claim(
run.stable_note,
claimed_step,
Some(run.total_non_tx_steps),
run.segment_index,
&stmt.label,
);
ledger::claim_non_tx_progress(ctx, run.ledger_id, &claim_note)
.await
.map_err(|e| RunnerError::LedgerWriteFailed {
version: run.version.to_string(),
source: e,
})?;
if let Err(e) = execute_runner_statement(ctx, &stmt.up, run.runner_ctx).await {
let total_so_far = run.prior_steps_completed.saturating_add(completed);
let note = format!(
"non-tx step {step} of segment {seg} failed: {label} — {e}",
step = step_idx + 1,
seg = run.segment_index,
label = stmt.label,
);
let _ = ledger::mark_partial(ctx, run.ledger_id, total_so_far, ¬e).await;
return Err(RunnerError::NonTransactionalSegmentFailed {
segment_index: run.segment_index,
step_index: step_idx,
statement_label: stmt.label.clone(),
applied_steps_count: total_so_far,
source: e,
});
}
completed = completed.saturating_add(1);
let total_so_far = run.prior_steps_completed.saturating_add(completed);
ledger::ack_non_tx_progress(ctx, run.ledger_id, total_so_far, run.stable_note)
.await
.map_err(|e| RunnerError::NonTransactionalProgressAckFailed {
segment_index: run.segment_index,
step_index: step_idx,
statement_label: stmt.label.clone(),
applied_steps_count: total_so_far,
source: e,
})?;
}
Ok(completed)
}
pub fn advisory_lock_key(bucket: &BucketKey) -> i64 {
let mut hasher = Sha256::new();
hasher.update(b"djogi:advisory_lock:");
hasher.update(bucket.database.as_bytes());
hasher.update(b"\x00");
hasher.update(bucket.app.as_bytes());
let digest = hasher.finalize();
let mut buf = [0u8; 8];
buf.copy_from_slice(&digest[..8]);
i64::from_be_bytes(buf)
}
pub(crate) async fn acquire_advisory_lock(
ctx: &mut PinnedCtx<'_>,
bucket: &BucketKey,
key: i64,
) -> Result<(), RunnerError> {
assert!(
!ctx.is_pool_backed(),
"acquire_advisory_lock called on a pool-backed context — \
the advisory lock would be acquired on an arbitrary pool \
connection and subsequent operations would run on different \
connections. Callers must use ctx.pin_for_migration() first \
(GH #274 / #331).",
);
const MAX_ATTEMPTS: u32 = 600; const RETRY_INTERVAL: std::time::Duration = std::time::Duration::from_millis(50);
for attempt in 0..MAX_ATTEMPTS {
let row = ctx
.query_one("SELECT pg_try_advisory_lock($1)", &[&key])
.await
.map_err(|e| RunnerError::AdvisoryLockQueryFailed {
app_label: bucket.app.clone(),
source: e,
})?;
let acquired: bool = row
.try_get(0)
.map_err(|e| RunnerError::AdvisoryLockQueryFailed {
app_label: bucket.app.clone(),
source: DjogiError::from(e),
})?;
if acquired {
return Ok(());
}
tokio::time::sleep(RETRY_INTERVAL).await;
if attempt + 1 == MAX_ATTEMPTS {
return Err(RunnerError::AdvisoryLockFailed {
bucket: bucket.clone(),
key,
attempts: MAX_ATTEMPTS,
});
}
}
Err(RunnerError::AdvisoryLockFailed {
bucket: bucket.clone(),
key,
attempts: MAX_ATTEMPTS,
})
}
pub(crate) async fn release_advisory_lock(ctx: &mut PinnedCtx<'_>, key: i64) -> bool {
assert!(
!ctx.is_pool_backed(),
"release_advisory_lock called on a pool-backed context — \
the unlock would run on a different connection than the one \
that holds the lock. Callers must use ctx.pin_for_migration() \
first (GH #274 / #331).",
);
let row = match ctx
.query_one("SELECT pg_advisory_unlock($1)", &[&key])
.await
{
Ok(r) => r,
Err(e) => {
tracing::warn!(
?e,
key,
"pg_advisory_unlock query failed; lock will auto-release on session close",
);
return true;
}
};
let released: bool = match row.try_get(0) {
Ok(v) => v,
Err(e) => {
tracing::warn!(
?e,
key,
"pg_advisory_unlock result could not be decoded; assuming released",
);
return true;
}
};
if !released {
tracing::error!(
key,
"pg_advisory_unlock returned false — lock was not held on this Postgres \
session. This indicates a session-pinning bug: the advisory lock was \
acquired on a different physical backend than the one executing the \
migration operations (GH #274/#280).",
);
}
released
}
#[allow(clippy::result_large_err)]
pub(crate) fn handle_release_result<T>(
result: Result<T, RunnerError>,
released: bool,
bucket: &BucketKey,
key: i64,
) -> Result<T, RunnerError> {
match result {
Ok(v) if released => Ok(v),
Ok(_) => Err(RunnerError::AdvisoryUnlockReturnedFalse {
key,
bucket: bucket.clone(),
}),
Err(e) => Err(e), }
}
async fn relpages_probe(
ctx: &mut DjogiContext,
runner_ctx: &RunnerCtx,
index_name: &str,
target_table: &str,
add_table_set: &BTreeSet<String>,
) -> Result<(), RunnerError> {
let row_opt = ctx
.query_opt(
"SELECT relpages FROM pg_class WHERE relname = $1 AND relkind = 'r'",
&[&target_table],
)
.await
.map_err(|e| RunnerError::CatalogQueryFailed {
query_label: "pg_class relpages",
source: e,
})?;
let relpages: i32 = match row_opt {
Some(r) => r
.try_get::<_, i32>(0)
.map_err(|e| RunnerError::CatalogQueryFailed {
query_label: "pg_class relpages",
source: DjogiError::from(e),
})?,
None => {
if add_table_set.contains(target_table) {
0
} else {
return Err(RunnerError::TargetTableNotFound {
bucket: runner_ctx.bucket.clone(),
index_name: index_name.to_string(),
target_table: target_table.to_string(),
});
}
}
};
let threshold = runner_ctx.config.concurrent_warn_relpages;
if (relpages as i64) > (threshold as i64) {
if runner_ctx.config.strict_concurrent_warnings {
return Err(RunnerError::RelpagesThresholdExceeded {
bucket: runner_ctx.bucket.clone(),
index_name: index_name.to_string(),
target_table: target_table.to_string(),
relpages,
threshold,
});
} else {
tracing::warn!(
bucket_database = %runner_ctx.bucket.database,
bucket_app = %runner_ctx.bucket.app,
index_name,
target_table,
relpages,
threshold,
"transactional CREATE INDEX on a large table will hold ACCESS EXCLUSIVE \
for the duration; consider opting into CREATE INDEX CONCURRENTLY",
);
}
}
Ok(())
}
async fn pk_flip_preflight(
ctx: &mut DjogiContext,
runner_ctx: &RunnerCtx,
plan: &MigrationPlan,
) -> Result<(), RunnerError> {
let mut tables: BTreeSet<String> = BTreeSet::new();
for seg in &plan.segments {
for stmt in &seg.statements {
if let Some(parent) = stmt.label.strip_prefix("PkFlipPrep ") {
tables.insert(parent.to_string());
} else if let Some(parent) = stmt.label.strip_prefix("PkFlipCutover ") {
tables.insert(parent.to_string());
} else if let Some(parent) = stmt.label.strip_prefix("PkFlipPartitionedPrep ") {
tables.insert(parent.to_string());
}
for child in scan_alter_table_targets(&stmt.up) {
tables.insert(child);
}
}
}
let walsender_rows = ctx
.query_all(
"SELECT COALESCE(application_name, ''), COALESCE(client_addr::text, '') \
FROM pg_stat_replication",
&[],
)
.await
.map_err(|e| RunnerError::CatalogQueryFailed {
query_label: "pg_stat_replication",
source: e,
})?;
let mut walsenders: Vec<(String, String)> = Vec::with_capacity(walsender_rows.len());
for r in &walsender_rows {
let app: String = r.try_get(0).unwrap_or_default();
let client: String = r.try_get(1).unwrap_or_default();
walsenders.push((app, client));
}
let sub_rows = ctx
.query_all(
"SELECT subname FROM pg_subscription \
WHERE subdbid = (SELECT oid FROM pg_database WHERE datname = current_database()) \
AND subenabled = true",
&[],
)
.await
.map_err(|e| RunnerError::CatalogQueryFailed {
query_label: "pg_subscription",
source: e,
})?;
let mut subscriptions: Vec<String> = Vec::with_capacity(sub_rows.len());
for r in &sub_rows {
let name: String = r.try_get(0).unwrap_or_default();
subscriptions.push(name);
}
if !walsenders.is_empty() || !subscriptions.is_empty() {
return Err(RunnerError::PkFlipHazardReplicaSessions {
walsenders,
subscriptions,
});
}
for table in &tables {
let zzz_rows = ctx
.query_all(
"SELECT tgname FROM pg_trigger \
WHERE tgrelid = (SELECT oid FROM pg_class WHERE relname = $1 AND relkind = 'r' LIMIT 1) \
AND NOT tgisinternal \
AND tgname LIKE 'zzz\\_%' ESCAPE '\\'",
&[table],
)
.await
.map_err(|e| RunnerError::CatalogQueryFailed {
query_label: "pg_trigger zzz scan",
source: e,
})?;
if !zzz_rows.is_empty() {
let names: Vec<String> = zzz_rows
.iter()
.map(|r| r.try_get::<_, String>(0).unwrap_or_default())
.collect();
return Err(RunnerError::PkFlipHazardPreexistingZzzTrigger {
table: table.clone(),
trigger_names: names,
});
}
let disabled_rows = ctx
.query_all(
"SELECT tgname, tgenabled FROM pg_trigger \
WHERE tgrelid = (SELECT oid FROM pg_class WHERE relname = $1 AND relkind = 'r' LIMIT 1) \
AND NOT tgisinternal \
AND tgenabled <> 'O'",
&[table],
)
.await
.map_err(|e| RunnerError::CatalogQueryFailed {
query_label: "pg_trigger disabled scan",
source: e,
})?;
if !disabled_rows.is_empty() {
let mut triggers: Vec<(String, char)> = Vec::with_capacity(disabled_rows.len());
for r in &disabled_rows {
let name: String = r.try_get(0).unwrap_or_default();
let raw: i8 = r.try_get(1).unwrap_or(0);
let ch = if raw >= 0 { (raw as u8) as char } else { '?' };
triggers.push((name, ch));
}
return Err(RunnerError::PkFlipHazardDisabledTriggers {
table: table.clone(),
triggers,
});
}
}
let threshold = runner_ctx.config.pk_flip_long_tx_threshold_secs;
if threshold > 0 {
let long_rows = ctx
.query_all(
"SELECT pid, EXTRACT(EPOCH FROM (now() - xact_start))::bigint AS age \
FROM pg_stat_activity \
WHERE pid <> pg_backend_pid() \
AND xact_start IS NOT NULL \
AND now() - xact_start > make_interval(secs => $1::int)",
&[&(threshold as i32)],
)
.await
.map_err(|e| RunnerError::CatalogQueryFailed {
query_label: "pg_stat_activity long tx scan",
source: e,
})?;
if !long_rows.is_empty() {
let mut offenders: Vec<(i32, i64)> = Vec::with_capacity(long_rows.len());
for r in &long_rows {
let pid: i32 = r.try_get(0).unwrap_or(0);
let age: i64 = r.try_get(1).unwrap_or(0);
offenders.push((pid, age));
}
return Err(RunnerError::PkFlipHazardLongRunningTx {
offenders,
threshold_secs: threshold,
});
}
}
Ok(())
}
fn scan_alter_table_targets(sql: &str) -> Vec<String> {
const MARKER: &[u8] = b"ALTER TABLE ";
let bytes = sql.as_bytes();
let mut out: Vec<String> = Vec::new();
let mut i = 0usize;
while i + MARKER.len() <= bytes.len() {
if &bytes[i..i + MARKER.len()] == MARKER {
let start = i + MARKER.len();
let (id, end) = if start < bytes.len() && bytes[start] == b'"' {
let id_start = start + 1;
let mut j = id_start;
while j < bytes.len() && bytes[j] != b'"' {
j += 1;
}
if j > id_start && j < bytes.len() {
(
String::from_utf8_lossy(&bytes[id_start..j]).into_owned(),
j + 1,
)
} else {
(String::new(), bytes.len())
}
} else {
let id_start = start;
let mut j = id_start;
let mut len = 0usize;
while j < bytes.len() && len < 63 {
let b = bytes[j];
let valid = if len == 0 {
b.is_ascii_alphabetic() || b == b'_'
} else {
b.is_ascii_alphanumeric() || b == b'_'
};
if !valid {
break;
}
j += 1;
len += 1;
}
if j > id_start {
(String::from_utf8_lossy(&bytes[id_start..j]).into_owned(), j)
} else {
(String::new(), j)
}
};
if !id.is_empty() {
out.push(id);
}
i = end;
} else {
i += 1;
}
}
out
}
fn segment_kind_name(kind: SegmentKind) -> &'static str {
match kind {
SegmentKind::Transactional => "transactional",
SegmentKind::NonTransactional => "non-transactional",
SegmentKind::MetadataOnly => "metadata-only",
}
}
#[allow(clippy::result_large_err)]
fn preflight_segment_sql_execution_compatibility(plan: &MigrationPlan) -> Result<(), RunnerError> {
for (segment_index, segment) in plan.segments.iter().enumerate() {
if segment.kind == SegmentKind::MetadataOnly {
continue;
}
for statement in &segment.statements {
if let Some(problem) =
classify_segment_sql_execution_mode_problem(segment.kind, &statement.up)
{
return Err(RunnerError::SegmentSqlExecutionModeConflict {
segment_index,
segment_kind: segment.kind,
statement_label: statement.label.clone(),
problem,
});
}
}
}
Ok(())
}
fn classify_segment_sql_execution_mode_problem(
segment_kind: SegmentKind,
sql: &str,
) -> Option<SegmentSqlExecutionModeProblem> {
let bytes = sql.as_bytes();
let mut idx = 0usize;
let first = next_sql_leading_keyword(bytes, &mut idx)?;
let second = next_sql_leading_keyword(bytes, &mut idx);
let third = next_sql_leading_keyword(bytes, &mut idx);
let fourth = next_sql_leading_keyword(bytes, &mut idx);
if token_eq(first, "BEGIN") {
return Some(SegmentSqlExecutionModeProblem::TransactionControl { keyword: "BEGIN" });
}
if token_eq(first, "START") && second.is_some_and(|tok| token_eq(tok, "TRANSACTION")) {
return Some(SegmentSqlExecutionModeProblem::TransactionControl {
keyword: "START TRANSACTION",
});
}
if token_eq(first, "COMMIT") {
return Some(SegmentSqlExecutionModeProblem::TransactionControl { keyword: "COMMIT" });
}
if token_eq(first, "ROLLBACK") {
return Some(SegmentSqlExecutionModeProblem::TransactionControl {
keyword: "ROLLBACK",
});
}
if token_eq(first, "SAVEPOINT") {
return Some(SegmentSqlExecutionModeProblem::TransactionControl {
keyword: "SAVEPOINT",
});
}
if token_eq(first, "RELEASE") && second.is_some_and(|tok| token_eq(tok, "SAVEPOINT")) {
return Some(SegmentSqlExecutionModeProblem::TransactionControl {
keyword: "RELEASE SAVEPOINT",
});
}
if segment_kind != SegmentKind::Transactional {
return None;
}
if token_eq(first, "CREATE")
&& second.is_some_and(|tok| token_eq(tok, "INDEX"))
&& third.is_some_and(|tok| token_eq(tok, "CONCURRENTLY"))
{
return Some(SegmentSqlExecutionModeProblem::RequiresNonTransactional {
statement_shape: "CREATE INDEX CONCURRENTLY",
});
}
if token_eq(first, "CREATE")
&& second.is_some_and(|tok| token_eq(tok, "UNIQUE"))
&& third.is_some_and(|tok| token_eq(tok, "INDEX"))
&& fourth.is_some_and(|tok| token_eq(tok, "CONCURRENTLY"))
{
return Some(SegmentSqlExecutionModeProblem::RequiresNonTransactional {
statement_shape: "CREATE UNIQUE INDEX CONCURRENTLY",
});
}
if token_eq(first, "DROP")
&& second.is_some_and(|tok| token_eq(tok, "INDEX"))
&& third.is_some_and(|tok| token_eq(tok, "CONCURRENTLY"))
{
return Some(SegmentSqlExecutionModeProblem::RequiresNonTransactional {
statement_shape: "DROP INDEX CONCURRENTLY",
});
}
None
}
fn next_sql_leading_keyword<'a>(bytes: &'a [u8], idx: &mut usize) -> Option<&'a [u8]> {
skip_sql_leading_ws_and_comments(bytes, idx);
if *idx >= bytes.len() || !is_sql_ident_start(bytes[*idx]) {
return None;
}
let start = *idx;
*idx += 1;
while *idx < bytes.len() && is_sql_ident_continue(bytes[*idx]) {
*idx += 1;
}
Some(&bytes[start..*idx])
}
fn skip_sql_leading_ws_and_comments(bytes: &[u8], idx: &mut usize) {
loop {
while *idx < bytes.len() && bytes[*idx].is_ascii_whitespace() {
*idx += 1;
}
if *idx + 1 < bytes.len() && bytes[*idx] == b'-' && bytes[*idx + 1] == b'-' {
*idx += 2;
while *idx < bytes.len() && bytes[*idx] != b'\n' {
*idx += 1;
}
continue;
}
if *idx + 1 < bytes.len() && bytes[*idx] == b'/' && bytes[*idx + 1] == b'*' {
*idx += 2;
let mut depth = 1usize;
while *idx < bytes.len() && depth > 0 {
if *idx + 1 < bytes.len() && bytes[*idx] == b'/' && bytes[*idx + 1] == b'*' {
depth += 1;
*idx += 2;
continue;
}
if *idx + 1 < bytes.len() && bytes[*idx] == b'*' && bytes[*idx + 1] == b'/' {
depth -= 1;
*idx += 2;
continue;
}
*idx += 1;
}
continue;
}
return;
}
}
fn is_sql_ident_start(byte: u8) -> bool {
byte.is_ascii_alphabetic() || byte == b'_'
}
fn is_sql_ident_continue(byte: u8) -> bool {
byte.is_ascii_alphanumeric() || byte == b'_'
}
fn token_eq(token: &[u8], expected: &str) -> bool {
token.eq_ignore_ascii_case(expected.as_bytes())
}
fn compute_checksum_for_plan_up(plan: &MigrationPlan) -> String {
let fragments: Vec<&str> = plan
.segments
.iter()
.flat_map(|s| s.statements.iter())
.map(|s| s.up.as_str())
.collect();
compute_checksum(fragments)
}
fn parse_create_index_statement(stmt: &OperationSql) -> Option<(String, String)> {
let label = stmt.label.as_str();
let index_name = label.strip_prefix("AddIndex ")?.to_string();
let bytes = stmt.up.as_bytes();
let needle = b" ON \"";
let mut i = 0usize;
while i + needle.len() <= bytes.len() {
if &bytes[i..i + needle.len()] == needle {
let start = i + needle.len();
let mut j = start;
while j < bytes.len() && bytes[j] != b'"' {
j += 1;
}
if j > start && j < bytes.len() {
let table = String::from_utf8_lossy(&bytes[start..j]).into_owned();
return Some((index_name, table));
}
return None;
}
i += 1;
}
None
}
fn compose_initial_note(
is_out_of_order: bool,
override_reason: Option<&str>,
conflicting_peer: Option<&(String, Option<String>)>,
) -> Option<String> {
match (is_out_of_order, override_reason) {
(true, Some(reason)) if !reason.is_empty() => {
let header = match conflicting_peer {
Some((peer, Some(ts))) => {
format!("out-of-order apply (peer {peer} applied at {ts}) override: {reason}")
}
Some((peer, None)) => {
format!("out-of-order apply (peer {peer}) override: {reason}")
}
None => format!("out-of-order apply override: {reason}"),
};
Some(header)
}
(true, None) => match conflicting_peer {
Some((peer, Some(ts))) => Some(format!(
"out-of-order apply: peer {peer} was already applied at {ts}"
)),
Some((peer, None)) => Some(format!("out-of-order apply: peer {peer}")),
None => None,
},
_ => None,
}
}
fn note_for_failed_transactional_segment(seg_idx: usize, e: &RunnerError) -> String {
match e {
RunnerError::RelpagesThresholdExceeded {
index_name,
target_table,
relpages,
threshold,
..
} => format!(
"relpages-probe failed at AddIndex {index_name} on table \
{target_table} (relpages={relpages} > threshold={threshold})",
),
RunnerError::TargetTableNotFound {
index_name,
target_table,
..
} => format!(
"relpages-probe failed at AddIndex {index_name}: target table \
`{target_table}` not found and not in plan's AddTable set",
),
RunnerError::TransactionalSegmentFailed {
statement_label,
source,
..
} => format!(
"transactional segment {seg_idx} failed at `{statement_label}`: \
{source}",
),
RunnerError::PkFlipVerificationFailed {
table,
count_violating,
} => format!(
"PK-flip verification halt at segment {seg_idx}: table `{table}` \
has {count_violating} row(s) with NULL or stale shadow values",
),
other => format!("transactional segment {seg_idx} failed: {other}"),
}
}
pub(crate) async fn bind_runner_node_identity(
ctx: &mut DjogiContext,
node_id: i32,
) -> Result<(), RunnerError> {
bind_runner_node_identity_steps(ctx, node_id)
.await
.map_err(|source| RunnerError::NodeIdentityBindingFailed { node_id, source })
}
async fn bind_runner_node_identity_steps(
ctx: &mut DjogiContext,
node_id: i32,
) -> Result<(), DjogiError> {
ctx.execute("SELECT set_heer_node_id($1)", &[&node_id])
.await
.map_err(|source| {
crate::DjogiError::Db(crate::DbError::other(format!(
"bind heer.node_id: {source}"
)))
})?;
ctx.execute("SELECT set_heer_ranj_node_id($1)", &[&node_id])
.await
.map_err(|source| {
crate::DjogiError::Db(crate::DbError::other(format!(
"bind heer.ranj_node_id: {source}"
)))
})?;
Ok(())
}
async fn provision_phase_zero_single_node_dev(
ctx: &mut DjogiContext,
node_id: i32,
) -> Result<(), RunnerError> {
ctx.batch_execute(heeranjid::postgres_schema::SEED_SQL)
.await
.map_err(|source| RunnerError::SingleNodeDevProvisioningFailed {
node_id,
step: "seed node rows",
source,
})?;
let node_seed_sql = super::bootstrap::compose_node_seed("", node_id).map_err(|source| {
RunnerError::SingleNodeDevProvisioningFailed {
node_id,
step: "compose node GUC defaults",
source: DjogiError::Db(DbError::other(source.to_string())),
}
})?;
ctx.batch_execute(&node_seed_sql).await.map_err(|source| {
RunnerError::SingleNodeDevProvisioningFailed {
node_id,
step: "set node GUC defaults",
source,
}
})?;
bind_runner_node_identity_steps(ctx, node_id)
.await
.map_err(|source| RunnerError::SingleNodeDevProvisioningFailed {
node_id,
step: "bind node identity",
source,
})
}
async fn generate_run_id(ctx: &mut DjogiContext, version: &str) -> Result<i64, RunnerError> {
if version == super::bootstrap::PHASE_ZERO_VERSION {
let row = ctx
.__query_one_for_macros(
"SELECT (EXTRACT(EPOCH FROM clock_timestamp()) * 1000000000)::BIGINT AS run_id",
&[],
)
.await
.map_err(|e| RunnerError::RunIdGenerationFailed { source: e })?;
let id: i64 = row
.try_get("run_id")
.map_err(|e| RunnerError::RunIdGenerationFailed {
source: crate::DjogiError::Db(crate::DbError::other(format!("decode run_id: {e}"))),
})?;
return Ok(id);
}
use crate::primary_key::PrimaryKeyDbGen;
let id = HeerId::generate(ctx)
.await
.map_err(|e| RunnerError::RunIdGenerationFailed { source: e })?;
Ok(id.as_i64())
}
fn elapsed_ms(t0: Instant) -> i64 {
t0.elapsed().as_millis().min(i64::MAX as u128) as i64
}
fn is_unique_violation(e: &DjogiError) -> bool {
use tokio_postgres::error::SqlState;
match e {
DjogiError::Db(db) => db_code_matches(db, &SqlState::UNIQUE_VIOLATION),
_ => false,
}
}
fn db_code_matches(db: &DbError, target: &tokio_postgres::error::SqlState) -> bool {
db.code().map(|c| c == target).unwrap_or(false)
}
async fn classify_duplicate_version_collision(
ctx: &mut DjogiContext,
version: &str,
app_label: &str,
) -> RunnerError {
let row = match load_ledger_row_for_version(ctx, version, app_label).await {
Ok(Some(row)) => row,
Ok(None) => {
return RunnerError::LedgerQueryFailed {
query_label: "load_row_for_version",
source: DjogiError::Db(DbError::other(format!(
"duplicate-version collision for `{version}` but \
no ledger row was returned by load_full_row_by_version",
))),
};
}
Err(source) => {
return RunnerError::LedgerQueryFailed {
query_label: "load_row_for_version",
source,
};
}
};
match row.status {
LedgerStatus::Applied | LedgerStatus::Faked | LedgerStatus::Baseline => {
let applied_at = load_applied_at(ctx, version, app_label).await;
RunnerError::VersionAlreadyApplied {
version: version.to_string(),
applied_at,
}
}
LedgerStatus::Pending | LedgerStatus::Failed | LedgerStatus::RolledBack => {
RunnerError::VersionCollisionNonTerminal {
version: row.version,
status: row.status,
run_id: row.run_id,
}
}
}
}
async fn find_higher_applied_version(
ctx: &mut DjogiContext,
bucket: &BucketKey,
candidate_version: &str,
) -> Result<Option<(String, Option<String>)>, DjogiError> {
let row_opt = ctx
.query_opt(
"SELECT version, \
to_char(applied_at AT TIME ZONE 'UTC', \
'YYYY-MM-DD\"T\"HH24:MI:SS\"Z\"') AS applied_at_rfc3339 \
FROM djogi_schema_migrations \
WHERE app_label = $1 \
AND status IN ('applied', 'faked', 'baseline') \
AND version > $2 \
ORDER BY version DESC \
LIMIT 1",
&[&bucket.app, &candidate_version],
)
.await?;
let Some(row) = row_opt else {
return Ok(None);
};
let conflicting_version: String = row.try_get(0)?;
let applied_at_rfc3339: Option<String> = row.try_get(1).ok();
Ok(Some((conflicting_version, applied_at_rfc3339)))
}
async fn bucket_has_applied_history(
ctx: &mut DjogiContext,
bucket: &BucketKey,
) -> Result<bool, DjogiError> {
let row = ctx
.query_one(
"SELECT EXISTS(SELECT 1 \
FROM djogi_schema_migrations \
WHERE app_label = $1 \
AND status IN ('applied', 'faked', 'baseline'))",
&[&bucket.app],
)
.await?;
row.try_get(0).map_err(Into::into)
}
async fn load_applied_at(
ctx: &mut DjogiContext,
version: &str,
app_label: &str,
) -> Option<OffsetDateTime> {
let row = ctx
.query_opt(
"SELECT applied_at FROM djogi_schema_migrations \
WHERE version = $1 AND app_label = $2",
&[&version, &app_label],
)
.await
.ok()??;
row.try_get::<_, OffsetDateTime>("applied_at").ok()
}
fn collect_add_table_targets(plan: &MigrationPlan) -> BTreeSet<String> {
let mut out = BTreeSet::new();
for seg in &plan.segments {
for stmt in &seg.statements {
if let Some(table) = stmt.label.strip_prefix("AddTable ") {
out.insert(table.to_string());
}
}
}
out
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum PartitionExpansionMode {
ApplyLenient,
ReplayStrict,
}
pub(crate) async fn materialize_execution_plan(
ctx: &mut DjogiContext,
plan: &MigrationPlan,
mode: PartitionExpansionMode,
) -> Result<
(
MigrationPlan,
std::collections::HashMap<String, Vec<String>>,
),
RunnerError,
> {
expand_partition_leaf_placeholders(ctx, plan, mode).await
}
pub(crate) fn serialize_leaf_identity(
leaves_cache: &std::collections::HashMap<String, Vec<String>>,
) -> Option<String> {
let mut entries: Vec<_> = leaves_cache
.iter()
.map(|(parent, leaves)| format!("{}:{}", parent, leaves.join(",")))
.collect();
entries.sort();
if entries.is_empty() {
None
} else {
Some(entries.join("\n"))
}
}
pub(crate) async fn compute_leaf_identity_cache(
ctx: &mut DjogiContext,
plan: &MigrationPlan,
) -> Result<std::collections::HashMap<String, Vec<String>>, RunnerError> {
let (_plan, cache) =
expand_partition_leaf_placeholders(ctx, plan, PartitionExpansionMode::ApplyLenient).await?;
Ok(cache)
}
async fn expand_partition_leaf_placeholders(
ctx: &mut DjogiContext,
plan: &MigrationPlan,
mode: PartitionExpansionMode,
) -> Result<
(
MigrationPlan,
std::collections::HashMap<String, Vec<String>>,
),
RunnerError,
> {
let mut leaves_cache: std::collections::HashMap<String, Vec<String>> =
std::collections::HashMap::new();
let mut new_plan = plan.clone();
for segment in &mut new_plan.segments {
let mut new_stmts: Vec<OperationSql> = Vec::with_capacity(segment.statements.len());
for stmt in std::mem::take(&mut segment.statements) {
let parent_for_label = partitioned_parent_from_label(&stmt.label);
match parent_for_label {
Some(parent) => {
if !leaves_cache.contains_key(&parent) {
let leaves = lookup_partition_leaves(ctx, &parent).await?;
leaves_cache.insert(parent.clone(), leaves);
}
let leaves = leaves_cache.get(&parent).expect("just inserted");
new_stmts.extend(expand_partition_statement(&stmt, &parent, leaves, mode)?);
}
None => new_stmts.push(stmt),
}
}
segment.statements = new_stmts;
}
Ok((new_plan, leaves_cache))
}
fn partitioned_parent_from_label(label: &str) -> Option<String> {
for prefix in [
"PkFlipPartitionedBackfill ",
"PkFlipPartitionedIndex ",
"PkFlipPartitionedSelfFkIndex ",
] {
if let Some(rest) = label.strip_prefix(prefix) {
let parent = rest.split_whitespace().next().unwrap_or(rest);
return Some(parent.to_string());
}
}
None
}
async fn lookup_partition_leaves(
ctx: &mut DjogiContext,
parent: &str,
) -> Result<Vec<String>, RunnerError> {
let oid_row = ctx
.query_one("SELECT to_regclass($1)::oid", &[&parent])
.await
.map_err(|e| RunnerError::CatalogQueryFailed {
query_label: "to_regclass",
source: e,
})?;
let oid_opt: Option<u32> = oid_row.try_get(0).ok();
let Some(oid) = oid_opt else {
return Ok(Vec::new());
};
let rows = ctx
.query_all(
"SELECT inhrelid::regclass::text \
FROM pg_inherits \
WHERE inhparent = $1 \
ORDER BY inhrelid::regclass::text",
&[&oid],
)
.await
.map_err(|e| RunnerError::CatalogQueryFailed {
query_label: "pg_inherits",
source: e,
})?;
let mut out: Vec<String> = Vec::with_capacity(rows.len());
for r in &rows {
let leaf: String = r.try_get(0).unwrap_or_default();
if !leaf.is_empty() {
out.push(leaf);
}
}
Ok(out)
}
#[allow(clippy::result_large_err)]
fn expand_partition_statement(
stmt: &OperationSql,
parent: &str,
leaves: &[String],
mode: PartitionExpansionMode,
) -> Result<Vec<OperationSql>, RunnerError> {
if leaves.is_empty() {
return match mode {
PartitionExpansionMode::ApplyLenient => Ok(vec![OperationSql {
label: format!("{} (no leaves)", stmt.label),
up: format!(
"-- pg_inherits returned 0 leaves for partitioned parent {parent}; nothing to expand"
),
down: stmt.down.clone(),
lossy: stmt.lossy.clone(),
}]),
PartitionExpansionMode::ReplayStrict => Err(RunnerError::PartitionExpansionNoLeaves {
parent: parent.to_string(),
statement_label: stmt.label.clone(),
}),
};
}
let mut out: Vec<OperationSql> = Vec::with_capacity(leaves.len());
if stmt.label.starts_with("PkFlipPartitionedBackfill ") {
for leaf in leaves {
let body = stmt.up.replace("<EACH_LEAF_TABLE>", leaf);
let body = strip_trailing_semicolon(&body);
let upper = format!(
"-- partitioned backfill, leaf {leaf}\n{body}",
leaf = leaf,
body = extract_call_line(&body),
);
out.push(OperationSql {
label: format!("PkFlipPartitionedBackfill {parent} leaf={leaf}"),
up: upper,
down: stmt.down.clone(),
lossy: stmt.lossy.clone(),
});
}
return Ok(out);
}
if stmt.label.starts_with("PkFlipPartitionedIndex ") {
let parent_stmt = extract_first_statement_starting_with(&stmt.up, "CREATE UNIQUE INDEX ");
let parent_index_name = recover_parent_index_name(&parent_stmt);
let (part_col, suffix) = recover_partition_columns(&parent_stmt);
out.push(OperationSql {
label: format!("PkFlipPartitionedIndex {parent} (parent-level)"),
up: parent_stmt.clone(),
down: stmt.down.clone(),
lossy: stmt.lossy.clone(),
});
for leaf in leaves {
let (leaf_schema, leaf_bare) = split_leaf_schema(leaf);
let leaf_idx = format!(
"{leaf_bare}_{pkey}_id{suffix}_idx",
pkey = part_col,
suffix = suffix
);
let leaf_idx_qualified = match leaf_schema {
Some(s) => format!("{s}.{leaf_idx}"),
None => leaf_idx.clone(),
};
let (parent_schema, _) = split_leaf_schema(parent);
let parent_idx_qualified = match parent_schema {
Some(s) => format!("{s}.{parent_index_name}"),
None => parent_index_name.clone(),
};
let create_concurrent = format!(
"CREATE UNIQUE INDEX CONCURRENTLY {leaf_idx} ON {leaf} ({pkey}, id{suffix})",
leaf_idx = leaf_idx,
leaf = leaf,
pkey = part_col,
suffix = suffix,
);
out.push(OperationSql {
label: format!("PkFlipPartitionedIndex {parent} leaf={leaf} (concurrent)"),
up: create_concurrent,
down: format!(
"DROP INDEX IF EXISTS {parent_idx_qualified}; DROP INDEX IF EXISTS {leaf_idx_qualified}",
),
lossy: None,
});
let attach = format!(
"ALTER INDEX {parent_idx_qualified} ATTACH PARTITION {leaf_idx_qualified}",
);
out.push(OperationSql {
label: format!("PkFlipPartitionedIndex {parent} leaf={leaf} (attach)"),
up: attach,
down: String::new(),
lossy: None,
});
}
return Ok(out);
}
if stmt.label.starts_with("PkFlipPartitionedSelfFkIndex ") {
let parent_stmt = extract_first_statement_starting_with(&stmt.up, "CREATE INDEX ");
let parent_index_name = recover_parent_index_name(&parent_stmt);
let (col, suffix) = recover_self_fk_column(&parent_stmt);
out.push(OperationSql {
label: format!("PkFlipPartitionedSelfFkIndex {parent} (parent-level)"),
up: parent_stmt.clone(),
down: stmt.down.clone(),
lossy: stmt.lossy.clone(),
});
for leaf in leaves {
let (leaf_schema, leaf_bare) = split_leaf_schema(leaf);
let leaf_idx = format!("{leaf_bare}_{col}{suffix}_idx", col = col, suffix = suffix);
let leaf_idx_qualified = match leaf_schema {
Some(s) => format!("{s}.{leaf_idx}"),
None => leaf_idx.clone(),
};
let (parent_schema, _) = split_leaf_schema(parent);
let parent_idx_qualified = match parent_schema {
Some(s) => format!("{s}.{parent_index_name}"),
None => parent_index_name.clone(),
};
let create_concurrent = format!(
"CREATE INDEX CONCURRENTLY {leaf_idx} ON {leaf} ({col}{suffix})",
leaf_idx = leaf_idx,
leaf = leaf,
col = col,
suffix = suffix,
);
out.push(OperationSql {
label: format!("PkFlipPartitionedSelfFkIndex {parent} leaf={leaf} (concurrent)"),
up: create_concurrent,
down: format!("DROP INDEX IF EXISTS {leaf_idx_qualified}"),
lossy: None,
});
let attach = format!(
"ALTER INDEX {parent_idx_qualified} ATTACH PARTITION {leaf_idx_qualified}",
);
out.push(OperationSql {
label: format!("PkFlipPartitionedSelfFkIndex {parent} leaf={leaf} (attach)"),
up: attach,
down: String::new(),
lossy: None,
});
}
return Ok(out);
}
Ok(vec![stmt.clone()])
}
fn strip_trailing_semicolon(s: &str) -> String {
let trimmed = s.trim_end();
trimmed.strip_suffix(';').unwrap_or(trimmed).to_string()
}
fn split_leaf_schema(qualified: &str) -> (Option<&str>, &str) {
let bytes = qualified.as_bytes();
let mut in_quote = false;
let mut sep = None;
for (i, &b) in bytes.iter().enumerate() {
match b {
b'"' => in_quote = !in_quote,
b'.' if !in_quote => sep = Some(i),
_ => {}
}
}
match sep {
Some(i) => (Some(&qualified[..i]), &qualified[i + 1..]),
None => (None, qualified),
}
}
fn extract_call_line(s: &str) -> String {
for line in s.lines() {
if line.contains("CALL heeranjid_bulk_backfill(") {
return strip_trailing_semicolon(line);
}
}
strip_trailing_semicolon(s)
}
fn extract_first_statement_starting_with(body: &str, start_marker: &str) -> String {
let bytes = body.as_bytes();
let mut i = 0usize;
while i < bytes.len() {
while i < bytes.len() && bytes[i] == b' ' {
i += 1;
}
if i + start_marker.len() <= bytes.len()
&& &bytes[i..i + start_marker.len()] == start_marker.as_bytes()
{
let start = i;
let mut j = i + start_marker.len();
while j < bytes.len() && bytes[j] != b';' {
j += 1;
}
let raw = String::from_utf8_lossy(&bytes[start..j]).into_owned();
return raw;
}
while i < bytes.len() && bytes[i] != b'\n' {
i += 1;
}
if i < bytes.len() {
i += 1;
}
}
String::new()
}
fn recover_parent_index_name(parent_stmt: &str) -> String {
let bytes = parent_stmt.as_bytes();
for marker in [b"CREATE UNIQUE INDEX " as &[u8], b"CREATE INDEX " as &[u8]] {
if bytes.len() < marker.len() {
continue;
}
let mut i = 0usize;
while i + marker.len() <= bytes.len() {
if &bytes[i..i + marker.len()] == marker {
let start = i + marker.len();
let mut name_start = start;
while name_start < bytes.len() && bytes[name_start] == b' ' {
name_start += 1;
}
const CONC: &[u8] = b"CONCURRENTLY";
if name_start + CONC.len() <= bytes.len()
&& &bytes[name_start..name_start + CONC.len()] == CONC
{
name_start += CONC.len();
while name_start < bytes.len() && bytes[name_start] == b' ' {
name_start += 1;
}
}
let mut j = name_start;
while j < bytes.len() && (bytes[j].is_ascii_alphanumeric() || bytes[j] == b'_') {
j += 1;
}
if j > name_start {
return String::from_utf8_lossy(&bytes[name_start..j]).into_owned();
}
break;
}
i += 1;
}
}
String::new()
}
fn recover_partition_columns(parent_stmt: &str) -> (String, String) {
let bytes = parent_stmt.as_bytes();
let mut open = None;
let mut close = None;
for (i, b) in bytes.iter().enumerate() {
if *b == b'(' && open.is_none() {
open = Some(i + 1);
} else if *b == b')' {
close = Some(i);
break;
}
}
let (Some(o), Some(c)) = (open, close) else {
return ("partition_key".to_string(), "_desc".to_string());
};
if o >= c {
return ("partition_key".to_string(), "_desc".to_string());
}
let inside = String::from_utf8_lossy(&bytes[o..c]).into_owned();
let mut parts = inside.splitn(2, ',');
let pkey = parts.next().unwrap_or("").trim().to_string();
let id_col = parts.next().unwrap_or("").trim();
let suffix = match id_col.strip_prefix("id") {
Some(s) => s.to_string(),
None => return ("partition_key".to_string(), "_desc".to_string()),
};
if pkey.is_empty() {
return ("partition_key".to_string(), "_desc".to_string());
}
(pkey, suffix)
}
fn recover_self_fk_column(parent_stmt: &str) -> (String, String) {
let bytes = parent_stmt.as_bytes();
let mut open = None;
let mut close = None;
for (i, b) in bytes.iter().enumerate() {
if *b == b'(' && open.is_none() {
open = Some(i + 1);
} else if *b == b')' {
close = Some(i);
break;
}
}
let (Some(o), Some(c)) = (open, close) else {
return ("col".to_string(), "_desc".to_string());
};
if o >= c {
return ("col".to_string(), "_desc".to_string());
}
let inside = parent_stmt[o..c].trim();
if let Some(col) = inside.strip_suffix("_desc")
&& !col.is_empty()
{
return (col.to_string(), "_desc".to_string());
}
("col".to_string(), "_desc".to_string())
}
#[cfg(test)]
mod tests {
#![allow(clippy::await_holding_lock)]
use super::*;
use std::collections::{BTreeMap, BTreeSet};
use std::path::PathBuf;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use crate::config::MigrateConfig;
use crate::migrate::bootstrap::PHASE_ZERO_VERSION;
use crate::migrate::diff::Classification;
use crate::migrate::projection::BucketKey;
use crate::migrate::schema::AppliedSchema;
use crate::migrate::segment::{MigrationPlan, Segment, SegmentKind};
use crate::migrate::sql::OperationSql;
use djogi_macros::djogi_test;
fn bucket(db: &str, app: &str) -> BucketKey {
BucketKey {
database: db.to_string(),
app: app.to_string(),
}
}
fn empty_snapshot() -> AppliedSchema {
AppliedSchema {
djogi_version: "0.1.0".to_string(),
enums: BTreeMap::new(),
format_version: SNAPSHOT_FORMAT_VERSION.to_string(),
generated_at: "2026-05-09T00:00:00Z".to_string(),
indexes: Vec::new(),
models: BTreeMap::new(),
registered_apps: vec!["".to_string()],
}
}
fn op(label: &str, up: &str) -> OperationSql {
OperationSql {
label: label.to_string(),
up: up.to_string(),
down: format!("-- down for {label}"),
lossy: None,
}
}
fn audit_plan() -> MigrationPlan {
MigrationPlan {
bucket: bucket("main", ""),
classification: Classification::Additive,
segments: vec![
Segment {
kind: SegmentKind::Transactional,
statements: vec![
op("AddTable audit_a", "CREATE TABLE audit_a (id bigint)"),
op("AddTable audit_b", "CREATE TABLE audit_b (id bigint)"),
],
},
Segment {
kind: SegmentKind::MetadataOnly,
statements: vec![op("RenameApp ignored", "-- metadata-only placeholder")],
},
Segment {
kind: SegmentKind::NonTransactional,
statements: vec![op(
"AddIndex audit_a_id_idx",
"CREATE INDEX CONCURRENTLY audit_a_id_idx ON audit_a (id)",
)],
},
],
}
}
fn single_table_plan(table: &str) -> MigrationPlan {
MigrationPlan {
bucket: bucket("main", ""),
classification: Classification::Additive,
segments: vec![Segment {
kind: SegmentKind::Transactional,
statements: vec![op(
&format!("AddTable {table}"),
&format!("CREATE TABLE {table} (id bigint)"),
)],
}],
}
}
fn single_segment_plan(kind: SegmentKind, label: &str, up: &str) -> MigrationPlan {
MigrationPlan {
bucket: bucket("main", ""),
classification: Classification::Additive,
segments: vec![Segment {
kind,
statements: vec![op(label, up)],
}],
}
}
fn runner_ctx_for_audit_with_snapshot_path(
plan: &MigrationPlan,
audit_pool: Option<deadpool_postgres::Pool>,
snapshot_path: Option<PathBuf>,
) -> RunnerCtx {
RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260509000000__audit_test".to_string(),
description: "audit test".to_string(),
checksum_up: compute_checksum_for_plan_up(plan),
checksum_down: None,
snapshot: Some(empty_snapshot()),
snapshot_path,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool,
runner_identity: Some(RunnerIdentity::SingleNodeDev), drift_baseline: DriftBaseline::Disabled,
}
}
fn runner_ctx_with_drift_baseline(
plan: &MigrationPlan,
version: &str,
drift_baseline: DriftBaseline,
) -> RunnerCtx {
RunnerCtx {
bucket: plan.bucket.clone(),
version: version.to_string(),
description: "drift gate test".to_string(),
checksum_up: compute_checksum_for_plan_up(plan),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: Some(RunnerIdentity::SingleNodeDev),
drift_baseline,
}
}
fn named_bucket_single_table_plan(db: &str, app: &str, table: &str) -> MigrationPlan {
MigrationPlan {
bucket: bucket(db, app),
classification: Classification::Additive,
segments: vec![Segment {
kind: SegmentKind::Transactional,
statements: vec![op(
&format!("AddTable {table}"),
&format!("CREATE TABLE {table} (id bigint)"),
)],
}],
}
}
fn drift_baseline_with_table(table: &str) -> AppliedSchema {
use crate::migrate::schema::{ColumnSchema, PkKindSchema, PrimaryKeySchema, TableSchema};
let id_col = ColumnSchema {
check: None,
codec: None,
comment: None,
default_sql: None,
foreign_key: None,
generated: None,
identity: None,
index_type: None,
indexed: false,
max_length: None,
name: "id".to_string(),
nullable: false,
on_delete: None,
outbox_exclude: false,
rationale: None,
relation_kind: None,
renamed_from: None,
sequence_within: None,
sql_type: "BIGINT".to_string(),
unique: false,
type_change_using: None,
};
let table_schema = TableSchema {
app: None,
columns: vec![id_col],
exclusion_constraints: Vec::new(),
fts: None,
is_through: false,
moved_from_app: None,
partition: None,
primary_key: PrimaryKeySchema {
columns: vec!["id".to_string()],
kind: PkKindSchema::HeerId,
},
rationale: None,
renamed_from: None,
rls_enabled: false,
storage_params: None,
table: table.to_string(),
table_comment: None,
tablespace: None,
tenant_key: None,
};
let mut snapshot = empty_snapshot();
snapshot.models.insert(table.to_string(), table_schema);
snapshot
}
fn unique_temp_path(tag: &str, ext: &str) -> PathBuf {
let stamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
std::env::temp_dir().join(format!("djogi-runner-{tag}-{stamp}.{ext}"))
}
fn acquire_test_workspace_guard() -> WorkspaceGuard {
crate::migrate::acquire_workspace_lock(
&unique_temp_path("audit", "lock"),
Duration::from_secs(2),
)
.expect("acquire workspace lock")
}
async fn unreachable_pool_context() -> DjogiContext {
let pool = crate::pg::pool::DjogiPool::builder("postgres://localhost:1/djogi_unreachable")
.timeout(Duration::from_millis(1))
.build()
.await
.expect("construct lazy unreachable pool");
DjogiContext::from_pool(pool)
}
fn generated_stale_phase_zero_plan() -> MigrationPlan {
let mut sql = crate::migrate::bootstrap::compose_phase_zero(
"main",
&BTreeSet::new(),
crate::migrate::bootstrap::DEFAULT_NODE_ID,
true,
)
.expect("compose seed-capable Phase 0");
let start = sql.find("DO $djogi$").expect("dynamic defaults block");
let end = sql[start..]
.find("SET heer.node_id = '1';")
.map(|offset| start + offset)
.expect("session SET after dynamic defaults");
sql.replace_range(
start..end,
"ALTER DATABASE \"main\" SET heer.node_id = '1';\n\
ALTER DATABASE \"main\" SET heer.ranj_node_id = '1';\n",
);
single_segment_plan(SegmentKind::Transactional, "Phase 0 bootstrap", &sql)
}
fn current_phase_zero_plan() -> MigrationPlan {
let sql = crate::migrate::bootstrap::compose_phase_zero(
"main",
&BTreeSet::new(),
crate::migrate::bootstrap::DEFAULT_NODE_ID,
false,
)
.expect("compose identity-free Phase 0");
single_segment_plan(SegmentKind::Transactional, "Phase 0 bootstrap", &sql)
}
fn markerless_seed_phase_zero_plan() -> MigrationPlan {
phase_zero_plan_with_seed_statement("INSERT INTO heer.heer_nodes (id) VALUES (1);")
}
fn phase_zero_plan_with_seed_statement(statement: &str) -> MigrationPlan {
let mut sql = crate::migrate::bootstrap::compose_phase_zero(
"main",
&BTreeSet::new(),
crate::migrate::bootstrap::DEFAULT_NODE_ID,
false,
)
.expect("compose identity-free Phase 0");
sql.push('\n');
sql.push_str(statement);
sql.push('\n');
single_segment_plan(SegmentKind::Transactional, "Phase 0 bootstrap", &sql)
}
fn extended_seed_statement_cases() -> [(&'static str, &'static str); 4] {
[
(
"cte_insert",
"WITH rows AS (SELECT 1) INSERT INTO heer.heer_nodes (id) VALUES (1);",
),
(
"cte_delete",
"WITH moved AS (DELETE FROM heer.heer_node_state RETURNING *) SELECT 1;",
),
(
"merge",
"MERGE INTO heer.heer_nodes AS target USING incoming ON false WHEN NOT MATCHED THEN INSERT (id) VALUES (1);",
),
(
"copy_from",
"COPY \"heer\".\"heer_ranj_node_state\" (\"node_id\") FROM STDIN;",
),
]
}
fn seed_capable_phase_zero_plan() -> MigrationPlan {
let sql = crate::migrate::bootstrap::compose_phase_zero(
"main",
&BTreeSet::new(),
crate::migrate::bootstrap::DEFAULT_NODE_ID,
true,
)
.expect("compose seed-capable Phase 0");
single_segment_plan(SegmentKind::Transactional, "Phase 0 bootstrap", &sql)
}
struct SigningKeyEnvUnsetGuard {
previous: Option<std::ffi::OsString>,
_guard: std::sync::MutexGuard<'static, ()>,
}
struct SigningKeyEnvReadGuard {
_guard: std::sync::MutexGuard<'static, ()>,
}
impl SigningKeyEnvReadGuard {
fn hold() -> Self {
Self {
_guard: crate::snapshot::sign::SIGNING_KEY_ENV_MUTEX
.lock()
.expect("signing-key env mutex"),
}
}
}
impl SigningKeyEnvUnsetGuard {
fn unset() -> Self {
let guard = crate::snapshot::sign::SIGNING_KEY_ENV_MUTEX
.lock()
.expect("signing-key env mutex");
let previous = std::env::var_os("DJOGI_SNAPSHOT_SIGNING_KEY");
unsafe {
std::env::remove_var("DJOGI_SNAPSHOT_SIGNING_KEY");
}
Self {
previous,
_guard: guard,
}
}
}
impl Drop for SigningKeyEnvUnsetGuard {
fn drop(&mut self) {
unsafe {
if let Some(previous) = &self.previous {
std::env::set_var("DJOGI_SNAPSHOT_SIGNING_KEY", previous);
} else {
std::env::remove_var("DJOGI_SNAPSHOT_SIGNING_KEY");
}
}
}
}
fn db_err(message: &'static str) -> DjogiError {
DjogiError::Db(DbError::other(message))
}
#[test]
fn runner_error_operator_actionable_variants() {
let cases = [
(
RunnerError::ChecksumMismatch(djogi::migrate::ledger::ChecksumMismatch {
version: "V20260101000000__add_users".to_string(),
expected:
"V1:11111111111111111111111111111111111111111111111111111111111111111111"
.to_string(),
actual:
"V1:22222222222222222222222222222222222222222222222222222222222222222222"
.to_string(),
}),
true,
),
(
RunnerError::ChecksumFormat(djogi::migrate::ledger::ChecksumFormatError {
side: djogi::migrate::ledger::ChecksumSide::Expected,
value: "V2:deadbeef".to_string(),
kind: djogi::migrate::ledger::ChecksumFormatErrorKind::WrongPrefix,
}),
true,
),
(
RunnerError::VersionAlreadyApplied {
version: "V20260101000000__add_users".to_string(),
applied_at: None,
},
true,
),
(
RunnerError::VersionCollisionNonTerminal {
version: "V20260101000000__add_users".to_string(),
status: LedgerStatus::Pending,
run_id: 1,
},
true,
),
(
RunnerError::TargetTableNotFound {
bucket: bucket("main", ""),
index_name: "users_email_idx".to_string(),
target_table: "users".to_string(),
},
true,
),
(
RunnerError::RelpagesThresholdExceeded {
bucket: bucket("main", ""),
index_name: "users_email_idx".to_string(),
target_table: "users".to_string(),
relpages: 4096,
threshold: 128,
},
true,
),
(
RunnerError::SegmentSqlExecutionModeConflict {
segment_index: 0,
segment_kind: SegmentKind::Transactional,
statement_label: "CREATE TRIGGER users_idx".to_string(),
problem: SegmentSqlExecutionModeProblem::RequiresNonTransactional {
statement_shape: "requires non-transactional execution",
},
},
true,
),
(
RunnerError::SnapshotPersistFailed {
path: "schema_snapshot.json".into(),
source: SnapshotError::Io {
path: None,
source: std::io::Error::other("disk full"),
},
},
true,
),
(RunnerError::BaselineSnapshotShouldNotBeProvided, true),
(
RunnerError::StalePhaseZeroArtifact {
version: "V00000000000000__phase_zero".to_string(),
refusal_reason: "generated-stale",
},
true,
),
(
RunnerError::DriftDetected {
bucket: bucket("main", "billing"),
report: crate::migrate::VerifyReport {
diagnostics: vec![crate::migrate::VerifyDiagnostic {
code: "D601".to_string(),
severity: crate::migrate::VerifySeverity::Error,
message: "Snapshot table missing from live DB".to_string(),
location: Some("billing.invoices".to_string()),
}],
latest_applied_version: Some("V20260601000000__billing".to_string()),
applied_count: 2,
unfinished_count: 0,
},
},
true,
),
(
RunnerError::DriftBaselineMissing {
bucket: bucket("main", "billing"),
},
true,
),
(
RunnerError::DriftBaselineCorrupted {
bucket: bucket("main", "billing"),
reason: String::new(),
},
true,
),
(
RunnerError::OutOfOrderRejected {
version: "V20260101000000__add_users".to_string(),
conflicting_version: "V20260102000000__add_orders".to_string(),
conflicting_applied_at: None,
},
true,
),
(
RunnerError::PkFlipHazardReplicaSessions {
walsenders: vec![("wal_sender".to_string(), "127.0.0.1:5432".to_string())],
subscriptions: vec!["replication_slot".to_string()],
},
true,
),
(
RunnerError::PkFlipHazardPreexistingZzzTrigger {
table: "public.events".to_string(),
trigger_names: vec!["zzz_rv_users_id".to_string()],
},
true,
),
(
RunnerError::PkFlipHazardDisabledTriggers {
table: "public.events".to_string(),
triggers: vec![("zzz_rv_events_id".to_string(), 'D')],
},
true,
),
(
RunnerError::PkFlipHazardLongRunningTx {
offenders: vec![(1234, 120)],
threshold_secs: 60,
},
true,
),
(
RunnerError::PkFlipVerificationFailed {
table: "public.events".to_string(),
count_violating: 2,
},
true,
),
(
RunnerError::AdvisoryUnlockReturnedFalse {
key: 0x0102_0304_0506_0708,
bucket: bucket("main", ""),
},
true,
),
(
RunnerError::PartitionExpansionNoLeaves {
parent: "events".to_string(),
statement_label: "migrate_partitions".to_string(),
},
true,
),
(
RunnerError::MissingApplyIdentity {
version: "V20260101000000__add_users".to_string(),
},
true,
),
];
for (err, expected) in cases {
assert_eq!(err.is_operator_actionable(), expected, "{err}");
}
}
#[test]
fn runner_error_non_operator_actionable_variants() {
let cases = [
(
RunnerError::LockTimeout {
path: PathBuf::from("/tmp/.djogi-migrations-lock"),
holder_pid: Some(99),
},
false,
),
(
RunnerError::GuardError(crate::migrate::guard::GuardError::Io {
path: PathBuf::from("/tmp/.djogi-migrations-lock"),
source: std::io::Error::other("permission denied"),
}),
false,
),
(
RunnerError::AdvisoryLockFailed {
bucket: bucket("main", ""),
key: 0x0000_0001_0000_0001,
attempts: 3,
},
false,
),
(
RunnerError::AdvisoryLockQueryFailed {
app_label: "main".to_string(),
source: db_err("advisory lock query failed"),
},
false,
),
(
RunnerError::LedgerWriteFailed {
version: "V20260101000000__add_users".to_string(),
source: db_err("insert failed"),
},
false,
),
(
RunnerError::LedgerQueryFailed {
query_label: "load_row_for_version",
source: db_err("query failed"),
},
false,
),
(
RunnerError::RunIdGenerationFailed {
source: db_err("sequence exhausted"),
},
false,
),
(
RunnerError::NodeIdentityBindingFailed {
node_id: 1,
source: db_err("identity bind failed"),
},
false,
),
(
RunnerError::SingleNodeDevProvisioningFailed {
node_id: 1,
step: "register node",
source: db_err("provisioning failed"),
},
false,
),
(
RunnerError::LedgerBootstrapFailed {
source: db_err("create table failed"),
},
false,
),
(
RunnerError::TransactionalSegmentFailed {
segment_index: 0,
statement_label: "AddTable users".to_string(),
source: db_err("transaction failed"),
},
false,
),
(
RunnerError::NonTransactionalSegmentFailed {
segment_index: 0,
step_index: 0,
statement_label: "INSERT user".to_string(),
applied_steps_count: 0,
source: db_err("statement failed"),
},
false,
),
(
RunnerError::NonTransactionalProgressAckFailed {
segment_index: 0,
step_index: 0,
statement_label: "INSERT user".to_string(),
applied_steps_count: 0,
source: db_err("ack failed"),
},
false,
),
(
RunnerError::ConfigLoadFailed {
source: figment::Error::from("invalid config"),
},
false,
),
(
RunnerError::BaselineProjectionFailed {
source: Box::new(djogi::migrate::verify::VerifyRunError::CatalogQueryFailed {
query_label: "pg_namespace",
source: db_err("verify catalog query failed"),
}),
},
false,
),
(
RunnerError::DriftPreflightFailed {
source: Box::new(djogi::migrate::verify::VerifyRunError::CatalogQueryFailed {
query_label: "columns",
source: db_err("drift preflight failed"),
}),
},
false,
),
(
RunnerError::CatalogQueryFailed {
query_label: "pg_class",
source: db_err("catalog query failed"),
},
false,
),
(
RunnerError::PinnedSessionCheckoutFailed {
source: db_err("pool checkout failed"),
},
false,
),
];
for (err, expected) in cases {
assert_eq!(err.is_operator_actionable(), expected, "{err}");
}
}
#[test]
fn drift_detected_display_names_bucket_and_error_count_and_next_step() {
let e = RunnerError::DriftDetected {
bucket: bucket("main", "billing"),
report: crate::migrate::verify::VerifyReport {
diagnostics: vec![crate::migrate::verify::VerifyDiagnostic {
code: "D601".to_string(),
severity: crate::migrate::verify::VerifySeverity::Error,
message: "snapshot table `users` missing in live".to_string(),
location: Some("users".to_string()),
}],
latest_applied_version: Some("V20260101000000__x".to_string()),
applied_count: 3,
unfinished_count: 0,
},
};
let msg = e.to_string();
assert!(msg.contains("database=main"), "got: {msg}");
assert!(msg.contains("app=billing"), "got: {msg}");
assert!(msg.contains("1 error-severity"), "got: {msg}");
assert!(msg.contains("djogi migrations verify"), "got: {msg}");
assert!(msg.contains("djogi migrations attune"), "got: {msg}");
assert!(std::error::Error::source(&e).is_none());
}
#[test]
fn drift_preflight_failed_exposes_verify_error_as_source() {
let e = RunnerError::DriftPreflightFailed {
source: Box::new(crate::migrate::verify::VerifyRunError::LedgerQueryFailed {
source: db_err("ledger read failed in test"),
}),
};
assert!(e.to_string().contains("drift pre-flight could not run"));
assert!(std::error::Error::source(&e).is_some());
}
#[test]
fn drift_baseline_missing_display_names_bucket_and_recovery_commands() {
let e = RunnerError::DriftBaselineMissing {
bucket: bucket("main", "billing"),
};
let msg = e.to_string();
assert!(msg.contains("database=main"), "got: {msg}");
assert!(msg.contains("app=billing"), "got: {msg}");
assert!(msg.contains("applied migration history"), "got: {msg}");
assert!(msg.contains("schema_snapshot.json"), "got: {msg}");
assert!(
msg.contains("djogi migrations repair snapshot-rebuild"),
"got: {msg}"
);
assert!(std::error::Error::source(&e).is_none());
}
#[test]
fn drift_baseline_corrupted_display_names_bucket_and_recovery_commands() {
let e = RunnerError::DriftBaselineCorrupted {
bucket: bucket("main", "billing"),
reason: "unexpected end of input".to_string(),
};
let msg = e.to_string();
assert!(msg.contains("database=main"), "got: {msg}");
assert!(msg.contains("app=billing"), "got: {msg}");
assert!(msg.contains("schema_snapshot.json"), "got: {msg}");
assert!(msg.contains("repair snapshot-rebuild"), "got: {msg}");
assert!(std::error::Error::source(&e).is_none());
}
#[test]
fn advisory_lock_key_is_deterministic_across_calls() {
let b = bucket("main", "");
let a = advisory_lock_key(&b);
let c = advisory_lock_key(&b);
assert_eq!(a, c, "same input must yield same lock key");
}
#[test]
fn advisory_lock_key_differs_on_database() {
let a = advisory_lock_key(&bucket("alpha", ""));
let b = advisory_lock_key(&bucket("beta", ""));
assert_ne!(a, b, "different database must yield different key");
}
#[test]
fn advisory_lock_key_differs_on_app() {
let a = advisory_lock_key(&bucket("main", "users"));
let b = advisory_lock_key(&bucket("main", "billing"));
assert_ne!(a, b, "different app must yield different key");
}
#[test]
fn advisory_lock_key_database_app_separator_prevents_collision() {
let a = advisory_lock_key(&bucket("ab", "c"));
let b = advisory_lock_key(&bucket("a", "bc"));
assert_ne!(a, b);
}
#[test]
fn advisory_lock_key_pins_big_endian_byte_decode() {
let bk = bucket("test_db", "test_app");
let key = advisory_lock_key(&bk);
assert_eq!(
key, 9_007_707_844_108_204_599_i64,
"advisory_lock_key must decode the first 8 SHA-256 bytes as big-endian"
);
assert_ne!(
key, 4_015_681_655_511_318_909_i64,
"advisory_lock_key must NOT decode bytes as little-endian"
);
}
#[test]
fn parse_index_extracts_name_and_table() {
let stmt = OperationSql {
label: "AddIndex users_email_idx".to_string(),
up: "CREATE INDEX \"users_email_idx\" ON \"users\" (\"email\")".to_string(),
down: String::new(),
lossy: None,
};
let parsed = parse_create_index_statement(&stmt).expect("parse");
assert_eq!(parsed.0, "users_email_idx");
assert_eq!(parsed.1, "users");
}
#[test]
fn parse_index_returns_none_for_non_index_label() {
let stmt = OperationSql {
label: "AddTable users".to_string(),
up: "CREATE TABLE \"users\" ()".to_string(),
down: String::new(),
lossy: None,
};
assert!(parse_create_index_statement(&stmt).is_none());
}
#[test]
fn parse_index_returns_none_when_marker_missing() {
let stmt = OperationSql {
label: "AddIndex weird_idx".to_string(),
up: "CREATE INDEX weird".to_string(),
down: String::new(),
lossy: None,
};
assert!(parse_create_index_statement(&stmt).is_none());
}
#[test]
fn plan_checksum_matches_manual_concatenation() {
let plan = MigrationPlan {
bucket: bucket("main", ""),
classification: Classification::Additive,
segments: vec![Segment {
kind: SegmentKind::Transactional,
statements: vec![
OperationSql {
label: "AddTable a".to_string(),
up: "CREATE TABLE a ()".to_string(),
down: "DROP TABLE a".to_string(),
lossy: None,
},
OperationSql {
label: "AddTable b".to_string(),
up: "CREATE TABLE b ()".to_string(),
down: "DROP TABLE b".to_string(),
lossy: None,
},
],
}],
};
let computed = compute_checksum_for_plan_up(&plan);
let manual = compute_checksum(["CREATE TABLE a ()", "CREATE TABLE b ()"]);
assert_eq!(computed, manual);
}
#[test]
fn plan_checksum_changes_on_segment_reorder() {
let make = |labels: [&str; 2]| MigrationPlan {
bucket: bucket("main", ""),
classification: Classification::Additive,
segments: vec![Segment {
kind: SegmentKind::Transactional,
statements: labels
.iter()
.map(|l| OperationSql {
label: format!("AddTable {l}"),
up: format!("CREATE TABLE {l} ()"),
down: format!("DROP TABLE {l}"),
lossy: None,
})
.collect(),
}],
};
let a = compute_checksum_for_plan_up(&make(["a", "b"]));
let b = compute_checksum_for_plan_up(&make(["b", "a"]));
assert_ne!(a, b);
}
#[test]
fn elapsed_ms_returns_non_negative() {
let t = Instant::now();
let ms = elapsed_ms(t);
assert!(ms >= 0);
}
#[test]
fn heer_id_zero_converts_to_zero_i64_directly() {
let z = HeerId::ZERO;
assert_eq!(z.as_i64(), 0);
let via_from: i64 = i64::from(z);
assert_eq!(via_from, 0);
}
#[test]
fn heer_id_non_zero_round_trips_directly() {
let v: i64 = 0x0123_4567_89AB_CDEF_i64;
let id = HeerId::try_from(v).expect("positive i64 round-trips through HeerId");
assert_eq!(id.as_i64(), v);
let via_from: i64 = i64::from(id);
assert_eq!(via_from, v);
}
#[test]
fn collect_add_table_targets_walks_all_segments() {
let plan = MigrationPlan {
bucket: bucket("main", ""),
classification: Classification::Additive,
segments: vec![
Segment {
kind: SegmentKind::Transactional,
statements: vec![
OperationSql {
label: "AddTable users".to_string(),
up: "CREATE TABLE users ()".to_string(),
down: "DROP TABLE users".to_string(),
lossy: None,
},
OperationSql {
label: "AddIndex users_email_idx".to_string(),
up: "CREATE INDEX...".to_string(),
down: String::new(),
lossy: None,
},
],
},
Segment {
kind: SegmentKind::Transactional,
statements: vec![OperationSql {
label: "AddTable orders".to_string(),
up: "CREATE TABLE orders ()".to_string(),
down: "DROP TABLE orders".to_string(),
lossy: None,
}],
},
],
};
let set = collect_add_table_targets(&plan);
assert!(set.contains("users"));
assert!(set.contains("orders"));
assert!(!set.contains("users_email_idx")); assert_eq!(set.len(), 2);
}
#[test]
fn collect_add_table_targets_empty_plan_returns_empty_set() {
let plan = MigrationPlan {
bucket: bucket("main", ""),
classification: Classification::NoOp,
segments: vec![],
};
let set = collect_add_table_targets(&plan);
assert!(set.is_empty());
}
#[test]
fn segment_sql_preflight_rejects_concurrent_index_in_transactional_segment() {
let plan = single_segment_plan(
SegmentKind::Transactional,
"AddIndex users_email_idx",
"CREATE INDEX CONCURRENTLY users_email_idx ON users (email);",
);
let err = preflight_segment_sql_execution_compatibility(&plan).expect_err("must reject");
match err {
RunnerError::SegmentSqlExecutionModeConflict {
segment_index,
segment_kind,
statement_label,
problem,
} => {
assert_eq!(segment_index, 0);
assert_eq!(segment_kind, SegmentKind::Transactional);
assert_eq!(statement_label, "AddIndex users_email_idx");
assert_eq!(
problem,
SegmentSqlExecutionModeProblem::RequiresNonTransactional {
statement_shape: "CREATE INDEX CONCURRENTLY",
}
);
}
other => panic!("expected SegmentSqlExecutionModeConflict, got {other:?}"),
}
}
#[test]
fn segment_sql_preflight_rejects_comment_separated_concurrent_index_keywords() {
let plan = single_segment_plan(
SegmentKind::Transactional,
"AddIndex users_email_idx",
"CREATE /* split */ UNIQUE INDEX /* still split */ CONCURRENTLY \
users_email_idx ON users (email);",
);
let err = preflight_segment_sql_execution_compatibility(&plan).expect_err("must reject");
match err {
RunnerError::SegmentSqlExecutionModeConflict {
segment_index,
segment_kind,
statement_label,
problem,
} => {
assert_eq!(segment_index, 0);
assert_eq!(segment_kind, SegmentKind::Transactional);
assert_eq!(statement_label, "AddIndex users_email_idx");
assert_eq!(
problem,
SegmentSqlExecutionModeProblem::RequiresNonTransactional {
statement_shape: "CREATE UNIQUE INDEX CONCURRENTLY",
}
);
}
other => panic!("expected SegmentSqlExecutionModeConflict, got {other:?}"),
}
}
#[test]
fn segment_sql_preflight_rejects_begin_in_transactional_segment() {
let plan = single_segment_plan(SegmentKind::Transactional, "manual begin", "BEGIN;");
let err = preflight_segment_sql_execution_compatibility(&plan).expect_err("must reject");
match err {
RunnerError::SegmentSqlExecutionModeConflict {
segment_index,
segment_kind,
statement_label,
problem,
} => {
assert_eq!(segment_index, 0);
assert_eq!(segment_kind, SegmentKind::Transactional);
assert_eq!(statement_label, "manual begin");
assert_eq!(
problem,
SegmentSqlExecutionModeProblem::TransactionControl { keyword: "BEGIN" }
);
}
other => panic!("expected SegmentSqlExecutionModeConflict, got {other:?}"),
}
}
#[test]
fn segment_sql_preflight_rejects_savepoint_in_non_transactional_segment() {
let plan = single_segment_plan(
SegmentKind::NonTransactional,
"manual savepoint",
"SAVEPOINT retry_guard;",
);
let err = preflight_segment_sql_execution_compatibility(&plan).expect_err("must reject");
match err {
RunnerError::SegmentSqlExecutionModeConflict {
segment_index,
segment_kind,
statement_label,
problem,
} => {
assert_eq!(segment_index, 0);
assert_eq!(segment_kind, SegmentKind::NonTransactional);
assert_eq!(statement_label, "manual savepoint");
assert_eq!(
problem,
SegmentSqlExecutionModeProblem::TransactionControl {
keyword: "SAVEPOINT",
}
);
}
other => panic!("expected SegmentSqlExecutionModeConflict, got {other:?}"),
}
}
#[test]
fn segment_sql_preflight_allows_set_constraints_in_transactional_segment() {
let plan = single_segment_plan(
SegmentKind::Transactional,
"defer constraints",
"SET CONSTRAINTS ALL DEFERRED;",
);
preflight_segment_sql_execution_compatibility(&plan).expect("set constraints allowed");
}
#[test]
fn segment_sql_preflight_ignores_function_body_begin_commit_tokens() {
let plan = single_segment_plan(
SegmentKind::Transactional,
"install function",
"-- leading comment mentioning BEGIN\n\
CREATE OR REPLACE FUNCTION public.bump_counter()\n\
RETURNS trigger AS $body$\n\
BEGIN\n\
NEW.counter := COALESCE(NEW.counter, 0) + 1;\n\
RETURN NEW;\n\
END;\n\
$body$ LANGUAGE plpgsql;",
);
preflight_segment_sql_execution_compatibility(&plan)
.expect("function body BEGIN/END must not trip preflight");
}
#[test]
fn is_unique_violation_rejects_non_db_errors() {
let nf = DjogiError::not_found("users");
assert!(!is_unique_violation(&nf));
}
#[test]
fn checksum_for_baseline_snapshot_is_deterministic() {
use crate::migrate::schema::{AppliedSchema, SNAPSHOT_FORMAT_VERSION};
use std::collections::BTreeMap;
let snap = AppliedSchema {
djogi_version: "0.1.0".to_string(),
enums: BTreeMap::new(),
format_version: SNAPSHOT_FORMAT_VERSION.to_string(),
generated_at: "2026-04-25T00:00:00Z".to_string(),
indexes: Vec::new(),
models: BTreeMap::new(),
registered_apps: vec!["".to_string()],
};
let a = checksum_for_baseline_snapshot(&snap);
let b = checksum_for_baseline_snapshot(&snap);
assert_eq!(a, b, "same input must produce same checksum");
assert!(a.starts_with(super::ledger::CHECKSUM_PREFIX));
assert_eq!(a.len(), super::ledger::CHECKSUM_LEN);
}
#[test]
fn checksum_for_baseline_snapshot_changes_on_schema_change() {
use crate::migrate::schema::{AppliedSchema, SNAPSHOT_FORMAT_VERSION};
use std::collections::BTreeMap;
let mut a = AppliedSchema {
djogi_version: "0.1.0".to_string(),
enums: BTreeMap::new(),
format_version: SNAPSHOT_FORMAT_VERSION.to_string(),
generated_at: "2026-04-25T00:00:00Z".to_string(),
indexes: Vec::new(),
models: BTreeMap::new(),
registered_apps: vec!["".to_string()],
};
let cs_a = checksum_for_baseline_snapshot(&a);
a.registered_apps.push("billing".to_string());
let cs_b = checksum_for_baseline_snapshot(&a);
assert_ne!(cs_a, cs_b, "schema change must yield different checksum");
}
#[test]
fn baseline_snapshot_should_not_be_provided_renders_message() {
let e = RunnerError::BaselineSnapshotShouldNotBeProvided;
let msg = format!("{e}");
assert!(msg.contains("baseline_plan rejects caller-supplied snapshots"));
assert!(msg.contains("snapshot = None"));
}
#[test]
fn version_collision_non_terminal_renders_status_and_run_id() {
for (status, guidance) in [
(
LedgerStatus::Pending,
"then `repair_partial_apply` to resolve it in place",
),
(
LedgerStatus::Failed,
"then `repair_resume_partial_apply` if it is still resumable or `repair_partial_apply` otherwise",
),
(
LedgerStatus::RolledBack,
"the rolled_back row is retained as audit history",
),
] {
let e = RunnerError::VersionCollisionNonTerminal {
version: "V20260524010101__example".to_string(),
status,
run_id: 4242,
};
let msg = format!("{e}");
assert!(msg.contains("V20260524010101__example"));
assert!(msg.contains(status.as_db_str()));
assert!(msg.contains("run_id 4242"));
assert!(msg.contains("djogi migrations status"));
assert!(msg.contains(guidance));
}
}
#[test]
fn version_collision_non_terminal_has_no_error_source() {
let e = RunnerError::VersionCollisionNonTerminal {
version: "V20260524010101__example".to_string(),
status: LedgerStatus::Pending,
run_id: 7,
};
assert!(std::error::Error::source(&e).is_none());
}
#[allow(clippy::await_holding_lock)]
#[djogi_test]
async fn apply_plan_writes_audit_rows_for_executed_segments_when_key_unset(
mut ctx: DjogiContext,
) {
let _signing_key_env = SigningKeyEnvUnsetGuard::unset();
let plan = audit_plan();
let audit_pool = ctx
.share_pool()
.expect("djogi_test context should be pool-backed")
.inner;
let snapshot_path = unique_temp_path("audit-happy-path", "json");
let cleanup_path = snapshot_path.clone();
let runner_ctx =
runner_ctx_for_audit_with_snapshot_path(&plan, Some(audit_pool), Some(snapshot_path));
let guard = acquire_test_workspace_guard();
let expected_sig = audit_signature_hex_for_snapshot(
runner_ctx.snapshot.as_ref().expect("test snapshot"),
[0u8; 32],
)
.expect("expected audit signature");
let report = apply_plan(&mut ctx, &plan, &runner_ctx, &guard)
.await
.expect("apply should write audit rows");
assert_eq!(report.transactional_segments, 1);
assert_eq!(report.non_transactional_segments, 1);
assert_eq!(report.metadata_segments, 1);
let rows = ctx
.query_all(
"SELECT target_database, app_label, ddl_sql, snapshot_signature_hex, \
applied_at <= lead(applied_at) OVER (ORDER BY id) AS applied_before_next \
FROM djogi_ddl_audit ORDER BY id",
&[],
)
.await
.expect("read audit rows");
assert_eq!(
rows.len(),
2,
"transactional and non-transactional segments should be audited; metadata-only skipped"
);
let first_sql: String = rows[0].try_get("ddl_sql").expect("first ddl_sql");
assert!(
first_sql.contains("CREATE TABLE audit_a")
&& first_sql.contains(";\n")
&& first_sql.contains("CREATE TABLE audit_b"),
"transactional segment should store concatenated statement SQL; got {first_sql}"
);
let second_sql: String = rows[1].try_get("ddl_sql").expect("second ddl_sql");
assert!(
second_sql.contains("CREATE INDEX CONCURRENTLY audit_a_id_idx"),
"non-transactional segment should be audited; got {second_sql}"
);
for row in rows {
let target_database: String = row
.try_get("target_database")
.expect("target_database column");
let app_label: String = row.try_get("app_label").expect("app_label column");
let sig: String = row
.try_get("snapshot_signature_hex")
.expect("snapshot signature column");
let applied_before_next: Option<bool> = row
.try_get("applied_before_next")
.expect("applied_at monotonic column");
assert_eq!(target_database, "main");
assert_eq!(app_label, "");
assert_eq!(sig, expected_sig);
assert_eq!(
sig,
"0".repeat(64),
"unset signing key should persist the no-op zero signature"
);
assert_ne!(
applied_before_next,
Some(false),
"applied_at should be monotonic in audit id order"
);
}
let _ = std::fs::remove_file(cleanup_path);
}
#[djogi_test]
async fn apply_plan_skips_audit_when_pool_none(mut ctx: DjogiContext) {
let plan = single_table_plan("audit_pool_none_applies");
let snapshot_path = unique_temp_path("audit-pool-none", "json");
let cleanup_path = snapshot_path.clone();
let runner_ctx = runner_ctx_for_audit_with_snapshot_path(&plan, None, Some(snapshot_path));
let guard = acquire_test_workspace_guard();
let report = apply_plan(&mut ctx, &plan, &runner_ctx, &guard)
.await
.expect("audit_pool = None should not fail app-side apply");
assert_eq!(report.transactional_segments, 1);
let app_table: Option<String> = ctx
.query_one(
"SELECT to_regclass('public.audit_pool_none_applies')::text",
&[],
)
.await
.expect("query app table existence")
.try_get(0)
.expect("decode app table existence");
assert_eq!(
app_table.as_deref(),
Some("audit_pool_none_applies"),
"audit opt-out should still apply app-side DDL"
);
let audit_table: Option<String> = ctx
.query_one("SELECT to_regclass('public.djogi_ddl_audit')::text", &[])
.await
.expect("query audit table existence")
.try_get(0)
.expect("decode audit table existence");
assert_eq!(
audit_table, None,
"audit_pool = None should not bootstrap or write the audit table"
);
let _ = std::fs::remove_file(cleanup_path);
}
#[test]
fn audit_signature_for_unset_key_path_is_zero_hex() {
let key = audit_signing_key_from_loaded(Ok(None));
let sig_hex = audit_signature_hex_for_snapshot(&empty_snapshot(), key)
.expect("render audit signature");
assert_eq!(
sig_hex,
"0".repeat(64),
"an unset signing key should use the no-op key and persist zero hex"
);
}
#[allow(clippy::await_holding_lock)]
#[djogi_test]
async fn apply_plan_writes_audit_rows_when_snapshot_none(mut ctx: DjogiContext) {
let _signing_key_env = SigningKeyEnvUnsetGuard::unset();
let plan = single_table_plan("audit_snapshot_none_applies");
let audit_pool = ctx
.share_pool()
.expect("djogi_test context should be pool-backed")
.inner;
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260509000001__reset_audit_test".to_string(),
description: "snapshot-none audit test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: Some(audit_pool),
runner_identity: Some(RunnerIdentity::SingleNodeDev), drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let report = apply_plan(&mut ctx, &plan, &runner_ctx, &guard)
.await
.expect("apply with audit_pool=Some and snapshot=None must succeed");
assert_eq!(report.transactional_segments, 1);
let row_count: i64 = ctx
.query_one(
"SELECT COUNT(*)::bigint FROM djogi_ddl_audit \
WHERE target_database = 'main' AND app_label = ''",
&[],
)
.await
.expect("count audit rows")
.try_get(0)
.expect("decode count");
assert!(
row_count >= 1,
"expected at least one audit row when audit_pool=Some, even with snapshot=None; got {row_count}"
);
let sig: Option<String> = ctx
.query_one(
"SELECT snapshot_signature_hex FROM djogi_ddl_audit \
WHERE target_database = 'main' AND app_label = '' \
ORDER BY id DESC LIMIT 1",
&[],
)
.await
.expect("read most recent audit row signature")
.try_get(0)
.expect("decode signature");
assert_eq!(
sig, None,
"snapshot=None must persist NULL signature, not the no-op zero hex"
);
}
#[djogi_test]
async fn apply_plan_audit_failure_does_not_roll_back_app_db(mut ctx: DjogiContext) {
let _signing_key_env = SigningKeyEnvReadGuard::hold();
let plan = single_table_plan("audit_failure_survives");
let snapshot_path = unique_temp_path("audit-failure-survives", "json");
let cleanup_path = snapshot_path.clone();
let audit_pool = crate::pg::pool::DjogiPool::builder(
"postgres://djogi:djogi@127.0.0.1:1/djogi_unreachable",
)
.max_size(1)
.timeout(Duration::from_millis(50))
.build()
.await
.expect("build unreachable audit pool")
.inner;
let runner_ctx =
runner_ctx_for_audit_with_snapshot_path(&plan, Some(audit_pool), Some(snapshot_path));
let guard = acquire_test_workspace_guard();
let report = apply_plan(&mut ctx, &plan, &runner_ctx, &guard)
.await
.expect("audit-side failure should not fail app-side apply");
assert_eq!(report.transactional_segments, 1);
let app_table: Option<String> = ctx
.query_one(
"SELECT to_regclass('public.audit_failure_survives')::text",
&[],
)
.await
.expect("query app table existence")
.try_get(0)
.expect("decode app table existence");
assert_eq!(
app_table.as_deref(),
Some("audit_failure_survives"),
"app-side DDL should remain committed when the audit DB is unavailable"
);
let _ = std::fs::remove_file(cleanup_path);
}
#[djogi::deliberately_bypass_convention_with_raw_sql]
#[djogi_test]
async fn b1_rolledback_is_non_terminal_collision_class(mut ctx: DjogiContext) {
let version = "V20260526000001__b1_rolledback_test";
ledger::bootstrap(&mut ctx).await.expect("bootstrap ledger");
ctx.raw_execute(
"INSERT INTO djogi_schema_migrations \
(version, description, checksum_up, status, run_id, \
snapshot_version, app_label) \
VALUES ($1, 'b1 test', 'V1:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855', \
'rolled_back', 99, '1.0', '')",
&[&version],
)
.await
.expect("insert RolledBack row");
let result = classify_duplicate_version_collision(&mut ctx, version, "").await;
match result {
RunnerError::VersionCollisionNonTerminal { ref status, .. } => {
assert!(
matches!(status, LedgerStatus::RolledBack),
"RolledBack should produce VersionCollisionNonTerminal, got {:?}",
status
);
}
RunnerError::VersionAlreadyApplied { .. } => {
panic!("RolledBack must NOT produce VersionAlreadyApplied (Decision 1)");
}
other => {
panic!("unexpected error: {other:?}");
}
}
}
#[djogi::deliberately_bypass_convention_with_raw_sql]
#[djogi_test]
async fn b1_regression_pending_row_not_auto_deleted(mut ctx: DjogiContext) {
let version = "V20260526000002__b1_pending_guard";
ledger::bootstrap(&mut ctx).await.expect("bootstrap ledger");
ctx.raw_execute(
"INSERT INTO djogi_schema_migrations \
(version, description, checksum_up, status, run_id, \
snapshot_version, app_label) \
VALUES ($1, 'b1 pending guard', 'V1:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855', \
'pending', 98, '1.0', '')",
&[&version],
)
.await
.expect("insert Pending row");
let result = classify_duplicate_version_collision(&mut ctx, version, "").await;
match result {
RunnerError::VersionCollisionNonTerminal { ref status, .. } => {
assert!(
matches!(status, LedgerStatus::Pending),
"Pending should produce VersionCollisionNonTerminal, got {:?}",
status
);
}
RunnerError::VersionAlreadyApplied { .. } => {
panic!("Pending must NOT produce VersionAlreadyApplied");
}
other => {
panic!("unexpected error: {other:?}");
}
}
}
#[djogi::deliberately_bypass_convention_with_raw_sql]
#[djogi_test]
async fn c1_fake_apply_rejects_ooo_with_reject_policy(mut ctx: DjogiContext) {
let higher_version = "V20260527000000__c1_higher_peer";
ledger::bootstrap(&mut ctx).await.expect("bootstrap ledger");
ctx.raw_execute(
"INSERT INTO djogi_schema_migrations \
(version, description, checksum_up, status, run_id, \
snapshot_version, app_label) \
VALUES ($1, 'c1 higher peer', 'V1:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855', \
'applied', 97, '1.0', '')",
&[&higher_version],
)
.await
.expect("insert higher applied row");
let plan = MigrationPlan {
bucket: bucket("main", ""),
classification: Classification::Additive,
segments: vec![Segment {
kind: SegmentKind::Transactional,
statements: vec![],
}],
};
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260526000000__c1_lower_version".to_string(),
description: "c1 lower test".to_string(),
checksum_up: "V1:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
.to_string(),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig {
concurrent_warn_relpages: 5000,
strict_concurrent_warnings: false,
pk_flip_long_tx_threshold_secs: 30,
pk_flip_join_table_option: Default::default(),
},
out_of_order_policy: crate::migrate::OutOfOrderPolicy::Reject,
audit_pool: None,
runner_identity: Some(RunnerIdentity::SingleNodeDev), drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = fake_apply_plan(&mut ctx, &plan, &runner_ctx, &guard, "test reason").await;
assert!(
matches!(result, Err(RunnerError::OutOfOrderRejected { .. })),
"fake-apply with Reject policy and higher peer should return OutOfOrderRejected, got: {result:?}"
);
}
#[djogi::deliberately_bypass_convention_with_raw_sql]
#[djogi_test]
async fn c1_fake_apply_allows_ooo_with_diagnostic_and_sets_flag(mut ctx: DjogiContext) {
let higher_version = "V20260527000001__c1_higher_peer2";
ledger::bootstrap(&mut ctx).await.expect("bootstrap ledger");
ctx.raw_execute(
"INSERT INTO djogi_schema_migrations \
(version, description, checksum_up, status, run_id, \
snapshot_version, app_label) \
VALUES ($1, 'c1 higher peer 2', 'V1:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855', \
'applied', 96, '1.0', '')",
&[&higher_version],
)
.await
.expect("insert higher applied row");
let plan = MigrationPlan {
bucket: bucket("main", ""),
classification: Classification::Additive,
segments: vec![Segment {
kind: SegmentKind::Transactional,
statements: vec![],
}],
};
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260526000001__c1_lower_version2".to_string(),
description: "c1 lower test 2".to_string(),
checksum_up: "V1:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
.to_string(),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig {
concurrent_warn_relpages: 5000,
strict_concurrent_warnings: false,
pk_flip_long_tx_threshold_secs: 30,
pk_flip_join_table_option: Default::default(),
},
out_of_order_policy: crate::migrate::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: Some(RunnerIdentity::SingleNodeDev), drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = fake_apply_plan(&mut ctx, &plan, &runner_ctx, &guard, "test reason").await;
assert!(
result.is_ok(),
"fake-apply with AllowWithDiagnostic should succeed: {result:?}"
);
let rows = ctx
.query_all(
"SELECT out_of_order_flag FROM djogi_schema_migrations \
WHERE version = 'V20260526000001__c1_lower_version2'",
&[],
)
.await
.expect("query ledger");
assert!(rows.len() == 1, "ledger row should exist");
let ooo_flag: bool = rows[0].try_get(0).expect("get out_of_order_flag");
assert!(
ooo_flag,
"out_of_order_flag should be true for OOO fake-apply"
);
}
#[djogi::deliberately_bypass_convention_with_raw_sql]
#[djogi_test]
async fn c1_fake_apply_ooo_note_contains_override_reason(mut ctx: DjogiContext) {
let higher_version = "V20260527000002__c1_higher_peer3";
ledger::bootstrap(&mut ctx).await.expect("bootstrap ledger");
ctx.raw_execute(
"INSERT INTO djogi_schema_migrations \
(version, description, checksum_up, status, run_id, \
snapshot_version, app_label) \
VALUES ($1, 'c1 higher peer 3', 'V1:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855', \
'applied', 95, '1.0', '')",
&[&higher_version],
)
.await
.expect("insert higher applied row");
let plan = MigrationPlan {
bucket: bucket("main", ""),
classification: Classification::Additive,
segments: vec![Segment {
kind: SegmentKind::Transactional,
statements: vec![],
}],
};
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260526000002__c1_lower_version3".to_string(),
description: "c1 lower test 3".to_string(),
checksum_up: "V1:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
.to_string(),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig {
concurrent_warn_relpages: 5000,
strict_concurrent_warnings: false,
pk_flip_long_tx_threshold_secs: 30,
pk_flip_join_table_option: Default::default(),
},
out_of_order_policy: crate::migrate::OutOfOrderPolicy::AllowExplicit {
override_reason: "merge window".to_string(),
},
audit_pool: None,
runner_identity: Some(RunnerIdentity::SingleNodeDev), drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result =
fake_apply_plan(&mut ctx, &plan, &runner_ctx, &guard, "schema pre-exists").await;
assert!(
result.is_ok(),
"fake-apply with AllowExplicit should succeed: {result:?}"
);
let rows = ctx
.query_all(
"SELECT partial_apply_note FROM djogi_schema_migrations \
WHERE version = 'V20260526000002__c1_lower_version3'",
&[],
)
.await
.expect("query ledger");
assert!(rows.len() == 1, "ledger row should exist");
let note: String = rows[0].try_get(0).expect("get partial_apply_note");
assert!(
note.contains("faked at"),
"note should contain 'faked at': {note}"
);
assert!(
note.contains("reason: schema pre-exists"),
"note should contain fake reason: {note}"
);
assert!(
note.contains("out-of-order apply"),
"note should contain out-of-order annotation: {note}"
);
assert!(
note.contains("override: merge window"),
"note should contain override reason: {note}"
);
}
#[test]
fn expand_partition_statement_partitioned_index_preserves_leaf_drop_down_sql() {
let stmt = OperationSql {
label: "PkFlipPartitionedIndex events_p".to_string(),
up: "CREATE UNIQUE INDEX events_p_ts_id_desc_idx ON ONLY events_p (ts, id_desc);\n\
-- Per leaf: CREATE UNIQUE INDEX CONCURRENTLY <leaf>_ts_id_desc_idx\n"
.to_string(),
down: "DROP INDEX IF EXISTS events_p_ts_id_desc_idx;".to_string(),
lossy: None,
};
let expanded = expand_partition_statement(
&stmt,
"events_p",
&["events_p_a".to_string(), "events_p_b".to_string()],
PartitionExpansionMode::ReplayStrict,
)
.expect("strict replay expansion with leaves");
assert_eq!(
expanded[0].label,
"PkFlipPartitionedIndex events_p (parent-level)"
);
assert_eq!(
expanded[0].down,
"DROP INDEX IF EXISTS events_p_ts_id_desc_idx;"
);
assert!(
expanded.iter().any(|s| {
s.label == "PkFlipPartitionedIndex events_p leaf=events_p_a (concurrent)"
&& s.down
== "DROP INDEX IF EXISTS events_p_ts_id_desc_idx; \
DROP INDEX IF EXISTS events_p_a_ts_id_desc_idx"
}),
"leaf A concurrent must drop parent then leaf: {expanded:?}",
);
assert!(
expanded.iter().any(|s| {
s.label == "PkFlipPartitionedIndex events_p leaf=events_p_b (concurrent)"
&& s.down
== "DROP INDEX IF EXISTS events_p_ts_id_desc_idx; \
DROP INDEX IF EXISTS events_p_b_ts_id_desc_idx"
}),
"leaf B concurrent must drop parent then leaf: {expanded:?}",
);
assert!(
expanded
.iter()
.filter(|s| s.label.contains("(attach)"))
.all(|s| s.down.is_empty()),
"attach statements remain forward-only",
);
}
#[test]
fn expand_partition_statement_schema_qualified_leaf_in_drop_down_sql() {
let parent = "myschema.events_p";
let leaves = vec![
"myschema.events_p_a".to_string(),
"myschema.events_p_b".to_string(),
];
let plan = OperationSql {
label: format!("PkFlipPartitionedIndex {parent}"),
up: "CREATE UNIQUE INDEX events_p_ts_id_desc_idx ON ONLY myschema.events_p (ts, id_desc)"
.to_string(),
down: "DROP INDEX IF EXISTS myschema.events_p_ts_id_desc_idx".to_string(),
lossy: None,
};
let expanded = expand_partition_statement(
&plan,
parent,
&leaves,
PartitionExpansionMode::ApplyLenient,
)
.expect("partition expansion should succeed");
assert_eq!(expanded.len(), 5);
let leaf_a_concurrent = expanded
.iter()
.find(|s| s.label == "PkFlipPartitionedIndex myschema.events_p leaf=myschema.events_p_a (concurrent)")
.expect("leaf A concurrent entry");
assert_eq!(
leaf_a_concurrent.up,
"CREATE UNIQUE INDEX CONCURRENTLY events_p_a_ts_id_desc_idx ON myschema.events_p_a (ts, id_desc)",
);
assert_eq!(
leaf_a_concurrent.down,
"DROP INDEX IF EXISTS myschema.events_p_ts_id_desc_idx; DROP INDEX IF EXISTS myschema.events_p_a_ts_id_desc_idx",
);
let leaf_a_attach = expanded
.iter()
.find(|s| {
s.label
== "PkFlipPartitionedIndex myschema.events_p leaf=myschema.events_p_a (attach)"
})
.expect("leaf A attach entry");
assert_eq!(
leaf_a_attach.up,
"ALTER INDEX myschema.events_p_ts_id_desc_idx ATTACH PARTITION myschema.events_p_a_ts_id_desc_idx",
);
assert_eq!(leaf_a_attach.down, "");
for entry in expanded.iter().filter(|s| s.label.contains("(concurrent)")) {
assert!(
!entry.up.contains("CONCURRENTLY myschema."),
"concurrent up must not have schema in index name: {}",
entry.up,
);
assert!(
!entry.up.contains("id DESC"),
"up must not contain space-form 'id DESC': {}",
entry.up,
);
}
}
#[test]
fn expand_partition_leaf_placeholders_replay_mode_refuses_empty_leaves() {
let stmt = OperationSql {
label: "PkFlipPartitionedIndex events_p".to_string(),
up: "CREATE UNIQUE INDEX events_p_ts_id_desc_idx ON ONLY events_p (ts, id_desc);"
.to_string(),
down: "DROP INDEX IF EXISTS events_p_ts_id_desc_idx;".to_string(),
lossy: None,
};
let err = expand_partition_statement(
&stmt,
"events_p",
&[],
PartitionExpansionMode::ReplayStrict,
)
.expect_err("strict replay must not replace a partition expansion with a no-op comment");
match err {
RunnerError::PartitionExpansionNoLeaves {
parent,
statement_label,
} => {
assert_eq!(parent, "events_p");
assert_eq!(statement_label, "PkFlipPartitionedIndex events_p");
}
other => panic!("expected PartitionExpansionNoLeaves, got {other:?}"),
}
let lenient = expand_partition_statement(
&stmt,
"events_p",
&[],
PartitionExpansionMode::ApplyLenient,
)
.expect("apply mode keeps the existing empty-leaf fallback");
assert_eq!(lenient.len(), 1);
assert!(lenient[0].up.contains("pg_inherits returned 0 leaves"));
}
#[test]
fn rollback_leaf_identity_mismatch_display() {
let err = RollbackError::LeafIdentityMismatch {
version: "001_create_users".to_string(),
stored_leaf_identity: "public.users:public.users_p2024_01,public.users_p2024_02\n"
.to_string(),
current_leaf_identity: "public.users:public.users_p2024_01,public.users_p2024_03\n"
.to_string(),
};
let msg = format!("{}", err);
assert!(msg.contains("[D624]"));
assert!(msg.contains("rollback refused"));
assert!(msg.contains("partition leaf identity mismatch"));
assert!(msg.contains("001_create_users"));
}
#[test]
fn rollback_stale_phase_zero_down_error_display() {
let phase_zero_version = crate::migrate::bootstrap::PHASE_ZERO_VERSION;
let err = RollbackError::StalePhaseZeroDown {
version: phase_zero_version.to_string(),
refusal_reason: "generated-stale",
};
let msg = format!("{}", err);
assert!(msg.contains("rollback refused"));
assert!(msg.contains("Phase 0 down SQL"));
assert!(msg.contains("generated-stale"));
assert!(msg.contains(phase_zero_version));
}
#[test]
fn rollback_stale_phase_zero_down_error_display_ambiguous() {
let phase_zero_version = crate::migrate::bootstrap::PHASE_ZERO_VERSION;
let err = RollbackError::StalePhaseZeroDown {
version: phase_zero_version.to_string(),
refusal_reason: "ambiguous",
};
let msg = format!("{}", err);
assert!(msg.contains("ambiguous"));
}
#[test]
fn rollback_missing_identity_error_display() {
let err = RollbackError::MissingRollbackIdentity {
version: "V20260101000000__add_users".to_string(),
};
let msg = format!("{}", err);
assert!(msg.contains("rollback refused"));
assert!(msg.contains("non-Phase-0 rollback"));
assert!(msg.contains("V20260101000000__add_users"));
assert!(msg.contains("binding-capable runner identity"));
}
#[test]
fn rollback_lossy_allow_reason_preserves_reason_without_lossy_markers() {
let plan = single_table_plan("lossy_reason_markerless");
let reason = "operator supplied a file-derived rollback reason".to_string();
let allow_reason = rollback_lossy_allow_reason(
&plan,
&LossyRollbackPolicy::Allow {
reason: reason.clone(),
},
)
.expect("allow policy should not refuse markerless rollback");
assert_eq!(
allow_reason.as_deref(),
Some(reason.as_str()),
"Allow {{ reason }} must survive even when the replay plan has no lossy markers"
);
}
#[test]
fn rollback_error_live_db_committed_signal_is_conservative() {
let pre_execution_runner = RollbackError::Runner {
source: RunnerError::LockTimeout {
path: PathBuf::from("/tmp/djogi-lock"),
holder_pid: None,
},
live_db_committed: false,
};
assert!(
!pre_execution_runner.live_db_committed(),
"pre-execution runner refusals must report non-committed"
);
let post_commit_runner = RollbackError::Runner {
source: RunnerError::LockTimeout {
path: PathBuf::from("/tmp/djogi-lock"),
holder_pid: None,
},
live_db_committed: true,
};
assert!(
post_commit_runner.live_db_committed(),
"post-commit runner errors must report committed"
);
let transactional_statement_failed = RollbackError::DownStatementFailed {
segment_index: 0,
statement_label: "DropTable tx".to_string(),
live_db_committed: false,
source: DjogiError::Db(DbError::other("transactional rollback failure")),
};
assert!(
!transactional_statement_failed.live_db_committed(),
"transactional statement failures before COMMIT should remain non-committed"
);
let non_transactional_statement_failed = RollbackError::DownStatementFailed {
segment_index: 1,
statement_label: "DropIndex concurrently".to_string(),
live_db_committed: true,
source: DjogiError::Db(DbError::other("non-transactional rollback failure")),
};
assert!(
non_transactional_statement_failed.live_db_committed(),
"non-transactional down failures must report committed once execution was attempted"
);
let snapshot_failed = RollbackError::SnapshotPersistFailed {
path: PathBuf::from("/tmp/djogi-snapshot"),
source: SnapshotError::Io {
path: None,
source: std::io::Error::other("snapshot write failed"),
},
};
assert!(
snapshot_failed.live_db_committed(),
"snapshot persistence happens after down SQL and ledger mutation"
);
let checksum_refusal = RollbackError::ChecksumDrift {
version: "V20260613000000__checksum_refusal".to_string(),
side: RollbackChecksumSide::Up,
ledger: Some(
"V1:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".to_string(),
),
on_disk: Some(
"V1:bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".to_string(),
),
};
assert!(
!checksum_refusal.live_db_committed(),
"checksum parity refusal must happen before destructive down SQL"
);
}
#[test]
fn migrate_surface_reexports_non_transactional_shape_helper() {
assert_eq!(
crate::migrate::find_non_transactional_statement_shape(
"CREATE INDEX CONCURRENTLY idx ON table_name (id)"
),
Some("CREATE INDEX CONCURRENTLY")
);
}
#[test]
fn phase_zero_down_classification_identifies_literal_database_default_as_stale() {
let down_sql = "ALTER DATABASE \"mydb\" SET heer.node_id = '1';\n\
ALTER DATABASE \"mydb\" SET heer.ranj_node_id = '1';";
let state = crate::migrate::phase_zero::classify_phase_zero_artifact(down_sql.as_bytes());
assert_ne!(
state,
crate::migrate::phase_zero::PhaseZeroArtifactState::IdentityFreeCurrent
);
}
#[test]
fn phase_zero_down_classification_identifies_complete_stale_as_generated_stale() {
let mut down_sql =
crate::migrate::bootstrap::compose_phase_zero("mydb", &BTreeSet::new(), 1, true)
.expect("compose seed-capable Phase 0");
let start = down_sql.find("DO $djogi$").expect("dynamic defaults");
let end = down_sql[start..]
.find("SET heer.node_id = '1';")
.map(|offset| start + offset)
.expect("session SET after dynamic defaults");
down_sql.replace_range(
start..end,
"ALTER DATABASE \"mydb\" SET heer.node_id = '1';\n\
ALTER DATABASE \"mydb\" SET heer.ranj_node_id = '1';\n",
);
let state = crate::migrate::phase_zero::classify_phase_zero_artifact(down_sql.as_bytes());
assert_eq!(
state,
crate::migrate::phase_zero::PhaseZeroArtifactState::GeneratedStale
);
}
#[test]
fn phase_zero_down_classification_identity_free_production_is_ok() {
let down_sql = "DROP TABLE IF EXISTS heer.heer_nodes;\n\
DROP FUNCTION IF EXISTS heer.heerid_next();";
let state = crate::migrate::phase_zero::classify_phase_zero_artifact(down_sql.as_bytes());
assert_ne!(
state,
crate::migrate::phase_zero::PhaseZeroArtifactState::IdentityFreeCurrent
);
}
#[test]
fn phase_zero_down_classification_comment_only_is_not_current() {
let down_sql = "-- Phase 0 down — no reverse operations for bootstrap";
let state = crate::migrate::phase_zero::classify_phase_zero_artifact(down_sql.as_bytes());
assert_ne!(
state,
crate::migrate::phase_zero::PhaseZeroArtifactState::IdentityFreeCurrent
);
}
fn phase_zero_plan_with_down_payload(down: &str) -> MigrationPlan {
MigrationPlan {
bucket: bucket("main", ""),
classification: Classification::Additive,
segments: vec![Segment {
kind: SegmentKind::Transactional,
statements: vec![OperationSql {
label: "Phase 0 rollback payload".to_string(),
up: String::new(),
down: down.to_string(),
lossy: None,
}],
}],
}
}
fn phase_zero_runner_ctx_for(plan: &MigrationPlan) -> RunnerCtx {
RunnerCtx {
bucket: plan.bucket.clone(),
version: PHASE_ZERO_VERSION.to_string(),
description: "Phase 0 rollback preflight".to_string(),
checksum_up: compute_checksum_for_plan_up(plan),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
}
}
#[test]
fn phase_zero_down_missing_payload_refuses_before_execution() {
for down in ["", " \n\t "] {
let plan = phase_zero_plan_with_down_payload(down);
let runner_ctx = phase_zero_runner_ctx_for(&plan);
let result = preflight_phase_zero_down_payload(&plan, &runner_ctx);
assert!(
matches!(
result,
Err(RollbackError::StalePhaseZeroDown {
refusal_reason: "missing",
..
})
),
"missing/whitespace Phase 0 down payload must refuse, got {result:?}"
);
}
}
#[test]
fn phase_zero_down_seed_capable_payload_refuses_before_execution() {
let down_sql = crate::migrate::bootstrap::compose_phase_zero(
"main",
&BTreeSet::new(),
crate::migrate::bootstrap::DEFAULT_NODE_ID,
true,
)
.expect("compose seed-capable Phase 0");
let plan = phase_zero_plan_with_down_payload(&down_sql);
let runner_ctx = phase_zero_runner_ctx_for(&plan);
let result = preflight_phase_zero_down_payload(&plan, &runner_ctx);
assert!(
matches!(
result,
Err(RollbackError::StalePhaseZeroDown {
refusal_reason: "seed-capable-runtime-only",
..
})
),
"seed-capable Phase 0 down payload must refuse before execution, got {result:?}"
);
}
#[test]
fn phase_zero_down_markerless_seed_payload_refuses_before_execution() {
let mut down_sql = crate::migrate::bootstrap::compose_phase_zero(
"main",
&BTreeSet::new(),
crate::migrate::bootstrap::DEFAULT_NODE_ID,
false,
)
.expect("compose identity-free Phase 0");
down_sql.push_str("\nINSERT INTO heer.heer_nodes (id) VALUES (1);\n");
let plan = phase_zero_plan_with_down_payload(&down_sql);
let runner_ctx = phase_zero_runner_ctx_for(&plan);
let result = preflight_phase_zero_down_payload(&plan, &runner_ctx);
assert!(
matches!(
result,
Err(RollbackError::StalePhaseZeroDown {
refusal_reason: "seed-dml-not-runtime-current",
..
})
),
"markerless seed Phase 0 down payload must refuse before execution, got {result:?}"
);
}
#[test]
fn phase_zero_down_extended_seed_payloads_refuse_before_execution() {
for (name, statement) in extended_seed_statement_cases() {
let mut down_sql = crate::migrate::bootstrap::compose_phase_zero(
"main",
&BTreeSet::new(),
crate::migrate::bootstrap::DEFAULT_NODE_ID,
false,
)
.expect("compose identity-free Phase 0");
down_sql.push('\n');
down_sql.push_str(statement);
down_sql.push('\n');
let plan = phase_zero_plan_with_down_payload(&down_sql);
let runner_ctx = phase_zero_runner_ctx_for(&plan);
let result = preflight_phase_zero_down_payload(&plan, &runner_ctx);
assert!(
matches!(
result,
Err(RollbackError::StalePhaseZeroDown {
refusal_reason: "seed-dml-not-runtime-current",
..
})
),
"extended seed Phase 0 down payload {name} must refuse before execution, got {result:?}"
);
}
}
#[test]
fn runner_identity_requires_binding_for_selected() {
assert!(RunnerIdentity::Selected { id: 7 }.requires_binding());
}
#[test]
fn runner_identity_requires_binding_for_single_node_dev() {
assert!(RunnerIdentity::SingleNodeDev.requires_binding());
}
#[test]
fn runner_identity_does_not_require_binding_for_identity_free() {
assert!(!RunnerIdentity::IdentityFree.requires_binding());
}
#[test]
fn runner_identity_selected_returns_node_id() {
assert_eq!(RunnerIdentity::Selected { id: 7 }.node_id(), Some(7));
}
#[test]
fn runner_identity_single_node_dev_returns_one() {
assert_eq!(RunnerIdentity::SingleNodeDev.node_id(), Some(1));
}
#[test]
fn runner_identity_free_returns_none() {
assert_eq!(RunnerIdentity::IdentityFree.node_id(), None);
}
#[tokio::test]
async fn apply_plan_missing_identity_refuses_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = single_table_plan("pre_pin_apply_no_identity");
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260602000000__pre_pin_apply_no_identity".to_string(),
description: "pre-pin apply identity test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: Some(empty_snapshot()),
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = apply_plan(&mut ctx, &plan, &runner_ctx, &guard).await;
assert!(
matches!(result, Err(RunnerError::MissingApplyIdentity { .. })),
"missing identity must refuse before pin/connection checkout, got {result:?}"
);
}
#[tokio::test]
async fn fake_apply_plan_missing_identity_refuses_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = single_segment_plan(SegmentKind::Transactional, "fake", "SELECT 1");
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260602000001__pre_pin_fake_no_identity".to_string(),
description: "pre-pin fake identity test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = fake_apply_plan(&mut ctx, &plan, &runner_ctx, &guard, "pre-pin").await;
assert!(
matches!(result, Err(RunnerError::MissingApplyIdentity { .. })),
"missing identity must refuse before fake-apply pins, got {result:?}"
);
}
#[tokio::test]
async fn baseline_plan_missing_identity_refuses_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = single_table_plan("pre_pin_baseline_no_identity");
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260602000002__pre_pin_baseline_no_identity".to_string(),
description: "pre-pin baseline identity test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = baseline_plan(&mut ctx, &plan.bucket, &runner_ctx, &guard, "pre-pin").await;
assert!(
matches!(result, Err(RunnerError::MissingApplyIdentity { .. })),
"missing identity must refuse before baseline pins, got {result:?}"
);
}
#[tokio::test]
async fn rollback_plan_missing_identity_refuses_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = single_table_plan("pre_pin_rollback_no_identity");
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260602000003__pre_pin_rollback_no_identity".to_string(),
description: "pre-pin rollback identity test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = rollback_plan(
&mut ctx,
&plan,
&runner_ctx,
&guard,
LossyRollbackPolicy::Refuse,
None,
)
.await;
assert!(
matches!(result, Err(RollbackError::MissingRollbackIdentity { .. })),
"missing identity must refuse before rollback pins, got {result:?}"
);
}
#[tokio::test]
async fn apply_plan_checksum_mismatch_precedes_missing_identity_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = single_table_plan("pre_pin_apply_bad_checksum");
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260602000004__pre_pin_apply_bad_checksum".to_string(),
description: "pre-pin apply checksum precedence test".to_string(),
checksum_up: compute_checksum(["different sql"]),
checksum_down: None,
snapshot: Some(empty_snapshot()),
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = apply_plan(&mut ctx, &plan, &runner_ctx, &guard).await;
assert!(
matches!(result, Err(RunnerError::ChecksumMismatch(_))),
"checksum mismatch must precede missing identity and pin, got {result:?}"
);
}
#[tokio::test]
async fn fake_apply_plan_checksum_mismatch_precedes_missing_identity_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = single_segment_plan(SegmentKind::Transactional, "fake", "SELECT 1");
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260602000005__pre_pin_fake_bad_checksum".to_string(),
description: "pre-pin fake checksum precedence test".to_string(),
checksum_up: compute_checksum(["different sql"]),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = fake_apply_plan(&mut ctx, &plan, &runner_ctx, &guard, "pre-pin").await;
assert!(
matches!(result, Err(RunnerError::ChecksumMismatch(_))),
"fake checksum mismatch must precede missing identity and pin, got {result:?}"
);
}
#[tokio::test]
async fn baseline_snapshot_refusal_precedes_missing_identity_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = single_table_plan("pre_pin_baseline_snapshot");
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260602000006__pre_pin_baseline_snapshot".to_string(),
description: "pre-pin baseline snapshot precedence test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: Some(empty_snapshot()),
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = baseline_plan(&mut ctx, &plan.bucket, &runner_ctx, &guard, "pre-pin").await;
assert!(
matches!(
result,
Err(RunnerError::BaselineSnapshotShouldNotBeProvided)
),
"snapshot refusal must precede missing identity and pin, got {result:?}"
);
}
#[tokio::test]
async fn rollback_prior_snapshot_refusal_precedes_missing_identity_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = single_table_plan("pre_pin_rollback_prior_snapshot");
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260602000007__pre_pin_rollback_prior_snapshot".to_string(),
description: "pre-pin rollback prior snapshot precedence test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: None,
snapshot_path: Some(unique_temp_path("pre-pin-prior", "json")),
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = rollback_plan(
&mut ctx,
&plan,
&runner_ctx,
&guard,
LossyRollbackPolicy::Refuse,
None,
)
.await;
assert!(
matches!(result, Err(RollbackError::PriorSnapshotMissing)),
"prior snapshot refusal must precede missing identity and pin, got {result:?}"
);
}
#[tokio::test]
async fn apply_plan_stale_phase_zero_artifact_refuses_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = generated_stale_phase_zero_plan();
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: PHASE_ZERO_VERSION.to_string(),
description: "pre-pin Phase 0 stale artifact test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: Some(empty_snapshot()),
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = apply_plan(&mut ctx, &plan, &runner_ctx, &guard).await;
assert!(
matches!(
result,
Err(RunnerError::StalePhaseZeroArtifact {
refusal_reason: "generated-stale",
..
})
),
"stale Phase 0 artifact must refuse before apply pins, got {result:?}"
);
}
#[tokio::test]
async fn apply_plan_seed_capable_phase_zero_artifact_refuses_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = seed_capable_phase_zero_plan();
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: PHASE_ZERO_VERSION.to_string(),
description: "pre-pin Phase 0 seed-capable artifact test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: Some(empty_snapshot()),
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = apply_plan(&mut ctx, &plan, &runner_ctx, &guard).await;
assert!(
matches!(
result,
Err(RunnerError::StalePhaseZeroArtifact {
refusal_reason: "seed-capable-runtime-only",
..
})
),
"seed-capable Phase 0 artifact must refuse before apply pins, got {result:?}"
);
}
#[tokio::test]
async fn apply_plan_markerless_seed_phase_zero_artifact_refuses_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = markerless_seed_phase_zero_plan();
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: PHASE_ZERO_VERSION.to_string(),
description: "pre-pin Phase 0 markerless seed artifact test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: Some(empty_snapshot()),
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = apply_plan(&mut ctx, &plan, &runner_ctx, &guard).await;
assert!(
matches!(
result,
Err(RunnerError::StalePhaseZeroArtifact {
refusal_reason: "seed-dml-not-runtime-current",
..
})
),
"markerless seed Phase 0 artifact must refuse before apply pins, got {result:?}"
);
}
#[tokio::test]
async fn apply_plan_extended_seed_phase_zero_artifacts_refuse_before_pin() {
for (name, statement) in extended_seed_statement_cases() {
let mut ctx = unreachable_pool_context().await;
let plan = phase_zero_plan_with_seed_statement(statement);
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: PHASE_ZERO_VERSION.to_string(),
description: format!("pre-pin Phase 0 extended seed artifact test: {name}"),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: Some(empty_snapshot()),
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = apply_plan(&mut ctx, &plan, &runner_ctx, &guard).await;
assert!(
matches!(
result,
Err(RunnerError::StalePhaseZeroArtifact {
refusal_reason: "seed-dml-not-runtime-current",
..
})
),
"extended seed Phase 0 artifact {name} must refuse before apply pins, got {result:?}"
);
}
}
#[tokio::test]
async fn fake_apply_plan_stale_phase_zero_artifact_refuses_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = generated_stale_phase_zero_plan();
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: PHASE_ZERO_VERSION.to_string(),
description: "pre-pin Phase 0 fake stale artifact test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = fake_apply_plan(&mut ctx, &plan, &runner_ctx, &guard, "pre-pin").await;
assert!(
matches!(
result,
Err(RunnerError::StalePhaseZeroArtifact {
refusal_reason: "generated-stale",
..
})
),
"stale Phase 0 artifact must refuse before fake-apply pins, got {result:?}"
);
}
#[tokio::test]
async fn fake_apply_plan_seed_capable_phase_zero_artifact_refuses_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = seed_capable_phase_zero_plan();
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: PHASE_ZERO_VERSION.to_string(),
description: "pre-pin Phase 0 fake seed-capable artifact test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = fake_apply_plan(&mut ctx, &plan, &runner_ctx, &guard, "pre-pin").await;
assert!(
matches!(
result,
Err(RunnerError::StalePhaseZeroArtifact {
refusal_reason: "seed-capable-runtime-only",
..
})
),
"seed-capable Phase 0 artifact must refuse before fake-apply pins, got {result:?}"
);
}
#[tokio::test]
async fn fake_apply_plan_markerless_seed_phase_zero_artifact_refuses_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = markerless_seed_phase_zero_plan();
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: PHASE_ZERO_VERSION.to_string(),
description: "pre-pin Phase 0 fake markerless seed artifact test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = fake_apply_plan(&mut ctx, &plan, &runner_ctx, &guard, "pre-pin").await;
assert!(
matches!(
result,
Err(RunnerError::StalePhaseZeroArtifact {
refusal_reason: "seed-dml-not-runtime-current",
..
})
),
"markerless seed Phase 0 artifact must refuse before fake-apply pins, got {result:?}"
);
}
#[tokio::test]
async fn fake_apply_plan_copy_from_seed_phase_zero_artifact_refuses_before_pin() {
let mut ctx = unreachable_pool_context().await;
let plan = phase_zero_plan_with_seed_statement(
"COPY \"heer\".\"heer_ranj_node_state\" (\"node_id\") FROM STDIN;",
);
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: PHASE_ZERO_VERSION.to_string(),
description: "pre-pin Phase 0 fake COPY FROM seed artifact test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None,
drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = fake_apply_plan(&mut ctx, &plan, &runner_ctx, &guard, "pre-pin").await;
assert!(
matches!(
result,
Err(RunnerError::StalePhaseZeroArtifact {
refusal_reason: "seed-dml-not-runtime-current",
..
})
),
"COPY FROM seed Phase 0 artifact must refuse before fake-apply pins, got {result:?}"
);
}
#[test]
fn truncate_for_log_handles_multibyte_boundary() {
let sql = format!("{}{}", "é".repeat(200), " SELECT 1");
let truncated = truncate_for_log(&sql);
assert!(truncated.ends_with("..."));
assert!(
truncated.len() <= 259,
"truncated message should keep the 256-byte budget plus ellipsis"
);
}
#[djogi_test]
async fn apply_plan_no_identity_non_phase_zero_refused(mut ctx: DjogiContext) {
ledger::bootstrap(&mut ctx).await.expect("bootstrap ledger");
let plan = single_table_plan("g_gate_apply_no_id");
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260601000000__g_gate_apply".to_string(),
description: "G-gate apply test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: Some(empty_snapshot()),
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None, drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = apply_plan(&mut ctx, &plan, &runner_ctx, &guard).await;
assert!(
matches!(result, Err(RunnerError::MissingApplyIdentity { .. })),
"non-P0 apply with no identity should refuse with MissingApplyIdentity, got: {result:?}",
);
let count: i64 = ctx
.query_one(
"SELECT COUNT(*) FROM djogi_schema_migrations WHERE version = 'V20260601000000__g_gate_apply'",
&[],
)
.await
.expect("count ledger rows")
.try_get(0)
.expect("count column");
assert_eq!(
count, 0,
"no ledger row should be inserted before the G-gate refusal"
);
}
#[djogi_test]
async fn fake_apply_no_identity_non_phase_zero_refused(mut ctx: DjogiContext) {
ledger::bootstrap(&mut ctx).await.expect("bootstrap ledger");
let plan = single_segment_plan(SegmentKind::Transactional, "G-gate fake-apply", "SELECT 1");
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260601000001__g_gate_fake".to_string(),
description: "G-gate fake-apply test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: None,
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None, drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = fake_apply_plan(&mut ctx, &plan, &runner_ctx, &guard, "G-gate test").await;
assert!(
matches!(result, Err(RunnerError::MissingApplyIdentity { .. })),
"non-P0 fake-apply with no identity should refuse with MissingApplyIdentity, got: {result:?}",
);
let count: i64 = ctx
.query_one(
"SELECT COUNT(*) FROM djogi_schema_migrations WHERE version = 'V20260601000001__g_gate_fake'",
&[],
)
.await
.expect("count ledger rows")
.try_get(0)
.expect("count column");
assert_eq!(
count, 0,
"no ledger row should be inserted before the G-gate refusal"
);
}
#[djogi::deliberately_bypass_convention_with_raw_sql]
#[djogi_test]
async fn baseline_no_identity_non_phase_zero_refused(mut ctx: DjogiContext) {
ledger::bootstrap(&mut ctx).await.expect("bootstrap ledger");
let plan = single_table_plan("g_gate_baseline_no_id");
ctx.raw_execute(
&format!("CREATE TABLE {} (id bigint)", "g_gate_baseline_no_id"),
&[],
)
.await
.expect("create table for baseline test");
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260601000002__g_gate_baseline".to_string(),
description: "G-gate baseline test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: None, snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None, drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let bucket = plan.bucket.clone();
let result = baseline_plan(&mut ctx, &bucket, &runner_ctx, &guard, "G-gate test").await;
assert!(
matches!(result, Err(RunnerError::MissingApplyIdentity { .. })),
"non-P0 baseline with no identity should refuse with MissingApplyIdentity, got: {result:?}",
);
let count: i64 = ctx
.query_one(
"SELECT COUNT(*) FROM djogi_schema_migrations WHERE version = 'V20260601000002__g_gate_baseline'",
&[],
)
.await
.expect("count ledger rows")
.try_get(0)
.expect("count column");
assert_eq!(
count, 0,
"no ledger row should be inserted before the G-gate refusal"
);
}
#[djogi_test]
async fn apply_plan_identity_free_non_phase_zero_refused(mut ctx: DjogiContext) {
ledger::bootstrap(&mut ctx).await.expect("bootstrap ledger");
let plan = single_table_plan("g_gate_apply_idfree");
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: "V20260601000003__g_gate_idfree".to_string(),
description: "G-gate IdentityFree test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: Some(empty_snapshot()),
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: Some(RunnerIdentity::IdentityFree), drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = apply_plan(&mut ctx, &plan, &runner_ctx, &guard).await;
assert!(
matches!(result, Err(RunnerError::MissingApplyIdentity { .. })),
"non-P0 apply with IdentityFree should refuse with MissingApplyIdentity, got: {result:?}",
);
let count: i64 = ctx
.query_one(
"SELECT COUNT(*) FROM djogi_schema_migrations WHERE version = 'V20260601000003__g_gate_idfree'",
&[],
)
.await
.expect("count ledger rows")
.try_get(0)
.expect("count column");
assert_eq!(
count, 0,
"no ledger row should be inserted before the G-gate refusal"
);
}
#[djogi_test]
async fn apply_phase_zero_no_identity_allowed(mut ctx: DjogiContext) {
ledger::bootstrap(&mut ctx).await.expect("bootstrap ledger");
let plan = current_phase_zero_plan();
let runner_ctx = RunnerCtx {
bucket: plan.bucket.clone(),
version: PHASE_ZERO_VERSION.to_string(),
description: "G-gate Phase 0 carve-out test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: Some(empty_snapshot()),
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: None, drift_baseline: DriftBaseline::Disabled,
};
let guard = acquire_test_workspace_guard();
let result = apply_plan(&mut ctx, &plan, &runner_ctx, &guard).await;
assert!(
result.is_ok(),
"Phase 0 apply with no identity should succeed (carve-out), got: {result:?}",
);
let count: i64 = ctx
.query_one(
&format!(
"SELECT COUNT(*) FROM djogi_schema_migrations WHERE version = '{}'",
PHASE_ZERO_VERSION
),
&[],
)
.await
.expect("count ledger rows")
.try_get(0)
.expect("count column");
assert_eq!(count, 1, "Phase 0 apply should insert a ledger row");
}
#[djogi_test]
async fn bind_runner_node_identity_second_guc_failure(mut ctx: DjogiContext) {
const TEST_NODE_ID: i32 = 1;
ctx.execute(
"DROP FUNCTION IF EXISTS set_heer_ranj_node_id(INTEGER)",
&[],
)
.await
.expect("drop original second GUC function");
ctx.execute(
"CREATE FUNCTION set_heer_ranj_node_id(INTEGER) RETURNS void AS $$ \
BEGIN RAISE EXCEPTION 'simulated second GUC failure'; END; $$ LANGUAGE plpgsql",
&[],
)
.await
.expect("create failing wrapper for second GUC");
let result = super::bind_runner_node_identity(&mut ctx, TEST_NODE_ID).await;
assert!(
matches!(&result, Err(RunnerError::NodeIdentityBindingFailed { node_id, .. }) if *node_id == TEST_NODE_ID),
"expected NodeIdentityBindingFailed for node_id={TEST_NODE_ID}, got: {result:?}",
);
let source = std::error::Error::source(&result.expect_err("should be error"))
.expect("NodeIdentityBindingFailed should have a source")
.to_string();
assert!(
source.contains("bind heer.ranj_node_id"),
"error source should identify the second GUC, got: {source}",
);
assert!(
!source.contains("bind heer.node_id"),
"error source should not reference the first GUC (that one succeeded), got: {source}",
);
}
fn single_statement_plan_for_bucket(
bucket_key: &BucketKey,
version: &str,
label: &str,
sql: &str,
) -> (MigrationPlan, RunnerCtx) {
let plan = MigrationPlan {
bucket: bucket_key.clone(),
classification: Classification::Additive,
segments: vec![Segment {
kind: SegmentKind::Transactional,
statements: vec![op(label, sql)],
}],
};
let runner_ctx = RunnerCtx {
bucket: bucket_key.clone(),
version: version.to_string(),
description: "per-app version stream test".to_string(),
checksum_up: compute_checksum_for_plan_up(&plan),
checksum_down: None,
snapshot: Some(empty_snapshot()),
snapshot_path: None,
config: MigrateConfig::default(),
out_of_order_policy: crate::migrate::policy::OutOfOrderPolicy::AllowWithDiagnostic,
audit_pool: None,
runner_identity: Some(RunnerIdentity::SingleNodeDev),
drift_baseline: DriftBaseline::Disabled,
};
(plan, runner_ctx)
}
#[djogi_test]
async fn same_version_applies_once_per_app_label(mut ctx: DjogiContext) {
let guard = acquire_test_workspace_guard();
let version = "V20260609000000__t397";
let bucket_a = BucketKey {
database: "main".into(),
app: "users".into(),
};
let bucket_b = BucketKey {
database: "main".into(),
app: "system".into(),
};
let _ = ctx.execute("DROP TABLE IF EXISTS t397_users", &[]).await;
let _ = ctx.execute("DROP TABLE IF EXISTS t397_system", &[]).await;
let (plan_a, runner_ctx_a) = single_statement_plan_for_bucket(
&bucket_a,
version,
"AddTable t397_users",
"CREATE TABLE t397_users (id BIGINT PRIMARY KEY)",
);
let (plan_b, runner_ctx_b) = single_statement_plan_for_bucket(
&bucket_b,
version,
"AddTable t397_system",
"CREATE TABLE t397_system (id BIGINT PRIMARY KEY)",
);
apply_plan(&mut ctx, &plan_a, &runner_ctx_a, &guard)
.await
.expect("bucket A applies");
apply_plan(&mut ctx, &plan_b, &runner_ctx_b, &guard)
.await
.expect("bucket B must apply on its own (version, app_label) stream");
let rows = super::ledger::select_all(&mut ctx)
.await
.expect("ledger readable");
let same_version: Vec<_> = rows.iter().filter(|r| r.version == version).collect();
assert_eq!(same_version.len(), 2, "one ledger row per app stream");
assert!(same_version.iter().any(|r| r.app_label == "users"));
assert!(same_version.iter().any(|r| r.app_label == "system"));
}
#[allow(clippy::await_holding_lock)]
#[djogi_test]
async fn apply_proceeds_with_gate_disabled(mut ctx: DjogiContext) {
static COUNTER: AtomicUsize = AtomicUsize::new(0);
let n = COUNTER.fetch_add(1, Ordering::SeqCst);
let first_table = format!("drift_disabled_a_{n}");
let drift_table = format!("drift_disabled_extra_{n}");
let second_table = format!("drift_disabled_b_{n}");
let guard = acquire_test_workspace_guard();
let first_plan = single_table_plan(&first_table);
let first_ctx = runner_ctx_with_drift_baseline(
&first_plan,
"V20260613000010__drift_disabled_a",
DriftBaseline::Disabled,
);
apply_plan(&mut ctx, &first_plan, &first_ctx, &guard)
.await
.expect("first apply establishes history");
ctx.batch_execute(&format!("CREATE TABLE {drift_table} (id bigint)"))
.await
.expect("create out-of-band drift table");
let second_plan = single_table_plan(&second_table);
let second_ctx = runner_ctx_with_drift_baseline(
&second_plan,
"V20260613000011__drift_disabled_b",
DriftBaseline::Disabled,
);
apply_plan(&mut ctx, &second_plan, &second_ctx, &guard)
.await
.expect("apply with DriftBaseline::Disabled must proceed despite live drift");
let created: Option<String> = ctx
.query_one(
&format!("SELECT to_regclass('public.{second_table}')::text"),
&[],
)
.await
.expect("regclass lookup")
.try_get(0)
.expect("regclass column");
assert!(
created.is_some(),
"second migration table must exist after a gate-disabled apply"
);
}
#[allow(clippy::await_holding_lock)]
#[djogi_test]
async fn apply_skips_drift_gate_on_never_applied_bucket(mut ctx: DjogiContext) {
static COUNTER: AtomicUsize = AtomicUsize::new(0);
let n = COUNTER.fetch_add(1, Ordering::SeqCst);
let table = format!("drift_skip_snapshot_{n}");
let plan = single_table_plan(&table);
let guard = acquire_test_workspace_guard();
let runner_ctx = runner_ctx_with_drift_baseline(
&plan,
"V20260613000020__drift_skip_snapshot",
DriftBaseline::Snapshot(drift_baseline_with_table("never_existed_table")),
);
apply_plan(&mut ctx, &plan, &runner_ctx, &guard)
.await
.expect(
"never-applied bucket must skip the gate even with a drifting Snapshot baseline",
);
let created: Option<String> = ctx
.query_one(&format!("SELECT to_regclass('public.{table}')::text"), &[])
.await
.expect("regclass lookup")
.try_get(0)
.expect("regclass column");
assert!(created.is_some(), "first migration table must be created");
}
#[allow(clippy::await_holding_lock)]
#[djogi_test]
async fn apply_skips_gate_for_missing_baseline_on_never_applied_bucket(mut ctx: DjogiContext) {
static COUNTER: AtomicUsize = AtomicUsize::new(0);
let n = COUNTER.fetch_add(1, Ordering::SeqCst);
let table = format!("drift_skip_missing_{n}");
let plan = single_table_plan(&table);
let guard = acquire_test_workspace_guard();
let runner_ctx = runner_ctx_with_drift_baseline(
&plan,
"V20260613000030__drift_skip_missing",
DriftBaseline::Missing,
);
apply_plan(&mut ctx, &plan, &runner_ctx, &guard)
.await
.expect("never-applied bucket must skip the gate for a Missing baseline");
}
#[allow(clippy::await_holding_lock)]
#[djogi_test]
async fn apply_skips_drift_gate_on_never_applied_bucket_with_corrupted_baseline(
mut ctx: DjogiContext,
) {
static COUNTER: AtomicUsize = AtomicUsize::new(0);
let n = COUNTER.fetch_add(1, Ordering::SeqCst);
let table = format!("drift_corrupted_skip_{n}");
let plan = single_table_plan(&table);
let guard = acquire_test_workspace_guard();
let runner_ctx = runner_ctx_with_drift_baseline(
&plan,
"V20260613000070__drift_corrupted_skip",
DriftBaseline::Corrupted("fake corruption error".to_string()),
);
apply_plan(&mut ctx, &plan, &runner_ctx, &guard)
.await
.expect("first apply with Corrupted baseline must succeed on never-applied bucket");
}
#[allow(clippy::await_holding_lock)]
#[djogi_test]
async fn apply_skips_gate_for_new_bucket_in_established_database(mut ctx: DjogiContext) {
static COUNTER: AtomicUsize = AtomicUsize::new(0);
let n = COUNTER.fetch_add(1, Ordering::SeqCst);
let alpha_table = format!("drift_new_bucket_alpha_{n}");
let beta_table = format!("drift_new_bucket_beta_{n}");
let beta_table_missing = format!("drift_new_bucket_beta_missing_{n}");
let guard = acquire_test_workspace_guard();
let alpha_plan = named_bucket_single_table_plan("djogi_test", "alpha", &alpha_table);
let alpha_ctx = runner_ctx_with_drift_baseline(
&alpha_plan,
"V20260613000040__drift_new_bucket_alpha",
DriftBaseline::Disabled,
);
apply_plan(&mut ctx, &alpha_plan, &alpha_ctx, &guard)
.await
.expect("alpha apply establishes history in the database");
let beta_plan = named_bucket_single_table_plan("djogi_test", "beta", &beta_table);
let beta_ctx = runner_ctx_with_drift_baseline(
&beta_plan,
"V20260613000041__drift_new_bucket_beta",
DriftBaseline::Snapshot(drift_baseline_with_table("never_existed_table")),
);
apply_plan(&mut ctx, &beta_plan, &beta_ctx, &guard)
.await
.expect("new bucket beta must skip the gate despite alpha's applied history");
let gamma_plan = named_bucket_single_table_plan("djogi_test", "gamma", &beta_table_missing);
let gamma_ctx = runner_ctx_with_drift_baseline(
&gamma_plan,
"V20260613000042__drift_new_bucket_gamma",
DriftBaseline::Missing,
);
apply_plan(&mut ctx, &gamma_plan, &gamma_ctx, &guard)
.await
.expect("new bucket gamma must skip the gate for a Missing baseline");
}
#[allow(clippy::await_holding_lock)]
#[djogi_test]
async fn apply_proceeds_on_advisory_only_divergence(mut ctx: DjogiContext) {
static COUNTER: AtomicUsize = AtomicUsize::new(0);
let n = COUNTER.fetch_add(1, Ordering::SeqCst);
let table = format!("drift_advisory_{n}");
let extra_index = format!("{table}_extra_idx");
let second_table = format!("drift_advisory_b_{n}");
let guard = acquire_test_workspace_guard();
let setup_plan = single_table_plan(&table);
let setup_ctx = runner_ctx_with_drift_baseline(
&setup_plan,
"V20260613000050__drift_advisory_setup",
DriftBaseline::Disabled,
);
apply_plan(&mut ctx, &setup_plan, &setup_ctx, &guard)
.await
.expect("setup apply establishes the table");
let baseline =
crate::migrate::verify::live_schema_for_repair(&mut ctx, &setup_plan.bucket, None)
.await
.expect("project baseline from live schema");
ctx.batch_execute(&format!("CREATE INDEX {extra_index} ON {table} (id)"))
.await
.expect("create out-of-band advisory index");
let next_plan = single_table_plan(&second_table);
let next_ctx = runner_ctx_with_drift_baseline(
&next_plan,
"V20260613000051__drift_advisory_next",
DriftBaseline::Snapshot(baseline),
);
apply_plan(&mut ctx, &next_plan, &next_ctx, &guard)
.await
.expect("advisory-only (warning) divergence must not refuse apply");
let created: Option<String> = ctx
.query_one(
&format!("SELECT to_regclass('public.{second_table}')::text"),
&[],
)
.await
.expect("regclass lookup")
.try_get(0)
.expect("regclass column");
assert!(
created.is_some(),
"second migration table must exist after an advisory-only proceed"
);
}
#[allow(clippy::await_holding_lock)]
#[djogi_test]
async fn apply_over_partial_non_tx_progress_refuses_as_drift(mut ctx: DjogiContext) {
static COUNTER: AtomicUsize = AtomicUsize::new(0);
let n = COUNTER.fetch_add(1, Ordering::SeqCst);
let setup_table = format!("drift_partial_setup_{n}");
let committed_table = format!("drift_partial_committed_{n}");
let broken_table = format!("drift_partial_broken_{n}");
let guard = acquire_test_workspace_guard();
let setup_plan = single_table_plan(&setup_table);
let setup_ctx = runner_ctx_with_drift_baseline(
&setup_plan,
"V20260613000060__drift_partial_setup",
DriftBaseline::Disabled,
);
apply_plan(&mut ctx, &setup_plan, &setup_ctx, &guard)
.await
.expect("setup apply establishes the bucket");
let baseline =
crate::migrate::verify::live_schema_for_repair(&mut ctx, &setup_plan.bucket, None)
.await
.expect("project baseline from live schema");
let partial_version = "V20260613000061__drift_partial_progress";
let partial_plan = MigrationPlan {
bucket: bucket("main", ""),
classification: Classification::Additive,
segments: vec![
Segment {
kind: SegmentKind::NonTransactional,
statements: vec![op(
&format!("AddTable {committed_table}"),
&format!("CREATE TABLE {committed_table} (id bigint)"),
)],
},
Segment {
kind: SegmentKind::NonTransactional,
statements: vec![op(
&format!("AddTable {broken_table}"),
&format!("CREATE TABLE {broken_table} (id THIS_IS_NOT_A_TYPE)"),
)],
},
],
};
let partial_ctx =
runner_ctx_with_drift_baseline(&partial_plan, partial_version, DriftBaseline::Disabled);
let partial_result = apply_plan(&mut ctx, &partial_plan, &partial_ctx, &guard).await;
assert!(
matches!(
partial_result,
Err(RunnerError::NonTransactionalSegmentFailed { .. })
),
"second non-tx segment must fail, leaving partial progress, got: {partial_result:?}"
);
let committed: Option<String> = ctx
.query_one(
&format!("SELECT to_regclass('public.{committed_table}')::text"),
&[],
)
.await
.expect("regclass lookup")
.try_get(0)
.expect("regclass column");
assert!(
committed.is_some(),
"the first non-tx statement must have committed its table"
);
ctx.batch_execute(&format!(
"DELETE FROM djogi_schema_migrations \
WHERE version = '{partial_version}' AND app_label = ''"
))
.await
.expect("clear failed partial ledger row");
let reapply_ctx = runner_ctx_with_drift_baseline(
&partial_plan,
partial_version,
DriftBaseline::Snapshot(baseline),
);
let reapply_result = apply_plan(&mut ctx, &partial_plan, &reapply_ctx, &guard).await;
assert!(
matches!(reapply_result, Err(RunnerError::DriftDetected { .. })),
"re-apply over committed partial non-tx progress must refuse as drift, got: {reapply_result:?}"
);
}
}