fix(blob_consistency): apply same repair logic as manifest_consistency

This commit is contained in:
Edouard Vanbelle
2026-09-02 22:13:55 +02:00
parent 8a63663209
commit 2a629c4e8b
@@ -73,6 +73,18 @@ pub struct BlobsConsistencyCheck {
/// fixed statement — same reasoning as `DedupService::manifest_reap_sql`. /// fixed statement — same reasoning as `DedupService::manifest_reap_sql`.
/// See `docs/plan/derived-blobs.md`. /// See `docs/plan/derived-blobs.md`.
chunk_page_sql: String, chunk_page_sql: String,
/// Per-row repair UPDATE, built from the SAME registry as
/// `chunk_page_sql` so detection and repair use identical formulas
/// by construction. A future `RefLevel::Chunk` ref source added to
/// the registry flows into both without a code change here.
///
/// Previously the repair query was inlined with a hardcoded
/// 2-term formula (accidentally matching detection today). Would
/// silently diverge the moment a new chunk-level ref source
/// landed — same class of latent bug the sibling
/// `manifests_consistency` service hit 2026-09-02. Preemptively
/// pulled from the registry here to keep the pair symmetric.
chunk_repair_sql: String,
} }
/// The chunk-level page query, with `actual_ref_count` summed from the /// The chunk-level page query, with `actual_ref_count` summed from the
@@ -126,11 +138,33 @@ fn chunk_page_sql(registry: &BlobReferenceRegistry) -> String {
) )
} }
/// Per-row corrective UPDATE for `storage.blobs.ref_count`, targeting
/// one blob by `hash`. Uses the SAME registry-derived expression as
/// [`chunk_page_sql`] so detection and repair agree on "actual" by
/// construction. A future `RefLevel::Chunk` ref source added to the
/// registry flows into both queries with no code change here.
///
/// The `<> (subquery)` guard makes the UPDATE a no-op when the value
/// is already correct — idempotent under concurrent-repair races and
/// under retry. The subquery re-reads inside the same statement, so a
/// concurrent write between page fetch and this UPDATE can't leave a
/// stale value.
fn chunk_repair_sql(registry: &BlobReferenceRegistry) -> String {
let expected = registry.ref_count_expr(RefLevel::Chunk, "b.hash");
format!(
"UPDATE storage.blobs b
SET ref_count = ({expected})::bigint
WHERE b.hash = $1
AND b.ref_count <> ({expected})::bigint"
)
}
impl BlobsConsistencyCheck { impl BlobsConsistencyCheck {
pub fn new(pool: Arc<PgPool>, reference_registry: Arc<BlobReferenceRegistry>) -> Self { pub fn new(pool: Arc<PgPool>, reference_registry: Arc<BlobReferenceRegistry>) -> Self {
Self { Self {
pool, pool,
chunk_page_sql: chunk_page_sql(&reference_registry), chunk_page_sql: chunk_page_sql(&reference_registry),
chunk_repair_sql: chunk_repair_sql(&reference_registry),
} }
} }
@@ -356,54 +390,51 @@ impl RecoverableJobHandler for BlobsConsistencyCheck {
// longer looks at. The refcount comparison reads one // longer looks at. The refcount comparison reads one
// consistent DB snapshot, so there is nothing to wait for. // consistent DB snapshot, so there is nothing to wait for.
for row in &rows { for row in &rows {
if row.ref_count as i64 != row.actual_ref_count { if row.ref_count as i64 == row.actual_ref_count {
continue;
}
finding_count += 1; finding_count += 1;
let affected = affected_files(self.pool.as_ref(), &row.hash).await; let affected = affected_files(self.pool.as_ref(), &row.hash).await;
record_or_log( let detail = serde_json::json!({
store,
BLOBS_CONSISTENCY_JOB_NAME,
"refcount_mismatch",
"inconsistent",
None, // hash isn't a UUID; resource identifier lives in detail
serde_json::json!({
"hash": row.hash, "hash": row.hash,
"stored": row.ref_count, "stored": row.ref_count,
"actual": row.actual_ref_count, "actual": row.actual_ref_count,
"delta": row.actual_ref_count - row.ref_count as i64, "delta": row.actual_ref_count - row.ref_count as i64,
"size": row.size, "size": row.size,
"affected_files": affected, "affected_files": affected,
}), });
)
.await;
// Repair pass — content-safe corrective UPDATE. Sets // Repair pass — content-safe corrective UPDATE. Sets
// `stored` to the value the auditor's two-term formula // `stored` to the value the auditor formula would
// would compute at UPDATE time (subquery mirrors // compute at UPDATE time (subquery matches
// `chunk_page_sql`'s `actual_ref_count`), so a // `chunk_page_sql`'s `actual_ref_count`), so a
// concurrent write between our page fetch and this // concurrent write between our page fetch and this
// UPDATE can't leave a stale value — the subquery // UPDATE can't leave a stale value — the subquery
// re-reads inside the same statement. The // re-reads inside the same statement. The `<>`
// `<> (subquery)` guard makes the UPDATE a no-op if // guard makes the UPDATE a no-op if the value is
// the drift has healed, making this idempotent under // already correct, so this is idempotent under retry.
// retry. //
if repair { // `self.chunk_repair_sql` is built once at construction
let expected = "( \ // from the same `BlobReferenceRegistry` as the page
(SELECT COUNT(*) FROM storage.files f \ // query — detection and repair use identical formulas
WHERE f.blob_hash = b.hash \ // by construction. See `chunk_repair_sql`.
AND NOT EXISTS ( \ //
SELECT 1 FROM storage.chunk_manifests m \ // Attempt repair FIRST, then record the finding with
WHERE m.file_hash = f.blob_hash \ // severity/kind reflecting the final state:
)) \ // * repair succeeded → severity "info", kind "refcount_repaired"
+ (SELECT COUNT(*) FROM storage.chunk_manifests m \ // * repair no-op → severity "info", kind "refcount_resolved"
WHERE b.hash = ANY(m.chunk_hashes)) \ // * repair failed → severity "inconsistent", kind "refcount_mismatch"
)"; // * no repair requested → severity "inconsistent", kind "refcount_mismatch"
let update_sql = format!( //
"UPDATE storage.blobs b \ // Parallels the WARN-then-INFO sequence in logs: an
SET ref_count = {expected} \ // unresolved drift raises attention ("inconsistent"),
WHERE b.hash = $1 \ // a repaired one records the fix at info level without
AND b.ref_count <> {expected}", // inflating the "needs action" tally the outcome UI
); // shows. The detail JSON still carries `stored/actual/
match sqlx::query(&update_sql) // delta/affected_files` so the audit trail is complete
// either way.
let (kind, severity) = if repair {
match sqlx::query(&self.chunk_repair_sql)
.bind(&row.hash) .bind(&row.hash)
.execute(self.pool.as_ref()) .execute(self.pool.as_ref())
.await .await
@@ -419,10 +450,14 @@ impl RecoverableJobHandler for BlobsConsistencyCheck {
actual = row.actual_ref_count, actual = row.actual_ref_count,
"🩹 blob ref_count repaired" "🩹 blob ref_count repaired"
); );
("refcount_repaired", "info")
} }
Ok(_) => { Ok(_) => {
// No row touched — concurrent repair or // Row not touched — either another
// self-healing drift. Silent no-op. // concurrent repair fixed it first, or
// drift healed between page fetch and
// UPDATE. Current state correct — info.
("refcount_resolved", "info")
} }
Err(e) => { Err(e) => {
tracing::warn!( tracing::warn!(
@@ -433,10 +468,22 @@ impl RecoverableJobHandler for BlobsConsistencyCheck {
error = %e, error = %e,
"blob ref_count repair UPDATE failed — finding stays" "blob ref_count repair UPDATE failed — finding stays"
); );
("refcount_mismatch", "inconsistent")
} }
} }
} } else {
} ("refcount_mismatch", "inconsistent")
};
record_or_log(
store,
BLOBS_CONSISTENCY_JOB_NAME,
kind,
severity,
None, // hash isn't a UUID; resource identifier lives in detail
detail,
)
.await;
} }
// Advance cursor + checkpoint. // Advance cursor + checkpoint.