feat(manifests-consistency): add safe repair mode
This commit is contained in:
@@ -45,11 +45,25 @@ use serde::{Deserialize, Serialize};
|
||||
/// of the entry to probe instead of the currently-active backend.
|
||||
/// `None` falls through to the live backend (today's behaviour).
|
||||
/// - Others — ignored.
|
||||
///
|
||||
/// Semantics of `repair` (added 2026-10-17 for the refcount fix):
|
||||
/// - `blobs_consistency` / `manifests_consistency` — when `true`,
|
||||
/// after each `refcount_mismatch` / `manifest_refcount_mismatch`
|
||||
/// finding is recorded, apply the corrective UPDATE that sets the
|
||||
/// stored counter to the auditor's computed `actual_ref_count`.
|
||||
/// Content-safe: the row itself is fine, only the counter is
|
||||
/// wrong. Race-safe: each UPDATE recomputes the auditor formula
|
||||
/// in the same statement, so a concurrent write can't leave a
|
||||
/// stale value. Default `false` preserves discovery-only
|
||||
/// behaviour. Also propagates through `consistency_batch` to
|
||||
/// both tenants — one `?repair=true` call fixes both counters.
|
||||
/// - Others — ignored.
|
||||
#[derive(Debug, Clone, Default)]
|
||||
pub struct JobRunArgs {
|
||||
pub force: bool,
|
||||
pub deep: bool,
|
||||
pub storage: Option<String>,
|
||||
pub repair: bool,
|
||||
}
|
||||
|
||||
/// Uniform outcome the supervisor logs and stores for every job dispatch.
|
||||
|
||||
@@ -347,6 +347,10 @@ impl RecoverableJobHandler for BlobsConsistencyCheck {
|
||||
// stats.finding_count — actual persistence happens in
|
||||
// `record_finding` on each emission).
|
||||
let mut finding_count = 0u64;
|
||||
// Only touched when `args.repair == true`. Symmetric with
|
||||
// `manifests_consistency`; reported in completion log +
|
||||
// `extra_stats` so operators see "found N, fixed M" in one line.
|
||||
let mut repaired_count = 0u64;
|
||||
|
||||
// Deep mode is a per-run flag with two consumers:
|
||||
// 1. This handler — decides whether to re-hash bytes.
|
||||
@@ -399,6 +403,41 @@ impl RecoverableJobHandler for BlobsConsistencyCheck {
|
||||
);
|
||||
}
|
||||
|
||||
// Repair mode: same shape as `deep` above so the admin run-
|
||||
// detail view can display `params.repair = "true"` alongside
|
||||
// `params.deep`. Fresh persists what the trigger asked for;
|
||||
// Resume reads back so a paused repair scan stays a repair
|
||||
// scan (a mid-scan crash mustn't silently downgrade to
|
||||
// discovery-only for the remaining rows).
|
||||
let repair = if is_fresh {
|
||||
let v = if args.repair { "true" } else { "false" };
|
||||
if let Err(e) = store.set_string_param("repair", v).await {
|
||||
return RunOutcome::Failed {
|
||||
message: format!("failed to persist repair flag to params: {e}"),
|
||||
};
|
||||
}
|
||||
args.repair
|
||||
} else {
|
||||
match store.get_string_param("repair").await {
|
||||
Ok(Some(v)) => v == "true",
|
||||
Ok(None) => false,
|
||||
Err(e) => {
|
||||
return RunOutcome::Failed {
|
||||
message: format!("read `repair` from params: {e}"),
|
||||
};
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
if repair {
|
||||
tracing::info!(
|
||||
target: "oxicloud::consistency",
|
||||
event = "blobs_consistency.repair_mode_active",
|
||||
run_id = %store.run_id(),
|
||||
"repair mode: refcount_mismatch findings will trigger corrective UPDATE"
|
||||
);
|
||||
}
|
||||
|
||||
loop {
|
||||
// Cooperative cancel poll between batches.
|
||||
match store.status().await {
|
||||
@@ -448,11 +487,17 @@ impl RecoverableJobHandler for BlobsConsistencyCheck {
|
||||
event = "blobs_consistency.completed",
|
||||
run_id = %store.run_id(),
|
||||
finding_count = finding_count,
|
||||
repaired_count = repaired_count,
|
||||
repair_requested = repair,
|
||||
deep = deep,
|
||||
"blobs_consistency completed with {} finding(s)",
|
||||
finding_count
|
||||
"blobs_consistency completed with {} finding(s), {} repaired",
|
||||
finding_count,
|
||||
repaired_count
|
||||
);
|
||||
return RunOutcome::completed();
|
||||
return RunOutcome::completed_with(serde_json::json!({
|
||||
"repair_requested": repair,
|
||||
"repaired_count": repaired_count,
|
||||
}));
|
||||
}
|
||||
|
||||
let grace_cutoff = Utc::now() - CREATE_GRACE;
|
||||
@@ -480,6 +525,67 @@ impl RecoverableJobHandler for BlobsConsistencyCheck {
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
|
||||
// Repair pass — content-safe corrective UPDATE. Sets
|
||||
// `stored` to the value the auditor's two-term formula
|
||||
// would compute at UPDATE time (subquery mirrors
|
||||
// `chunk_page_sql`'s `actual_ref_count`), so a
|
||||
// concurrent write between our page fetch and this
|
||||
// UPDATE can't leave a stale value — the subquery
|
||||
// re-reads inside the same statement. The
|
||||
// `<> (subquery)` guard makes the UPDATE a no-op if
|
||||
// the drift has healed, making this idempotent under
|
||||
// retry.
|
||||
if repair {
|
||||
let expected = "( \
|
||||
(SELECT COUNT(*) FROM storage.files f \
|
||||
WHERE f.blob_hash = b.hash \
|
||||
AND NOT EXISTS ( \
|
||||
SELECT 1 FROM storage.chunk_manifests m \
|
||||
WHERE m.file_hash = f.blob_hash \
|
||||
)) \
|
||||
+ (SELECT COUNT(*) FROM storage.chunk_manifests m \
|
||||
WHERE b.hash = ANY(m.chunk_hashes)) \
|
||||
)";
|
||||
let update_sql = format!(
|
||||
"UPDATE storage.blobs b \
|
||||
SET ref_count = {expected} \
|
||||
WHERE b.hash = $1 \
|
||||
AND b.ref_count <> {expected}",
|
||||
);
|
||||
match sqlx::query(&update_sql)
|
||||
.bind(&row.hash)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
{
|
||||
Ok(res) if res.rows_affected() > 0 => {
|
||||
repaired_count += 1;
|
||||
tracing::info!(
|
||||
target: "audit",
|
||||
event = "blobs_consistency.repaired",
|
||||
run_id = %store.run_id(),
|
||||
hash = %row.hash,
|
||||
stored_was = row.ref_count,
|
||||
actual = row.actual_ref_count,
|
||||
"🩹 blob ref_count repaired"
|
||||
);
|
||||
}
|
||||
Ok(_) => {
|
||||
// No row touched — concurrent repair or
|
||||
// self-healing drift. Silent no-op.
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
target: "oxicloud::consistency",
|
||||
event = "blobs_consistency.repair_failed",
|
||||
run_id = %store.run_id(),
|
||||
hash = %row.hash,
|
||||
error = %e,
|
||||
"blob ref_count repair UPDATE failed — finding stays"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Skip physical probes for rows within the write
|
||||
@@ -628,11 +734,17 @@ impl RecoverableJobHandler for BlobsConsistencyCheck {
|
||||
event = "blobs_consistency.completed",
|
||||
run_id = %store.run_id(),
|
||||
finding_count = finding_count,
|
||||
repaired_count = repaired_count,
|
||||
repair_requested = repair,
|
||||
deep = deep,
|
||||
"blobs_consistency completed with {} finding(s)",
|
||||
finding_count
|
||||
"blobs_consistency completed with {} finding(s), {} repaired",
|
||||
finding_count,
|
||||
repaired_count
|
||||
);
|
||||
return RunOutcome::completed();
|
||||
return RunOutcome::completed_with(serde_json::json!({
|
||||
"repair_requested": repair,
|
||||
"repaired_count": repaired_count,
|
||||
}));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -178,6 +178,7 @@ impl JobHandler for ConsistencyBatch {
|
||||
"per_check": per_check,
|
||||
"deep": args.deep,
|
||||
"force": args.force,
|
||||
"repair": args.repair,
|
||||
"ok": ok_count,
|
||||
"err": err_count,
|
||||
}),
|
||||
|
||||
@@ -161,9 +161,11 @@ impl RecoverableJobHandler for ManifestsConsistencyCheck {
|
||||
async fn run_resumable(
|
||||
&self,
|
||||
store: &dyn JobStore,
|
||||
_args: &JobRunArgs,
|
||||
args: &JobRunArgs,
|
||||
resume_cursor: Option<Vec<u8>>,
|
||||
) -> RunOutcome {
|
||||
let is_fresh = resume_cursor.is_none();
|
||||
|
||||
// Cursor: the last `file_hash` as UTF-8. Same convention as
|
||||
// `blobs_consistency`, which also pages a hash-keyed table.
|
||||
let mut cursor: Option<String> = match resume_cursor {
|
||||
@@ -179,7 +181,48 @@ impl RecoverableJobHandler for ManifestsConsistencyCheck {
|
||||
},
|
||||
};
|
||||
|
||||
// Persist the repair flag into `params.repair` so the admin
|
||||
// run-detail view can display whether the run was a discovery
|
||||
// scan or an active repair. Fresh takes it from args; Resume
|
||||
// reads back so a paused repair scan stays a repair scan (a
|
||||
// mid-scan crash mustn't silently downgrade the remaining
|
||||
// rows to discovery-only). Same shape as
|
||||
// `blobs_consistency_service.rs`'s `deep` handling — see the
|
||||
// reasoning documented there.
|
||||
let repair = if is_fresh {
|
||||
let v = if args.repair { "true" } else { "false" };
|
||||
if let Err(e) = store.set_string_param("repair", v).await {
|
||||
return RunOutcome::Failed {
|
||||
message: format!("failed to persist repair flag to params: {e}"),
|
||||
};
|
||||
}
|
||||
args.repair
|
||||
} else {
|
||||
match store.get_string_param("repair").await {
|
||||
Ok(Some(v)) => v == "true",
|
||||
Ok(None) => false,
|
||||
Err(e) => {
|
||||
return RunOutcome::Failed {
|
||||
message: format!("read `repair` from params: {e}"),
|
||||
};
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
if repair {
|
||||
tracing::info!(
|
||||
target: "oxicloud::consistency",
|
||||
event = "manifests_consistency.repair_mode_active",
|
||||
run_id = %store.run_id(),
|
||||
"repair mode: manifest_refcount_mismatch findings will trigger corrective UPDATE"
|
||||
);
|
||||
}
|
||||
|
||||
let mut finding_count = 0u64;
|
||||
// Only relevant when `repair == true`. Reported inline in
|
||||
// the completion log + the `extra_stats` payload so operators
|
||||
// can see "we found N and fixed M" in one line.
|
||||
let mut repaired_count = 0u64;
|
||||
|
||||
loop {
|
||||
// Cooperative cancel poll between batches.
|
||||
@@ -227,10 +270,16 @@ impl RecoverableJobHandler for ManifestsConsistencyCheck {
|
||||
event = "manifests_consistency.completed",
|
||||
run_id = %store.run_id(),
|
||||
finding_count = finding_count,
|
||||
"manifests_consistency completed with {} finding(s)",
|
||||
finding_count
|
||||
repaired_count = repaired_count,
|
||||
repair_requested = repair,
|
||||
"manifests_consistency completed with {} finding(s), {} repaired",
|
||||
finding_count,
|
||||
repaired_count
|
||||
);
|
||||
return RunOutcome::completed();
|
||||
return RunOutcome::completed_with(serde_json::json!({
|
||||
"repair_requested": repair,
|
||||
"repaired_count": repaired_count,
|
||||
}));
|
||||
}
|
||||
|
||||
for row in &rows {
|
||||
@@ -258,6 +307,64 @@ impl RecoverableJobHandler for ManifestsConsistencyCheck {
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
|
||||
// Repair pass — content-safe corrective UPDATE. The
|
||||
// stored counter is set to what the auditor formula
|
||||
// would compute at UPDATE time (subquery matches
|
||||
// `manifest_page_sql`'s `actual_ref_count` predicate),
|
||||
// so a concurrent file insert/delete between our page
|
||||
// fetch and this UPDATE can't leave a stale value —
|
||||
// the subquery re-reads inside the same statement.
|
||||
// The `<> (subquery)` guard makes the UPDATE a no-op
|
||||
// if the value is already correct, so this is
|
||||
// idempotent under retry.
|
||||
if repair {
|
||||
match sqlx::query(
|
||||
"UPDATE storage.chunk_manifests m \
|
||||
SET ref_count = ( \
|
||||
SELECT COUNT(*) FROM storage.files \
|
||||
WHERE blob_hash = m.file_hash \
|
||||
) \
|
||||
WHERE m.file_hash = $1 \
|
||||
AND m.ref_count <> ( \
|
||||
SELECT COUNT(*) FROM storage.files \
|
||||
WHERE blob_hash = m.file_hash \
|
||||
)",
|
||||
)
|
||||
.bind(&row.file_hash)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
{
|
||||
Ok(res) if res.rows_affected() > 0 => {
|
||||
repaired_count += 1;
|
||||
tracing::info!(
|
||||
target: "audit",
|
||||
event = "manifests_consistency.repaired",
|
||||
run_id = %store.run_id(),
|
||||
file_hash = %row.file_hash,
|
||||
stored_was = row.ref_count,
|
||||
actual = row.actual_ref_count,
|
||||
"🩹 manifest ref_count repaired"
|
||||
);
|
||||
}
|
||||
Ok(_) => {
|
||||
// Row not touched — either another concurrent
|
||||
// repair fixed it first, or the drift healed
|
||||
// itself between page fetch and UPDATE.
|
||||
// Silent no-op.
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
target: "oxicloud::consistency",
|
||||
event = "manifests_consistency.repair_failed",
|
||||
run_id = %store.run_id(),
|
||||
file_hash = %row.file_hash,
|
||||
error = %e,
|
||||
"manifest ref_count repair UPDATE failed — finding stays"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Advance cursor + checkpoint.
|
||||
@@ -279,10 +386,16 @@ impl RecoverableJobHandler for ManifestsConsistencyCheck {
|
||||
event = "manifests_consistency.completed",
|
||||
run_id = %store.run_id(),
|
||||
finding_count = finding_count,
|
||||
"manifests_consistency completed with {} finding(s)",
|
||||
finding_count
|
||||
repaired_count = repaired_count,
|
||||
repair_requested = repair,
|
||||
"manifests_consistency completed with {} finding(s), {} repaired",
|
||||
finding_count,
|
||||
repaired_count
|
||||
);
|
||||
return RunOutcome::completed();
|
||||
return RunOutcome::completed_with(serde_json::json!({
|
||||
"repair_requested": repair,
|
||||
"repaired_count": repaired_count,
|
||||
}));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2575,6 +2575,12 @@ pub async fn list_jobs(State(state): State<Arc<AppState>>) -> impl IntoResponse
|
||||
/// `deep=true` opts into slow variants — `consistency_batch` fans it
|
||||
/// out to sub-jobs; `storage_consistency` (when implemented) will
|
||||
/// re-BLAKE3 each blob for bitrot detection. See `JobRunArgs.deep`.
|
||||
///
|
||||
/// `repair=true` opts into corrective action on the refcount
|
||||
/// consistency tenants (`blobs_consistency`, `manifests_consistency`,
|
||||
/// and `consistency_batch` which fans out to both). Default `false`
|
||||
/// preserves discovery-only. See `JobRunArgs.repair` for the
|
||||
/// content-safety and race-safety guarantees.
|
||||
#[derive(serde::Deserialize)]
|
||||
pub struct TriggerJobQuery {
|
||||
#[serde(default)]
|
||||
@@ -2591,6 +2597,8 @@ pub struct TriggerJobQuery {
|
||||
/// `AppConfig.storage_entries`.
|
||||
#[serde(default)]
|
||||
pub storage: Option<String>,
|
||||
#[serde(default)]
|
||||
pub repair: bool,
|
||||
}
|
||||
|
||||
/// `POST /api/admin/jobs/{name}/trigger` — dispatch one run off-schedule.
|
||||
@@ -2629,15 +2637,18 @@ pub async fn trigger_job(
|
||||
job = %name,
|
||||
force = query.force,
|
||||
deep = query.deep,
|
||||
"👮🏻♂️ Admin triggered job {} (force={}, deep={})",
|
||||
repair = query.repair,
|
||||
"👮🏻♂️ Admin triggered job {} (force={}, deep={}, repair={})",
|
||||
name,
|
||||
query.force,
|
||||
query.deep,
|
||||
query.repair,
|
||||
);
|
||||
let args = JobRunArgs {
|
||||
force: query.force,
|
||||
deep: query.deep,
|
||||
storage: query.storage.clone(),
|
||||
repair: query.repair,
|
||||
};
|
||||
|
||||
// Jobs that can run for hours (backend_migration, future
|
||||
|
||||
Reference in New Issue
Block a user