use std::path::PathBuf;
use crate::__bypass::guarded_batch_execute;
use crate::context::{DjogiContext, PinnedCtx};
use crate::error::DjogiError;
use super::guard::WorkspaceGuard;
use super::ledger::{
self, ChecksumFormatErrorKind, LedgerRow, LedgerStatus, compute_checksum,
load_full_row_by_version, validate_checksum_format,
};
use super::projection::BucketKey;
use super::runner::{
PartitionExpansionMode, RunnerError, 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,
},
}
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}`",
),
}
}
}
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),
_ => 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::AdvisoryLockQueryFailed {
app_label: bucket.app.clone(),
source: DjogiError::Db(crate::error::DbError::other(format!(
"unexpected acquire error: {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),
}
}
pub async fn repair_checksum_drift(
ctx: &mut DjogiContext,
_guard: &WorkspaceGuard,
bucket: &BucketKey,
version: &str,
new_checksum_up: &str,
new_checksum_down: Option<&str>,
confirmation: RepairConfirmation,
) -> Result<RepairReport, RepairError> {
if confirmation != RepairConfirmation::OperatorAcknowledged {
return Err(RepairError::InsufficientConfirmation);
}
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).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",
&[&version, &new_checksum_up, &new_down_owned],
)
.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,
}
pub async fn repair_partial_apply(
ctx: &mut DjogiContext,
_guard: &WorkspaceGuard,
bucket: &BucketKey,
version: &str,
resolution: PartialApplyResolution,
note: &str,
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_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).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",
&[&version, &status_str, &target_steps, ¬e],
)
.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,
version: &str,
plan: &MigrationPlan,
confirmation: RepairConfirmation,
) -> Result<RepairReport, RepairError> {
if confirmation != RepairConfirmation::OperatorAcknowledged {
return Err(RepairError::InsufficientConfirmation);
}
let lock_key = advisory_lock_key(&plan.bucket);
let mut pinned = ctx
.pin_for_migration()
.await
.map_err(|e| RepairError::PinnedSessionCheckoutFailed { source: e })?;
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).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).await?;
if let Some(ref stored_identity) = row.leaf_identity {
let pre_cache =
compute_leaf_identity_cache(ctx, plan)
.await
.map_err(|e| RepairError::LedgerIo {
source: DjogiError::Db(crate::error::DbError::other(format!(
"leaf-identity pre-check failed: {e}"
))),
})?;
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::LedgerIo {
source: DjogiError::Db(crate::error::DbError::other(format!(
"materialize_execution_plan failed: {other}"
))),
},
})?;
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",
&[&row_version, &applied],
)
.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,
) -> Result<i64, RepairError> {
let row = ctx
.query_one(
"SELECT id FROM djogi_schema_migrations WHERE version = $1",
&[&version],
)
.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)
.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) -> Result<LedgerRow, RepairError> {
load_full_row_by_version(ctx, version)
.await
.map_err(|e| RepairError::LedgerIo { source: e })?
.ok_or_else(|| 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::*;
#[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 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_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"));
}
}