use std::path::{Path, PathBuf};
use crate::__bypass::guarded_batch_execute;
use crate::context::{DjogiContext, PinnedCtx};
use crate::error::DjogiError;
use super::bootstrap::PHASE_ZERO_VERSION;
use super::guard::WorkspaceGuard;
use super::ledger::{
self, ChecksumFormatErrorKind, LEDGER_SELECT_COLS, LedgerRow, LedgerStatus, compute_checksum,
validate_checksum_format,
};
use super::projection::BucketKey;
use super::runner::{
PartitionExpansionMode, RunnerError, RunnerIdentity, acquire_advisory_lock, advisory_lock_key,
compute_leaf_identity_cache, materialize_execution_plan, release_advisory_lock,
serialize_leaf_identity,
};
use super::segment::{MigrationPlan, SegmentKind};
use super::snapshot::{SnapshotError, save_snapshot};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RepairConfirmation {
OperatorAcknowledged,
}
#[derive(Debug, Clone)]
pub struct RepairReport {
pub actions_taken: Vec<String>,
pub ledger_changes: Vec<LedgerChange>,
pub snapshot_changes: Vec<SnapshotChange>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LedgerChange {
pub version: String,
pub column: &'static str,
pub before: String,
pub after: String,
}
impl LedgerChange {
pub(crate) fn new(version: &str, column: &'static str, before: String, after: String) -> Self {
Self {
version: version.to_string(),
column,
before,
after,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SnapshotChange {
pub path: PathBuf,
pub description: String,
}
#[derive(Debug)]
pub enum RepairError {
VersionNotFound { version: String },
InsufficientConfirmation,
InvalidChecksum {
value: String,
kind: ChecksumFormatErrorKind,
},
InvalidResolution {
version: String,
current_status: LedgerStatus,
attempted: PartialApplyResolution,
},
BucketAppMismatch {
version: String,
row_app_label: String,
supplied_app: String,
},
LedgerIo { source: DjogiError },
SnapshotIo {
path: PathBuf,
source: SnapshotError,
},
PlanVersionMismatch { ledger: String, plan: String },
PlanChecksumMismatch {
version: String,
ledger_checksum: String,
plan_checksum: String,
},
LeafIdentityMismatch {
version: String,
stored_leaf_identity: String,
current_leaf_identity: String,
},
NothingToResume {
version: String,
applied: i32,
total: Option<i32>,
},
ResumeStepFailed {
version: String,
step_index: usize,
statement_label: String,
applied_steps_count: i32,
source: DjogiError,
},
ResumeBlockedByNonTxProgressClaim { version: String, note: String },
ResumeProgressAckFailed {
version: String,
step_index: usize,
statement_label: String,
applied_steps_count: i32,
source: DjogiError,
},
SuppliedSnapshotDiverges { differences: Vec<String> },
AdvisoryLockFailed {
bucket: BucketKey,
key: i64,
attempts: u32,
},
AdvisoryLockQueryFailed {
app_label: String,
source: DjogiError,
},
AdvisoryUnlockReturnedFalse {
key: i64,
},
PinnedSessionCheckoutFailed {
source: DjogiError,
},
ResumePlanShapeMismatch {
version: String,
ledger_total_steps: usize,
replay_total_steps: usize,
},
ReplayPlanShapeMismatch {
version: String,
expected_step_count: usize,
actual_step_count: usize,
},
PhaseZeroArtifactRefused {
version: String,
},
MissingResumeIdentity {
version: String,
},
Runner(RunnerError),
}
impl std::fmt::Display for RepairError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
RepairError::VersionNotFound { version } => write!(
f,
"repair could not find a ledger row for version `{version}`"
),
RepairError::InsufficientConfirmation => f.write_str(
"repair refused: caller did not supply RepairConfirmation::OperatorAcknowledged",
),
RepairError::InvalidChecksum { value, kind } => {
write!(f, "repair rejected new checksum `{value}`: {kind}")
}
RepairError::InvalidResolution {
version,
current_status,
attempted,
} => write!(
f,
"repair resolution {attempted:?} is not valid for version `{version}` \
(current status: {current})",
current = current_status.as_db_str(),
),
RepairError::BucketAppMismatch {
version,
row_app_label,
supplied_app,
} => write!(
f,
"repair refused for version `{version}`: supplied bucket app `{supplied_app}` \
does not match ledger row app_label `{row_app_label}`; repair would hold \
the advisory lock for the wrong bucket",
),
RepairError::LedgerIo { source } => write!(f, "repair ledger I/O failed: {source}"),
RepairError::SnapshotIo { path, source } => {
write!(
f,
"repair snapshot I/O at {} failed: {source}",
path.display()
)
}
RepairError::PlanVersionMismatch { ledger, plan } => write!(
f,
"repair_resume_partial_apply: plan version `{plan}` does not match \
the ledger row's version `{ledger}`",
),
RepairError::PlanChecksumMismatch {
version,
ledger_checksum,
plan_checksum,
} => write!(
f,
"repair_resume_partial_apply on `{version}`: plan recomputed \
checksum_up={plan_checksum} differs from ledger {ledger_checksum}; \
the supplied plan does not match the original apply",
),
RepairError::NothingToResume {
version,
applied,
total,
} => match total {
Some(t) => write!(
f,
"repair_resume_partial_apply on `{version}`: applied_steps_count={applied} \
already equals total_steps={t}; nothing to resume",
),
None => write!(
f,
"repair_resume_partial_apply on `{version}`: total_steps is NULL — \
this row was a transactional-only apply with no resumable steps",
),
},
RepairError::ResumeStepFailed {
version,
step_index,
statement_label,
applied_steps_count,
source,
} => write!(
f,
"repair_resume_partial_apply on `{version}`: step {step_index} `{statement_label}` \
failed after {applied_steps_count} successful step(s): {source}",
),
RepairError::ResumeBlockedByNonTxProgressClaim { version, note } => write!(
f,
"repair_resume_partial_apply on `{version}` refused: the ledger row carries an \
outstanding non-tx progress claim, so the next step may already have \
committed. Reconcile the row with `repair_partial_apply` or manual \
inspection before resuming. Current note: {note}",
),
RepairError::ResumeProgressAckFailed {
version,
step_index,
statement_label,
applied_steps_count,
source,
} => write!(
f,
"repair_resume_partial_apply on `{version}`: step {} `{statement_label}` \
committed, but the ledger failed to durably acknowledge \
applied_steps_count={applied_steps_count}; the claim note was preserved and \
automatic resume is now blocked: {source}",
step_index + 1,
),
RepairError::SuppliedSnapshotDiverges { differences } => write!(
f,
"repair_snapshot_rebuild: supplied / rebuilt snapshot diverges from \
the live catalog projection: {differences:?}",
),
RepairError::AdvisoryLockFailed {
bucket,
key,
attempts,
} => write!(
f,
"D274 repair advisory lock for bucket database={db} app={app} \
(key=0x{key:016x}) could not be acquired after {attempts} attempts; \
a concurrent runner or repair invocation holds the lock (GH #274)",
db = bucket.database,
app = bucket.app,
),
RepairError::AdvisoryLockQueryFailed { app_label, source } => write!(
f,
"D274 repair pg_try_advisory_lock query failed for app `{app_label}`: \
{source} (GH #274)",
),
RepairError::AdvisoryUnlockReturnedFalse { key } => write!(
f,
"D274 repair pg_advisory_unlock returned false for key=0x{key:016x}; \
the advisory lock was not held on this session — session-pinning \
correctness failure (GH #274/#280)",
),
RepairError::PinnedSessionCheckoutFailed { source } => write!(
f,
"D274 repair failed to check out a pinned Postgres session from the \
pool before the repair operation began (GH #274): {source}",
),
RepairError::ResumePlanShapeMismatch {
version,
ledger_total_steps,
replay_total_steps,
} => write!(
f,
"D317 repair resume on version {version}: plan shape mismatch — \
ledger total_steps={ledger_total_steps}, \
expanded replay non-transactional statements={replay_total_steps}; \
refusing to replay a divergent plan (GH #317)",
),
RepairError::ReplayPlanShapeMismatch {
version,
expected_step_count,
actual_step_count,
} => write!(
f,
"D317 repair replay finalization on version {version}: \
completed {actual_step_count} step(s) but the materialized plan \
expected {expected_step_count}; plan shape invariant violated (GH #317)",
),
RepairError::LeafIdentityMismatch {
version,
stored_leaf_identity: _,
current_leaf_identity: _,
} => write!(
f,
"[D623] repair refused: partition leaf identity mismatch for `{version}`",
),
RepairError::PhaseZeroArtifactRefused { version } => write!(
f,
"repair refused: Phase 0 artifact for `{version}` is not identity-free \
replay-current; replace with an identity-free Phase 0 artifact before repairing"
),
RepairError::MissingResumeIdentity { version } => write!(
f,
"repair resume-partial refused: version '{version}' requires a binding-capable \
runner identity for SQL replay",
),
RepairError::Runner(e) => write!(f, "repair failed at runner level: {e}"),
}
}
}
impl std::error::Error for RepairError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
RepairError::LedgerIo { source } => Some(source),
RepairError::SnapshotIo { source, .. } => Some(source),
RepairError::ResumeStepFailed { source, .. } => Some(source),
RepairError::ResumeProgressAckFailed { source, .. } => Some(source),
RepairError::AdvisoryLockQueryFailed { source, .. } => Some(source),
RepairError::PinnedSessionCheckoutFailed { source } => Some(source),
RepairError::Runner(e) => Some(e),
_ => None,
}
}
}
async fn acquire_advisory_lock_repair(
ctx: &mut PinnedCtx<'_>,
bucket: &BucketKey,
key: i64,
) -> Result<(), RepairError> {
acquire_advisory_lock(ctx, bucket, key)
.await
.map_err(|e| match e {
super::runner::RunnerError::AdvisoryLockFailed {
bucket,
key,
attempts,
} => RepairError::AdvisoryLockFailed {
bucket,
key,
attempts,
},
super::runner::RunnerError::AdvisoryLockQueryFailed { app_label, source } => {
RepairError::AdvisoryLockQueryFailed { app_label, source }
}
other => {
RepairError::Runner(other)
}
})
}
#[allow(clippy::result_large_err)]
fn handle_repair_release<T>(
result: Result<T, RepairError>,
released: bool,
key: i64,
) -> Result<T, RepairError> {
match result {
Ok(v) if released => Ok(v),
Ok(_) => Err(RepairError::AdvisoryUnlockReturnedFalse { key }),
Err(e) => Err(e),
}
}
#[expect(
clippy::result_large_err,
reason = "repair guard returns the public RepairError enum unchanged"
)]
fn check_phase_zero_repair(
workspace: &Path,
bucket: &BucketKey,
version: &str,
) -> Result<(), RepairError> {
if version != PHASE_ZERO_VERSION {
return Ok(());
}
let bucket_dir = super::target::bucket_dir(workspace, bucket);
let up_path = bucket_dir.join(super::naming::up_filename(version));
let bytes = std::fs::read(&up_path).map_err(|e| RepairError::LedgerIo {
source: DjogiError::Db(crate::error::DbError::other(format!(
"read Phase 0 up-file at {}: {e}",
up_path.display()
))),
})?;
super::phase_zero::require_identity_free_phase_zero_migration_artifact(&bytes).map_err(|_| {
RepairError::PhaseZeroArtifactRefused {
version: version.to_string(),
}
})
}
#[expect(
clippy::result_large_err,
reason = "repair guard returns the public RepairError enum unchanged"
)]
fn check_phase_zero_materialized_resume_stream(
version: &str,
plan: &MigrationPlan,
) -> Result<(), RepairError> {
if version != PHASE_ZERO_VERSION {
return Ok(());
}
let combined_up = plan
.segments
.iter()
.flat_map(|seg| seg.statements.iter())
.map(|stmt| stmt.up.as_str())
.collect::<Vec<_>>()
.join("\n");
super::phase_zero::require_identity_free_phase_zero_migration_artifact(combined_up.as_bytes())
.map_err(|_| RepairError::PhaseZeroArtifactRefused {
version: version.to_string(),
})
}
#[expect(
clippy::too_many_arguments,
reason = "public repair API keeps explicit bucket, workspace, checksum, and confirmation inputs"
)]
pub async fn repair_checksum_drift(
ctx: &mut DjogiContext,
_guard: &WorkspaceGuard,
bucket: &BucketKey,
version: &str,
workspace: &Path,
new_checksum_up: &str,
new_checksum_down: Option<&str>,
confirmation: RepairConfirmation,
) -> Result<RepairReport, RepairError> {
if confirmation != RepairConfirmation::OperatorAcknowledged {
return Err(RepairError::InsufficientConfirmation);
}
check_phase_zero_repair(workspace, bucket, version)?;
if let Err(kind) = validate_checksum_format(new_checksum_up) {
return Err(RepairError::InvalidChecksum {
value: new_checksum_up.to_string(),
kind,
});
}
if let Some(down) = new_checksum_down
&& let Err(kind) = validate_checksum_format(down)
{
return Err(RepairError::InvalidChecksum {
value: down.to_string(),
kind,
});
}
let lock_key = advisory_lock_key(bucket);
let mut pinned = ctx
.pin_for_migration()
.await
.map_err(|e| RepairError::PinnedSessionCheckoutFailed { source: e })?;
repair_checksum_drift_pinned(
&mut pinned,
version,
new_checksum_up,
new_checksum_down,
bucket,
lock_key,
)
.await
}
async fn repair_checksum_drift_pinned(
ctx: &mut PinnedCtx<'_>,
version: &str,
new_checksum_up: &str,
new_checksum_down: Option<&str>,
bucket: &BucketKey,
lock_key: i64,
) -> Result<RepairReport, RepairError> {
acquire_advisory_lock_repair(ctx, bucket, lock_key).await?;
let row = match load_row(ctx, version, &bucket.app).await {
Ok(r) => r,
Err(e) => {
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(e);
}
};
if let Err(e) = ensure_row_matches_bucket_app(&row, bucket, version) {
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(e);
}
let before_up = row.checksum_up.clone();
let before_down = row.checksum_down.clone();
let new_down_owned = new_checksum_down.map(|s| s.to_string());
let mutation_result = ctx
.execute(
"UPDATE djogi_schema_migrations \
SET checksum_up = $2, checksum_down = $3 \
WHERE version = $1 AND app_label = $4",
&[&version, &new_checksum_up, &new_down_owned, &bucket.app],
)
.await
.map_err(|e| RepairError::LedgerIo { source: e })
.map(|_| ());
let released = release_advisory_lock(ctx, lock_key).await;
handle_repair_release(mutation_result, released, lock_key)?;
ctx.mark_clean();
let mut actions = vec![format!(
"checksum_up of `{version}` updated from {before_up} to {new_checksum_up}"
)];
let mut changes = vec![LedgerChange::new(
version,
"checksum_up",
before_up,
new_checksum_up.to_string(),
)];
let after_down_render = new_down_owned
.clone()
.unwrap_or_else(|| "<none>".to_string());
let before_down_render = before_down.clone().unwrap_or_else(|| "<none>".to_string());
actions.push(format!(
"checksum_down of `{version}` updated from {before_down_render} to {after_down_render}"
));
changes.push(LedgerChange::new(
version,
"checksum_down",
before_down_render,
after_down_render,
));
Ok(RepairReport {
actions_taken: actions,
ledger_changes: changes,
snapshot_changes: Vec::new(),
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PartialApplyResolution {
MarkRolledBack,
MarkFaked,
MarkApplied,
}
#[expect(
clippy::too_many_arguments,
reason = "public repair API keeps explicit bucket, workspace, resolution, note, and confirmation inputs"
)]
pub async fn repair_partial_apply(
ctx: &mut DjogiContext,
_guard: &WorkspaceGuard,
bucket: &BucketKey,
version: &str,
workspace: &Path,
resolution: PartialApplyResolution,
note: &str,
confirmation: RepairConfirmation,
) -> Result<RepairReport, RepairError> {
if confirmation != RepairConfirmation::OperatorAcknowledged {
return Err(RepairError::InsufficientConfirmation);
}
check_phase_zero_repair(workspace, bucket, version)?;
let lock_key = advisory_lock_key(bucket);
let mut pinned = ctx
.pin_for_migration()
.await
.map_err(|e| RepairError::PinnedSessionCheckoutFailed { source: e })?;
repair_partial_apply_pinned(&mut pinned, version, resolution, note, bucket, lock_key).await
}
async fn repair_partial_apply_pinned(
ctx: &mut PinnedCtx<'_>,
version: &str,
resolution: PartialApplyResolution,
note: &str,
bucket: &BucketKey,
lock_key: i64,
) -> Result<RepairReport, RepairError> {
acquire_advisory_lock_repair(ctx, bucket, lock_key).await?;
let row = match load_row(ctx, version, &bucket.app).await {
Ok(r) => r,
Err(e) => {
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(e);
}
};
if let Err(e) = ensure_row_matches_bucket_app(&row, bucket, version) {
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(e);
}
if !matches!(row.status, LedgerStatus::Failed | LedgerStatus::Pending) {
let current_status = row.status;
let _released = release_advisory_lock(ctx, lock_key).await;
return Err(RepairError::InvalidResolution {
version: version.to_string(),
current_status,
attempted: resolution,
});
}
let target_status = match resolution {
PartialApplyResolution::MarkRolledBack => LedgerStatus::RolledBack,
PartialApplyResolution::MarkFaked => LedgerStatus::Faked,
PartialApplyResolution::MarkApplied => LedgerStatus::Applied,
};
let target_steps = match resolution {
PartialApplyResolution::MarkApplied => row.total_steps.unwrap_or(row.applied_steps_count),
_ => row.applied_steps_count,
};
let status_str = target_status.as_db_str();
let mutation_result = ctx
.execute(
"UPDATE djogi_schema_migrations \
SET status = $2, applied_steps_count = $3, partial_apply_note = $4 \
WHERE version = $1 AND app_label = $5",
&[&version, &status_str, &target_steps, ¬e, &bucket.app],
)
.await
.map_err(|e| RepairError::LedgerIo { source: e })
.map(|_| ());
let released = release_advisory_lock(ctx, lock_key).await;
handle_repair_release(mutation_result, released, lock_key)?;
ctx.mark_clean();
Ok(RepairReport {
actions_taken: vec![format!(
"partial-apply repair of `{version}`: status {old} -> {new}; \
applied_steps_count {old_steps} -> {new_steps}; note set",
old = row.status.as_db_str(),
new = target_status.as_db_str(),
old_steps = row.applied_steps_count,
new_steps = target_steps,
)],
ledger_changes: vec![
LedgerChange::new(
version,
"status",
row.status.as_db_str().to_string(),
target_status.as_db_str().to_string(),
),
LedgerChange::new(
version,
"applied_steps_count",
row.applied_steps_count.to_string(),
target_steps.to_string(),
),
LedgerChange::new(
version,
"partial_apply_note",
row.partial_apply_note.clone().unwrap_or_default(),
note.to_string(),
),
],
snapshot_changes: Vec::new(),
})
}
pub async fn repair_resume_partial_apply(
ctx: &mut DjogiContext,
_guard: &WorkspaceGuard,
workspace: &Path,
version: &str,
plan: &MigrationPlan,
runner_identity: Option<RunnerIdentity>,
confirmation: RepairConfirmation,
) -> Result<RepairReport, RepairError> {
if confirmation != RepairConfirmation::OperatorAcknowledged {
return Err(RepairError::InsufficientConfirmation);
}
check_phase_zero_repair(workspace, &plan.bucket, version)?;
if version != PHASE_ZERO_VERSION {
let identity = match runner_identity {
Some(id) => id,
None => {
return Err(RepairError::MissingResumeIdentity {
version: version.to_string(),
});
}
};
if !identity.requires_binding() {
return Err(RepairError::MissingResumeIdentity {
version: version.to_string(),
});
}
}
let lock_key = advisory_lock_key(&plan.bucket);
let mut pinned = ctx
.pin_for_migration()
.await
.map_err(|e| RepairError::PinnedSessionCheckoutFailed { source: e })?;
if version != PHASE_ZERO_VERSION {
let identity = runner_identity
.expect("INVARIANT: non-Phase-0 runner identity was validated before pinning");
super::runner::bind_runner_node_identity(
&mut pinned,
identity.node_id().expect(
"INVARIANT: node_id is Some when requires_binding is true; \
IdentityFree ruled out by refusal gate above",
),
)
.await
.map_err(RepairError::Runner)?;
}
repair_resume_pinned(&mut pinned, version, plan, &plan.bucket, lock_key).await
}
async fn repair_resume_pinned(
ctx: &mut PinnedCtx<'_>,
version: &str,
plan: &MigrationPlan,
bucket: &BucketKey,
lock_key: i64,
) -> Result<RepairReport, RepairError> {
acquire_advisory_lock_repair(ctx, bucket, lock_key).await?;
let result = repair_resume_body(ctx, version, plan, bucket).await;
let released = release_advisory_lock(ctx, lock_key).await;
let result = handle_repair_release(result, released, lock_key);
if result.is_ok() {
ctx.mark_clean();
}
result
}
async fn repair_resume_body(
ctx: &mut DjogiContext,
version: &str,
plan: &MigrationPlan,
bucket: &BucketKey,
) -> Result<RepairReport, RepairError> {
let row = load_row(ctx, version, &bucket.app).await?;
ensure_row_matches_bucket_app(&row, bucket, version)?;
let plan_checksum = compute_plan_checksum_up(plan);
if plan_checksum != row.checksum_up {
return Err(RepairError::PlanChecksumMismatch {
version: version.to_string(),
ledger_checksum: row.checksum_up.clone(),
plan_checksum,
});
}
if !matches!(row.status, LedgerStatus::Failed | LedgerStatus::Pending) {
return Err(RepairError::InvalidResolution {
version: version.to_string(),
current_status: row.status,
attempted: PartialApplyResolution::MarkApplied,
});
}
if ledger::note_has_non_tx_progress_claim(row.partial_apply_note.as_deref()) {
return Err(RepairError::ResumeBlockedByNonTxProgressClaim {
version: version.to_string(),
note: row.partial_apply_note.clone().unwrap_or_default(),
});
}
let total = match row.total_steps {
Some(t) => t,
None => {
return Err(RepairError::NothingToResume {
version: version.to_string(),
applied: row.applied_steps_count,
total: None,
});
}
};
if row.applied_steps_count >= total {
return Err(RepairError::NothingToResume {
version: version.to_string(),
applied: row.applied_steps_count,
total: Some(total),
});
}
let ledger_id = lookup_ledger_id_by_version(ctx, &row.version, &bucket.app).await?;
if let Some(ref stored_identity) = row.leaf_identity {
let pre_cache = compute_leaf_identity_cache(ctx, plan)
.await
.map_err(RepairError::Runner)?;
let pre_identity = serialize_leaf_identity(&pre_cache).unwrap_or_default();
if pre_identity != *stored_identity {
return Err(RepairError::LeafIdentityMismatch {
version: row.version.clone(),
stored_leaf_identity: stored_identity.clone(),
current_leaf_identity: pre_identity,
});
}
}
let (materialized_plan, leaves_cache_repair) =
materialize_execution_plan(ctx, plan, PartitionExpansionMode::ReplayStrict)
.await
.map_err(|e| match e {
RunnerError::PartitionExpansionNoLeaves { .. } => {
RepairError::ResumePlanShapeMismatch {
version: row.version.clone(),
ledger_total_steps: total as usize,
replay_total_steps: 0,
}
}
other => RepairError::Runner(other),
})?;
check_phase_zero_materialized_resume_stream(version, &materialized_plan)?;
if let Some(ref stored_identity) = row.leaf_identity {
let current_identity = serialize_leaf_identity(&leaves_cache_repair).unwrap_or_default();
if current_identity != *stored_identity {
return Err(RepairError::LeafIdentityMismatch {
version: row.version.clone(),
stored_leaf_identity: stored_identity.clone(),
current_leaf_identity: current_identity,
});
}
}
let replay_total_steps: usize = materialized_plan
.segments
.iter()
.filter(|s| s.kind == SegmentKind::NonTransactional)
.map(|s| s.statements.len())
.sum();
if total as usize != replay_total_steps {
return Err(RepairError::ResumePlanShapeMismatch {
version: row.version.clone(),
ledger_total_steps: total as usize,
replay_total_steps,
});
}
let mut remaining_to_skip = row.applied_steps_count as usize;
let mut applied = row.applied_steps_count;
let mut actions: Vec<String> = Vec::new();
for (seg_idx, segment) in materialized_plan.segments.iter().enumerate() {
if segment.kind != SegmentKind::NonTransactional {
continue;
}
for (step_within, stmt) in segment.statements.iter().enumerate() {
if remaining_to_skip > 0 {
remaining_to_skip -= 1;
continue;
}
let claimed_step = applied.saturating_add(1);
let claim_note = ledger::format_non_tx_progress_claim(
row.partial_apply_note.as_deref(),
claimed_step,
row.total_steps,
seg_idx,
&stmt.label,
);
ledger::claim_non_tx_progress(ctx, ledger_id, &claim_note)
.await
.map_err(|e| RepairError::LedgerIo { source: e })?;
if let Err(e) = guarded_batch_execute(ctx, &stmt.up).await {
let note = format!(
"resume failed at segment {seg_idx} step {step_within}: \
{label} — {e}",
label = stmt.label,
);
let _ = ledger::mark_partial(ctx, ledger_id, applied, ¬e).await;
return Err(RepairError::ResumeStepFailed {
version: row.version.clone(),
step_index: applied as usize,
statement_label: stmt.label.clone(),
applied_steps_count: applied,
source: e,
});
}
applied = applied.saturating_add(1);
ledger::ack_non_tx_progress(ctx, ledger_id, applied, row.partial_apply_note.as_deref())
.await
.map_err(|e| RepairError::ResumeProgressAckFailed {
version: row.version.clone(),
step_index: claimed_step.saturating_sub(1) as usize,
statement_label: stmt.label.clone(),
applied_steps_count: applied,
source: e,
})?;
actions.push(format!(
"resumed step {applied} of {total}: {label}",
label = stmt.label,
));
}
}
if applied as usize != replay_total_steps {
return Err(RepairError::ReplayPlanShapeMismatch {
version: row.version.clone(),
expected_step_count: replay_total_steps,
actual_step_count: applied as usize,
});
}
let before_status = row.status.as_db_str().to_string();
let before_steps = row.applied_steps_count.to_string();
let row_version = row.version.clone();
ctx.execute(
"UPDATE djogi_schema_migrations \
SET status = 'applied', applied_steps_count = $2, partial_apply_note = NULL \
WHERE version = $1 AND app_label = $3",
&[&row_version, &applied, &bucket.app],
)
.await
.map_err(|e| RepairError::LedgerIo { source: e })
.map(|_| RepairReport {
actions_taken: actions,
ledger_changes: vec![
LedgerChange::new(
&row_version,
"status",
before_status,
LedgerStatus::Applied.as_db_str().to_string(),
),
LedgerChange::new(
&row_version,
"applied_steps_count",
before_steps,
applied.to_string(),
),
],
snapshot_changes: Vec::new(),
})
}
fn compute_plan_checksum_up(plan: &MigrationPlan) -> String {
let frags: Vec<&str> = plan
.segments
.iter()
.flat_map(|s| s.statements.iter())
.map(|s| s.up.as_str())
.collect();
compute_checksum(frags)
}
async fn lookup_ledger_id_by_version(
ctx: &mut DjogiContext,
version: &str,
app_label: &str,
) -> Result<i64, RepairError> {
let row = ctx
.query_one(
"SELECT id FROM djogi_schema_migrations \
WHERE version = $1 AND app_label = $2",
&[&version, &app_label],
)
.await
.map_err(|e| RepairError::LedgerIo { source: e })?;
let id: i64 = row.try_get(0).map_err(io_err)?;
Ok(id)
}
pub async fn repair_snapshot_rebuild(
ctx: &mut DjogiContext,
_guard: &WorkspaceGuard,
bucket: &BucketKey,
snapshot_path: &std::path::Path,
confirmation: RepairConfirmation,
) -> Result<RepairReport, RepairError> {
if confirmation != RepairConfirmation::OperatorAcknowledged {
return Err(RepairError::InsufficientConfirmation);
}
let lock_key = advisory_lock_key(bucket);
let mut pinned = ctx
.pin_for_migration()
.await
.map_err(|e| RepairError::PinnedSessionCheckoutFailed { source: e })?;
repair_snapshot_rebuild_pinned(&mut pinned, bucket, snapshot_path, lock_key).await
}
async fn repair_snapshot_rebuild_pinned(
ctx: &mut PinnedCtx<'_>,
bucket: &BucketKey,
snapshot_path: &std::path::Path,
lock_key: i64,
) -> Result<RepairReport, RepairError> {
acquire_advisory_lock_repair(ctx, bucket, lock_key).await?;
let projected = super::verify::live_schema_for_repair(ctx, bucket, None)
.await
.map_err(|e| RepairError::LedgerIo {
source: DjogiError::Db(crate::error::DbError::other(format!(
"live-DB projection failed: {e}"
))),
});
let applied_for_bucket = count_applied_for_app(ctx, &bucket.app)
.await
.ok()
.unwrap_or(-1);
let released = release_advisory_lock(ctx, lock_key).await;
let projected = handle_repair_release(projected, released, lock_key)?;
ctx.mark_clean();
save_snapshot(&projected, snapshot_path).map_err(|e| RepairError::SnapshotIo {
path: snapshot_path.to_path_buf(),
source: e,
})?;
let mut actions = Vec::new();
actions.push(format!(
"snapshot rebuilt for bucket database={} app={} -> {}",
bucket.database,
bucket.app,
snapshot_path.display(),
));
actions.push(format!(
"projected {n} table(s), {idx} index(es) from live catalog",
n = projected.models.len(),
idx = projected.indexes.len(),
));
match applied_for_bucket {
-1 => actions
.push("advisory: ledger table not readable; bucket apply-count unknown".to_string()),
0 => actions.push(
"advisory: bucket has 0 applied ledger rows; rebuild recorded \
as the snapshot for a fresh / empty migration history"
.to_string(),
),
n => actions.push(format!("advisory: bucket has {n} applied ledger row(s)")),
}
let description = match applied_for_bucket {
-1 => "rebuilt from live-DB projection (ledger unreadable)".to_string(),
n => format!("rebuilt from live-DB projection ({n} applied rows)"),
};
Ok(RepairReport {
actions_taken: actions,
ledger_changes: Vec::new(),
snapshot_changes: vec![SnapshotChange {
path: snapshot_path.to_path_buf(),
description,
}],
})
}
async fn load_row(
ctx: &mut DjogiContext,
version: &str,
app_label: &str,
) -> Result<LedgerRow, RepairError> {
let pg_row = ctx
.query_opt(
&format!("{LEDGER_SELECT_COLS} WHERE version = $1 AND app_label = $2"),
&[&version, &app_label],
)
.await
.map_err(|e| RepairError::LedgerIo { source: e })?;
if let Some(r) = pg_row {
return LedgerRow::try_from(&r).map_err(io_err);
}
let any_row = ctx
.query_opt(
&format!("{LEDGER_SELECT_COLS} WHERE version = $1 ORDER BY app_label LIMIT 1"),
&[&version],
)
.await
.map_err(|e| RepairError::LedgerIo { source: e })?;
match any_row {
Some(r) => LedgerRow::try_from(&r).map_err(io_err),
None => Err(RepairError::VersionNotFound {
version: version.to_string(),
}),
}
}
#[allow(clippy::result_large_err)]
fn ensure_row_matches_bucket_app(
row: &LedgerRow,
bucket: &BucketKey,
version: &str,
) -> Result<(), RepairError> {
if row.app_label == bucket.app {
return Ok(());
}
Err(RepairError::BucketAppMismatch {
version: version.to_string(),
row_app_label: row.app_label.clone(),
supplied_app: bucket.app.clone(),
})
}
fn io_err(e: tokio_postgres::Error) -> RepairError {
RepairError::LedgerIo {
source: DjogiError::from(e),
}
}
async fn count_applied_for_app(ctx: &mut DjogiContext, app_label: &str) -> Result<i64, DjogiError> {
let row = ctx
.query_one(
"SELECT COUNT(*)::bigint FROM djogi_schema_migrations \
WHERE app_label = $1 AND status = 'applied'",
&[&app_label],
)
.await?;
let n: i64 = row.try_get(0)?;
Ok(n)
}
#[cfg(test)]
mod tests {
use super::*;
use djogi_macros::djogi_test;
use std::collections::BTreeSet;
use std::fs;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
fn temp_root(tag: &str) -> PathBuf {
static COUNTER: AtomicUsize = AtomicUsize::new(0);
let n = COUNTER.fetch_add(1, Ordering::SeqCst);
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
let p = std::env::temp_dir().join(format!("djogi-repair-{tag}-{nanos}-{n}"));
fs::create_dir_all(&p).unwrap();
p
}
fn phase_zero_bucket() -> BucketKey {
BucketKey {
database: "main".to_string(),
app: String::new(),
}
}
async fn unreachable_pool_context() -> DjogiContext {
let pool = crate::pg::pool::DjogiPool::builder("postgres://localhost:1/djogi_unreachable")
.timeout(Duration::from_millis(1))
.build()
.await
.expect("construct lazy unreachable pool");
DjogiContext::from_pool(pool)
}
fn acquire_test_workspace_guard() -> WorkspaceGuard {
crate::migrate::acquire_workspace_lock(
&temp_root("repair-lock").join("djogi.lock"),
Duration::from_secs(2),
)
.expect("acquire workspace lock")
}
fn resume_plan(bucket: BucketKey) -> MigrationPlan {
MigrationPlan {
bucket,
classification: super::super::diff::Classification::Additive,
segments: vec![super::super::segment::Segment {
kind: SegmentKind::NonTransactional,
statements: vec![super::super::sql::OperationSql {
label: "resume step".to_string(),
up: "SELECT 1".to_string(),
down: String::new(),
lossy: None,
}],
}],
}
}
fn current_production_phase_zero_sql() -> String {
super::super::bootstrap::compose_phase_zero(
"main",
&BTreeSet::new(),
super::super::bootstrap::DEFAULT_NODE_ID,
false,
)
.expect("compose production Phase 0")
}
fn markerless_seed_phase_zero_sql() -> String {
phase_zero_with_seed_statement("INSERT INTO heer.heer_nodes (id) VALUES (1);")
}
fn phase_zero_with_seed_statement(statement: &str) -> String {
let mut sql = current_production_phase_zero_sql();
sql.push('\n');
sql.push_str(statement);
sql.push('\n');
sql
}
fn extended_seed_statement_cases() -> [(&'static str, &'static str); 4] {
[
(
"cte_insert",
"WITH rows AS (SELECT 1) INSERT INTO heer.heer_nodes (id) VALUES (1);",
),
(
"cte_delete",
"WITH moved AS (DELETE FROM heer.heer_node_state RETURNING *) SELECT 1;",
),
(
"merge",
"MERGE INTO heer.heer_nodes AS target USING incoming ON false WHEN NOT MATCHED THEN INSERT (id) VALUES (1);",
),
(
"copy_from",
"COPY \"heer\".\"heer_ranj_node_state\" (\"node_id\") FROM STDIN;",
),
]
}
fn incomplete_phase_zero_sql() -> String {
current_production_phase_zero_sql()
.replace("-- HeeRanjID base schema + functions (idempotent).\n", "")
}
fn generated_stale_phase_zero_sql() -> String {
let mut sql = super::super::bootstrap::compose_phase_zero(
"main",
&BTreeSet::new(),
super::super::bootstrap::DEFAULT_NODE_ID,
true,
)
.expect("compose seed-capable Phase 0");
let start = sql.find("DO $djogi$").expect("dynamic defaults block");
let end = sql[start..]
.find("SET heer.node_id = '1';")
.map(|offset| start + offset)
.expect("session SET after dynamic defaults");
sql.replace_range(
start..end,
"ALTER DATABASE \"main\" SET heer.node_id = '1';\n\
ALTER DATABASE \"main\" SET heer.ranj_node_id = '1';\n",
);
sql
}
fn write_phase_zero_artifact(root: &Path, bucket: &BucketKey, sql: &str) {
let dir = super::super::target::bucket_dir(root, bucket);
fs::create_dir_all(&dir).unwrap();
fs::write(
dir.join(super::super::naming::up_filename(PHASE_ZERO_VERSION)),
sql,
)
.unwrap();
}
fn phase_zero_resume_plan_with_up(up: &str) -> MigrationPlan {
MigrationPlan {
bucket: phase_zero_bucket(),
classification: super::super::diff::Classification::Additive,
segments: vec![super::super::segment::Segment {
kind: SegmentKind::NonTransactional,
statements: vec![super::super::sql::OperationSql {
label: "resume Phase 0 step".to_string(),
up: up.to_string(),
down: String::new(),
lossy: None,
}],
}],
}
}
#[test]
fn confirmation_value_is_only_constructible_via_variant_name() {
let c = RepairConfirmation::OperatorAcknowledged;
match c {
RepairConfirmation::OperatorAcknowledged => (),
}
assert_eq!(c, RepairConfirmation::OperatorAcknowledged);
}
#[test]
fn ledger_change_round_trips_through_clone_and_eq() {
let a = LedgerChange {
version: "V1".to_string(),
column: "checksum_up",
before: "old".to_string(),
after: "new".to_string(),
};
assert_eq!(a, a.clone());
}
#[test]
fn snapshot_change_round_trips_through_clone_and_eq() {
let a = SnapshotChange {
path: PathBuf::from("/tmp/x.json"),
description: "rebuilt".to_string(),
};
assert_eq!(a, a.clone());
}
#[test]
fn partial_apply_resolution_distinct_variants() {
let kinds = [
PartialApplyResolution::MarkRolledBack,
PartialApplyResolution::MarkFaked,
PartialApplyResolution::MarkApplied,
];
for (i, a) in kinds.iter().enumerate() {
for (j, b) in kinds.iter().enumerate() {
if i == j {
assert_eq!(a, b);
} else {
assert_ne!(a, b);
}
}
}
}
#[test]
fn invalid_checksum_carries_kind() {
let bad = "V2:notvalid";
let kind = validate_checksum_format(bad).unwrap_err();
match kind {
ChecksumFormatErrorKind::WrongPrefix
| ChecksumFormatErrorKind::WrongLength { .. }
| ChecksumFormatErrorKind::NonLowercaseHex { .. } => (),
}
}
#[test]
fn repair_refuses_missing_phase_zero_artifact() {
let root = temp_root("missing_p0");
let bucket = phase_zero_bucket();
write_phase_zero_artifact(&root, &bucket, " \n\t ");
let result = check_phase_zero_repair(&root, &bucket, PHASE_ZERO_VERSION);
assert!(
matches!(result, Err(RepairError::PhaseZeroArtifactRefused { .. })),
"repair must refuse missing Phase 0 artifacts; got {result:?}"
);
let _ = fs::remove_dir_all(&root);
}
#[test]
fn repair_refuses_incomplete_phase_zero_artifact() {
let root = temp_root("incomplete_p0");
let bucket = phase_zero_bucket();
write_phase_zero_artifact(&root, &bucket, &incomplete_phase_zero_sql());
let result = check_phase_zero_repair(&root, &bucket, PHASE_ZERO_VERSION);
assert!(
matches!(result, Err(RepairError::PhaseZeroArtifactRefused { .. })),
"repair must refuse incomplete Phase 0 artifacts; got {result:?}"
);
let _ = fs::remove_dir_all(&root);
}
#[test]
fn repair_refuses_markerless_seed_phase_zero_artifact() {
let root = temp_root("markerless_seed_p0");
let bucket = phase_zero_bucket();
write_phase_zero_artifact(&root, &bucket, &markerless_seed_phase_zero_sql());
let result = check_phase_zero_repair(&root, &bucket, PHASE_ZERO_VERSION);
assert!(
matches!(result, Err(RepairError::PhaseZeroArtifactRefused { .. })),
"repair must refuse markerless seed Phase 0 artifacts; got {result:?}"
);
let _ = fs::remove_dir_all(&root);
}
#[test]
fn repair_refuses_extended_seed_phase_zero_artifacts() {
for (name, statement) in extended_seed_statement_cases() {
let root = temp_root(&format!("extended_seed_p0_{name}"));
let bucket = phase_zero_bucket();
write_phase_zero_artifact(&root, &bucket, &phase_zero_with_seed_statement(statement));
let result = check_phase_zero_repair(&root, &bucket, PHASE_ZERO_VERSION);
assert!(
matches!(result, Err(RepairError::PhaseZeroArtifactRefused { .. })),
"repair must refuse extended seed Phase 0 artifact {name}; got {result:?}"
);
let _ = fs::remove_dir_all(&root);
}
}
#[test]
fn repair_materialized_resume_stream_refuses_markerless_seed_phase_zero() {
let plan = phase_zero_resume_plan_with_up(&markerless_seed_phase_zero_sql());
let result = check_phase_zero_materialized_resume_stream(PHASE_ZERO_VERSION, &plan);
assert!(
matches!(result, Err(RepairError::PhaseZeroArtifactRefused { .. })),
"repair materialized resume stream must refuse markerless seed Phase 0, got {result:?}"
);
}
#[test]
fn repair_materialized_resume_stream_refuses_extended_seed_phase_zero() {
for (name, statement) in extended_seed_statement_cases() {
let plan = phase_zero_resume_plan_with_up(&phase_zero_with_seed_statement(statement));
let result = check_phase_zero_materialized_resume_stream(PHASE_ZERO_VERSION, &plan);
assert!(
matches!(result, Err(RepairError::PhaseZeroArtifactRefused { .. })),
"repair materialized resume stream must refuse extended seed Phase 0 {name}, got {result:?}"
);
}
}
#[tokio::test]
async fn repair_resume_missing_identity_refuses_before_pin() {
let mut ctx = unreachable_pool_context().await;
let root = temp_root("resume_no_identity");
let bucket = BucketKey {
database: "main".to_string(),
app: String::new(),
};
let plan = resume_plan(bucket);
let guard = acquire_test_workspace_guard();
let result = repair_resume_partial_apply(
&mut ctx,
&guard,
&root,
"V20260602000008__resume_no_identity",
&plan,
None,
RepairConfirmation::OperatorAcknowledged,
)
.await;
assert!(
matches!(result, Err(RepairError::MissingResumeIdentity { .. })),
"missing resume identity must refuse before pin/connection checkout, got {result:?}"
);
let _ = fs::remove_dir_all(&root);
}
#[tokio::test]
async fn repair_resume_phase_zero_artifact_preflight_refuses_before_pin() {
let cases = [
("generated_stale", generated_stale_phase_zero_sql()),
("markerless_seed", markerless_seed_phase_zero_sql()),
(
"cte_seed",
phase_zero_with_seed_statement(
"WITH rows AS (SELECT 1) INSERT INTO heer.heer_nodes (id) VALUES (1);",
),
),
(
"merge_seed",
phase_zero_with_seed_statement(
"MERGE INTO heer.heer_nodes AS target USING incoming ON false WHEN NOT MATCHED THEN INSERT (id) VALUES (1);",
),
),
(
"copy_seed",
phase_zero_with_seed_statement(
"COPY \"heer\".\"heer_ranj_node_state\" (\"node_id\") FROM STDIN;",
),
),
("missing", " \n\t ".to_string()),
("incomplete", incomplete_phase_zero_sql()),
("ambiguous", "SELECT 1".to_string()),
];
for (name, sql) in cases {
let mut ctx = unreachable_pool_context().await;
let root = temp_root(&format!("resume_p0_{name}"));
let bucket = phase_zero_bucket();
write_phase_zero_artifact(&root, &bucket, &sql);
let plan = resume_plan(bucket);
let guard = acquire_test_workspace_guard();
let result = repair_resume_partial_apply(
&mut ctx,
&guard,
&root,
PHASE_ZERO_VERSION,
&plan,
None,
RepairConfirmation::OperatorAcknowledged,
)
.await;
assert!(
matches!(result, Err(RepairError::PhaseZeroArtifactRefused { .. })),
"{name} Phase 0 artifact must refuse before resume pins, got {result:?}"
);
let _ = fs::remove_dir_all(&root);
}
}
#[test]
fn compute_plan_checksum_up_matches_runner_path() {
use crate::migrate::diff::Classification;
use crate::migrate::projection::BucketKey;
use crate::migrate::segment::{Segment, SegmentKind};
use crate::migrate::sql::OperationSql;
let plan = MigrationPlan {
bucket: BucketKey {
database: "main".to_string(),
app: "".to_string(),
},
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,
}],
}],
};
let from_repair = compute_plan_checksum_up(&plan);
let manual = compute_checksum(["CREATE TABLE users ()"]);
assert_eq!(
from_repair, manual,
"B-10: repair's resume checksum must match the runner's plan checksum"
);
}
#[test]
fn plan_checksum_mismatch_message_names_versions() {
let e = RepairError::PlanChecksumMismatch {
version: "V123".to_string(),
ledger_checksum: "V1:aaaaa".to_string(),
plan_checksum: "V1:bbbbb".to_string(),
};
let msg = format!("{e}");
assert!(msg.contains("V123"));
assert!(msg.contains("V1:aaaaa"));
assert!(msg.contains("V1:bbbbb"));
}
#[test]
fn nothing_to_resume_message_distinguishes_total_states() {
let with_total = RepairError::NothingToResume {
version: "V1".to_string(),
applied: 5,
total: Some(5),
};
let m = format!("{with_total}");
assert!(m.contains("applied_steps_count=5"));
assert!(m.contains("total_steps=5"));
let no_total = RepairError::NothingToResume {
version: "V1".to_string(),
applied: 0,
total: None,
};
let m = format!("{no_total}");
assert!(m.contains("total_steps is NULL"));
}
fn minimal_ledger_row(app_label: &str) -> LedgerRow {
use crate::migrate::ledger::ExecutionMode;
LedgerRow {
version: "V1__test".to_string(),
description: "unit test row".to_string(),
checksum_up: "V1:0000000000000000000000000000000000000000000000000000000000000000"
.to_string(),
checksum_down: None,
execution_mode: ExecutionMode::Transactional,
status: LedgerStatus::Applied,
execution_time_ms: 0,
out_of_order_flag: false,
applied_steps_count: 0,
total_steps: None,
partial_apply_note: None,
run_id: 0,
snapshot_version: "1".to_string(),
app_label: app_label.to_string(),
leaf_identity: None,
}
}
#[test]
fn bucket_app_mismatch_display_names_all_fields() {
let e = RepairError::BucketAppMismatch {
version: "V20260425__add_users".to_string(),
row_app_label: "blog".to_string(),
supplied_app: "store".to_string(),
};
let msg = format!("{e}");
assert!(
msg.contains("V20260425__add_users"),
"Display must include the version; msg={msg}"
);
assert!(
msg.contains("blog"),
"Display must include the row's app_label; msg={msg}"
);
assert!(
msg.contains("store"),
"Display must include the supplied (wrong) app; msg={msg}"
);
}
#[test]
fn ensure_row_matches_bucket_app_ok_when_apps_agree() {
let row = minimal_ledger_row("blog");
let bucket = BucketKey {
database: "main".to_string(),
app: "blog".to_string(),
};
assert!(
ensure_row_matches_bucket_app(&row, &bucket, "V1__test").is_ok(),
"matching apps must return Ok"
);
}
#[test]
fn ensure_row_matches_bucket_app_err_when_apps_differ() {
let row = minimal_ledger_row("blog");
let bucket = BucketKey {
database: "main".to_string(),
app: "store".to_string(),
};
let err = ensure_row_matches_bucket_app(&row, &bucket, "V1__test")
.expect_err("mismatched apps must return Err");
match err {
RepairError::BucketAppMismatch {
version,
row_app_label,
supplied_app,
} => {
assert_eq!(version, "V1__test", "error must carry the version");
assert_eq!(
row_app_label, "blog",
"error must carry the row's app_label"
);
assert_eq!(supplied_app, "store", "error must carry the supplied app");
}
other => panic!("expected BucketAppMismatch, got {other:?}"),
}
}
#[test]
fn repair_error_resume_plan_shape_mismatch_display_names_counts() {
let e = RepairError::ResumePlanShapeMismatch {
version: "V20260526031700__shape".to_string(),
ledger_total_steps: 5,
replay_total_steps: 1,
};
let msg = e.to_string();
assert!(msg.contains("V20260526031700__shape"));
assert!(msg.contains("ledger total_steps=5"));
assert!(msg.contains("expanded replay non-transactional statements=1"));
}
#[test]
fn repair_error_runner_preserves_node_identity_binding_source_chain() {
let err = RepairError::Runner(RunnerError::NodeIdentityBindingFailed {
node_id: 7,
source: DjogiError::Db(crate::error::DbError::other("bind statement failed")),
});
let level1 = std::error::Error::source(&err)
.expect("RepairError::Runner must expose its RunnerError source");
let runner = level1
.downcast_ref::<RunnerError>()
.expect("source must downcast to RunnerError");
match runner {
RunnerError::NodeIdentityBindingFailed { node_id, .. } => {
assert_eq!(*node_id, 7);
}
other => panic!("expected NodeIdentityBindingFailed, got {other:?}"),
}
let level2 = std::error::Error::source(runner)
.expect("NodeIdentityBindingFailed must expose its DjogiError source");
assert!(level2.downcast_ref::<DjogiError>().is_some());
}
#[test]
fn repair_error_runner_preserves_catalog_query_source_chain() {
let err = RepairError::Runner(RunnerError::CatalogQueryFailed {
query_label: "pg_class relpages",
source: DjogiError::Db(crate::error::DbError::other("connection reset by peer")),
});
let level1 = std::error::Error::source(&err)
.expect("RepairError::Runner must expose its RunnerError source");
let runner = level1
.downcast_ref::<RunnerError>()
.expect("source must downcast to RunnerError");
match runner {
RunnerError::CatalogQueryFailed { query_label, .. } => {
assert_eq!(*query_label, "pg_class relpages");
}
other => panic!("expected CatalogQueryFailed, got {other:?}"),
}
let level2 = std::error::Error::source(runner)
.expect("CatalogQueryFailed must expose its DjogiError source");
assert!(level2.downcast_ref::<DjogiError>().is_some());
}
#[test]
fn repair_error_replay_plan_shape_mismatch_display_names_counts() {
let e = RepairError::ReplayPlanShapeMismatch {
version: "V20260526031700__shape".to_string(),
expected_step_count: 5,
actual_step_count: 1,
};
let msg = e.to_string();
assert!(msg.contains("V20260526031700__shape"));
assert!(msg.contains("completed 1 step(s)"));
assert!(msg.contains("expected 5"));
}
#[test]
fn repair_leaf_identity_mismatch_display() {
let err = RepairError::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("[D623]"));
assert!(msg.contains("repair refused"));
assert!(msg.contains("partition leaf identity mismatch"));
assert!(msg.contains("001_create_users"));
}
#[test]
fn repair_error_missing_resume_identity_display() {
let err = RepairError::MissingResumeIdentity {
version: "V20260601000004__g_gate_resume".to_string(),
};
let msg = format!("{err}");
assert!(
msg.contains("repair resume-partial refused"),
"Display must contain refusal message; msg={msg}"
);
assert!(
msg.contains("V20260601000004__g_gate_resume"),
"Display must include the version; msg={msg}"
);
assert!(
msg.contains("binding-capable"),
"Display must mention binding-capable identity requirement; msg={msg}"
);
assert!(
msg.contains("SQL replay"),
"Display must mention SQL replay context; msg={msg}"
);
}
#[test]
fn repair_gate_refuses_none_identity_for_non_phase_zero() {
let version = "V20260601000005__g_gate_resume_unit";
let runner_identity: Option<RunnerIdentity> = None;
let result: Result<(), RepairError> = if version != PHASE_ZERO_VERSION {
match runner_identity {
Some(id) => {
if !id.requires_binding() {
Err(RepairError::MissingResumeIdentity {
version: version.to_string(),
})
} else {
Ok(()) }
}
None => Err(RepairError::MissingResumeIdentity {
version: version.to_string(),
}),
}
} else {
Ok(()) };
assert!(
matches!(result, Err(RepairError::MissingResumeIdentity { .. })),
"Non-P0 resume with no identity should produce MissingResumeIdentity"
);
}
#[test]
fn repair_gate_refuses_identity_free_for_non_phase_zero() {
let version = "V20260601000006__g_gate_resume_idfree";
let runner_identity = Some(RunnerIdentity::IdentityFree);
let result: Result<(), RepairError> = if version != PHASE_ZERO_VERSION {
match runner_identity {
Some(id) => {
if !id.requires_binding() {
Err(RepairError::MissingResumeIdentity {
version: version.to_string(),
})
} else {
Ok(()) }
}
None => Err(RepairError::MissingResumeIdentity {
version: version.to_string(),
}),
}
} else {
Ok(()) };
assert!(
matches!(result, Err(RepairError::MissingResumeIdentity { .. })),
"Non-P0 resume with IdentityFree should produce MissingResumeIdentity"
);
}
#[test]
fn repair_gate_allows_phase_zero_without_identity() {
let version = PHASE_ZERO_VERSION;
let runner_identity: Option<RunnerIdentity> = None;
let result: Result<(), RepairError> = if version != PHASE_ZERO_VERSION {
match runner_identity {
Some(id) => {
if !id.requires_binding() {
Err(RepairError::MissingResumeIdentity {
version: version.to_string(),
})
} else {
Ok(()) }
}
None => Err(RepairError::MissingResumeIdentity {
version: version.to_string(),
}),
}
} else {
Ok(()) };
assert!(
result.is_ok(),
"Phase 0 resume with no identity should NOT refuse (carve-out)"
);
}
#[test]
fn repair_gate_allows_non_phase_zero_with_binding_identity() {
let version = "V20260601000007__g_gate_resume_binding";
let runner_identity = Some(RunnerIdentity::SingleNodeDev);
let result: Result<(), RepairError> = if version != PHASE_ZERO_VERSION {
match runner_identity {
Some(id) => {
if !id.requires_binding() {
Err(RepairError::MissingResumeIdentity {
version: version.to_string(),
})
} else {
Ok(()) }
}
None => Err(RepairError::MissingResumeIdentity {
version: version.to_string(),
}),
}
} else {
Ok(()) };
assert!(
result.is_ok(),
"Non-P0 resume with binding-capable identity should pass the gate"
);
}
#[djogi::deliberately_bypass_convention_with_raw_sql]
#[djogi_test]
async fn load_row_two_phase_lookup(mut ctx: DjogiContext) {
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, 'test', 'V1:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855', \
'applied', 1, '1', 'app_a')",
&[&"V1"],
)
.await
.expect("seed (V1, app_a) row");
let result = load_row(&mut ctx, "V1", "app_a").await;
let row = result.expect("Phase 1: exact composite hit must return Ok");
assert_eq!(
row.app_label, "app_a",
"Phase 1: returned row must belong to app_a"
);
ctx.raw_execute(
"INSERT INTO djogi_schema_migrations \
(version, description, checksum_up, status, run_id, \
snapshot_version, app_label) \
VALUES ($1, 'test', 'V1:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855', \
'applied', 2, '1', 'app_b')",
&[&"V1"],
)
.await
.expect("seed (V1, app_b) row");
let a_row = load_row(&mut ctx, "V1", "app_a")
.await
.expect("multi-row: (V1, app_a) must still resolve to app_a's row");
assert_eq!(
a_row.app_label, "app_a",
"multi-row: app_a key returns app_a row"
);
let b_row = load_row(&mut ctx, "V1", "app_b")
.await
.expect("multi-row: (V1, app_b) must resolve to app_b's row");
assert_eq!(
b_row.app_label, "app_b",
"multi-row: app_b key returns app_b row"
);
ctx.raw_execute(
"INSERT INTO djogi_schema_migrations \
(version, description, checksum_up, status, run_id, \
snapshot_version, app_label) \
VALUES ($1, 'test', 'V2:e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855', \
'applied', 3, '1', 'app_c')",
&[&"V2"],
)
.await
.expect("seed (V2, app_c) row for fallback test");
let fallback_result = load_row(&mut ctx, "V2", "app_d").await;
let fallback_row = fallback_result.expect(
"Phase 2 fallback: version exists under a different app — must return Ok(row), not VersionNotFound",
);
assert_eq!(
fallback_row.app_label, "app_c",
"Phase 2 fallback: returned row must belong to app_c (the actual owner)",
);
let wrong_bucket = BucketKey {
database: "main".to_string(),
app: "app_d".to_string(),
};
let mismatch = ensure_row_matches_bucket_app(&fallback_row, &wrong_bucket, "V2")
.expect_err("wrong-app fallback row must produce BucketAppMismatch");
assert!(
matches!(mismatch, RepairError::BucketAppMismatch { .. }),
"expected BucketAppMismatch from ensure_row_matches_bucket_app, got {mismatch:?}",
);
let missing_result = load_row(&mut ctx, "MISSING", "app_a").await;
assert!(
matches!(missing_result, Err(RepairError::VersionNotFound { .. })),
"absent version must return VersionNotFound, got {missing_result:?}",
);
}
}