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;
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 },
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,
OutOfOrderRejected {
version: String,
conflicting_version: String,
conflicting_applied_at: Option<String>,
},
BaselineProjectionFailed {
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,
},
}
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::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; re-running `djogi migrations apply` will remove the rolled-back row and re-apply the migration"
}
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::BaselineProjectionFailed { source } => write!(
f,
"baseline live-DB projection failed before ledger insert: {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",
),
}
}
}
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::BaselineProjectionFailed { source } => Some(source.as_ref()),
RunnerError::CatalogQueryFailed { source, .. } => Some(source),
RunnerError::LedgerQueryFailed { source, .. } => Some(source),
RunnerError::RunIdGenerationFailed { source } => Some(source),
RunnerError::LedgerBootstrapFailed { source } => Some(source),
RunnerError::PinnedSessionCheckoutFailed { source } => Some(source),
RunnerError::VersionCollisionNonTerminal { .. } => None,
_ => None,
}
}
}
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>,
}
#[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,
}
pub async fn apply_plan(
ctx: &mut DjogiContext,
plan: &MigrationPlan,
runner_ctx: &RunnerCtx,
_guard: &WorkspaceGuard,
) -> Result<RunReport, RunnerError> {
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",
);
}
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
};
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).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;
}
}
}
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)]
pub enum RollbackError {
Runner(RunnerError),
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,
source: DjogiError,
},
PriorSnapshotMissing,
LeafIdentityMismatch {
version: String,
stored_leaf_identity: String,
current_leaf_identity: String,
},
SnapshotPersistFailed {
path: PathBuf,
source: SnapshotError,
},
}
impl std::fmt::Display for RollbackError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
RollbackError::Runner(e) => write!(f, "rollback failed at runner level: {e}"),
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::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}`",
),
}
}
}
impl std::error::Error for RollbackError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
RollbackError::Runner(e) => Some(e),
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);
}
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()) {
(_, 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 }, false) => Ok(Some(reason.clone())),
}
}
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 row_result = load_ledger_row_for_version(ctx, &runner_ctx.version)
.await
.map_err(|e| {
RollbackError::Runner(RunnerError::LedgerQueryFailed {
query_label: "load_row_for_version",
source: e,
})
});
let row_opt = match row_result {
Ok(r) => r,
Err(e) => {
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(e);
}
};
let row = match row_opt {
Some(r) => r,
None => {
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(RollbackError::VersionNotFound {
version: runner_ctx.version.clone(),
});
}
};
if row.app_label != plan.bucket.app {
let e = RollbackError::BucketAppMismatch {
version: runner_ctx.version.clone(),
row_app_label: row.app_label.clone(),
supplied_app: plan.bucket.app.clone(),
};
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(e);
}
if !matches!(row.status, LedgerStatus::Applied | LedgerStatus::Faked) {
let current_status = row.status;
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(RollbackError::VersionNotRollbackable {
version: runner_ctx.version.clone(),
current_status,
});
}
if let Some(ref stored_identity) = row.leaf_identity {
let pre_cache = match compute_leaf_identity_cache(ctx, plan).await {
Ok(c) => c,
Err(e) => {
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(RollbackError::Runner(e));
}
};
let pre_identity = serialize_leaf_identity(&pre_cache).unwrap_or_default();
if pre_identity != *stored_identity {
let e = RollbackError::LeafIdentityMismatch {
version: runner_ctx.version.clone(),
stored_leaf_identity: stored_identity.clone(),
current_leaf_identity: pre_identity,
};
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(e);
}
}
let (replay_plan, leaves_cache_rollback) =
match materialize_execution_plan(ctx, plan, PartitionExpansionMode::ReplayStrict).await {
Ok((p, cache)) => (p, cache),
Err(e) => {
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(RollbackError::Runner(e));
}
};
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 {
let e = RollbackError::LeafIdentityMismatch {
version: runner_ctx.version.clone(),
stored_leaf_identity: row.leaf_identity.clone().unwrap(),
current_leaf_identity: current_rollback_identity,
};
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(e);
}
let allow_reason = match rollback_lossy_allow_reason(&replay_plan, &lossy_policy) {
Ok(r) => r,
Err(e) => {
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(e);
}
};
let result = rollback_inner(ctx, &replay_plan, runner_ctx, prior_snapshot, allow_reason).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(
RunnerError::AdvisoryUnlockReturnedFalse {
key: lock_key,
bucket: plan.bucket.clone(),
},
)),
(Err(e), _) => Err(e),
}
}
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(),
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(),
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(),
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",
&[&runner_ctx.version, ¬e],
)
.await
.map_err(|e| {
RollbackError::Runner(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(),
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> {
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,
};
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).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);
}
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)
.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,
);
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).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,
) -> Result<Option<LedgerRow>, DjogiError> {
load_full_row_by_version(ctx, version).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 {
ctx.batch_execute(sql).await
} else {
guarded_batch_execute(ctx, sql).await
}
}
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}"),
}
}
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,
) -> RunnerError {
let row = match load_ledger_row_for_version(ctx, version).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).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 load_applied_at(ctx: &mut DjogiContext, version: &str) -> Option<OffsetDateTime> {
let row = ctx
.query_opt(
"SELECT applied_at FROM djogi_schema_migrations WHERE version = $1",
&[&version],
)
.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;
use std::path::PathBuf;
use std::time::Duration;
use crate::config::MigrateConfig;
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,
}
}
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")
}
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");
}
}
}
}
#[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,
"re-running `djogi migrations apply` will remove the rolled-back row and re-apply the migration",
),
] {
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),
};
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,
};
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,
};
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,
};
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"));
}
}