From 09a5c4317d73d366e8d2b72a1a6ed181847c2205 Mon Sep 17 00:00:00 2001 From: coder 0 Date: Tue, 11 Aug 2026 21:55:31 -0400 Subject: [PATCH] fix(media): make layout cleanup idempotent Co-authored-by: coder 0 Signed-off-by: coder 0 --- .../src/bin/buzz-media-layout-backfill.rs | 111 +++++++++-- .../bin/buzz-media-layout-delete-legacy.rs | 96 ++++++++- .../src/bin/media_layout_common/mod.rs | 187 ++++++++++++++++-- 3 files changed, 357 insertions(+), 37 deletions(-) diff --git a/crates/buzz-media/src/bin/buzz-media-layout-backfill.rs b/crates/buzz-media/src/bin/buzz-media-layout-backfill.rs index 6c915c1c3..d63a2940b 100644 --- a/crates/buzz-media/src/bin/buzz-media-layout-backfill.rs +++ b/crates/buzz-media/src/bin/buzz-media-layout-backfill.rs @@ -5,7 +5,28 @@ use buzz_media::migration::{objects_for_sidecar, parse_sidecar_key, RequestPacer use clap::Parser; mod media_layout_common; -use media_layout_common::{verify_destination, CommonArgs}; +use media_layout_common::{verify_destination, CommonArgs, DestinationVerification}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum BackfillDecision { + Skip, + Copy, +} + +fn decide_backfill( + verification: DestinationVerification, + sharded_key: &str, +) -> Result { + match verification { + DestinationVerification::LegacySourceAbsent | DestinationVerification::Verified => { + Ok(BackfillDecision::Skip) + } + DestinationVerification::DestinationMissing => Ok(BackfillDecision::Copy), + DestinationVerification::ByteMismatch => { + anyhow::bail!("destination verification failed: {sharded_key}") + } + } +} #[derive(Debug, Parser)] #[command(name = "buzz-media-layout-backfill")] @@ -57,26 +78,36 @@ async fn main() -> Result<()> { let meta = serde_json::from_slice(&bytes).context("parse media sidecar")?; for object in objects_for_sidecar(community, sha, &meta)? { processed += 1; - pacer.wait().await; - if verify_destination(&storage, &mut pacer, &object).await? { - skipped += 1; - continue; + let verification = verify_destination(&storage, &mut pacer, &object).await?; + if matches!(verification, DestinationVerification::LegacySourceAbsent) { + tracing::info!(source = %object.legacy, destination = %object.sharded, "legacy source absent; skipping"); } - pacer.wait().await; - let source_exists = storage.head(&object.legacy).await?; - if !source_exists { - anyhow::bail!( - "legacy source missing: {} (checkpoint: {sidecar_key})", - object.legacy - ); + match decide_backfill(verification, &object.sharded)? { + BackfillDecision::Skip => { + skipped += 1; + continue; + } + BackfillDecision::Copy => {} } if args.dry_run { tracing::info!(source = %object.legacy, destination = %object.sharded, "would copy"); } else { pacer.wait().await; storage.copy(&object.legacy, &object.sharded).await?; - if !verify_destination(&storage, &mut pacer, &object).await? { - anyhow::bail!("destination verification failed: {}", object.sharded); + match verify_destination(&storage, &mut pacer, &object).await? { + DestinationVerification::Verified => {} + DestinationVerification::LegacySourceAbsent => { + anyhow::bail!( + "legacy source disappeared before verification: {}", + object.legacy + ); + } + DestinationVerification::DestinationMissing => { + anyhow::bail!("destination verification failed: {}", object.sharded); + } + DestinationVerification::ByteMismatch => { + anyhow::bail!("destination verification failed: {}", object.sharded); + } } copied += 1; } @@ -94,3 +125,55 @@ async fn main() -> Result<()> { tracing::info!(processed, copied, skipped, dry_run = args.dry_run, checkpoint = ?checkpoint, "backfill complete"); Ok(()) } + +#[cfg(test)] +mod tests { + use super::{decide_backfill, BackfillDecision}; + use crate::media_layout_common::DestinationVerification; + + #[test] + fn backfill_skips_absent_legacy_source() { + let decision = decide_backfill( + DestinationVerification::LegacySourceAbsent, + "media/aa/bb/hash.png", + ) + .expect("absent legacy source should be skipped"); + + assert_eq!(decision, BackfillDecision::Skip); + } + + #[test] + fn backfill_skips_already_verified_destination() { + let decision = decide_backfill(DestinationVerification::Verified, "media/aa/bb/hash.png") + .expect("verified destination should be skipped"); + + assert_eq!(decision, BackfillDecision::Skip); + } + + #[test] + fn backfill_copies_when_destination_is_missing() { + let decision = decide_backfill( + DestinationVerification::DestinationMissing, + "media/aa/bb/hash.png", + ) + .expect("missing destination should be copied"); + + assert_eq!(decision, BackfillDecision::Copy); + } + + #[test] + fn backfill_fails_closed_on_byte_mismatch() { + let error = decide_backfill( + DestinationVerification::ByteMismatch, + "media/aa/bb/hash.png", + ) + .expect_err("byte mismatch must fail"); + + assert!( + error + .to_string() + .contains("destination verification failed"), + "unexpected error: {error}" + ); + } +} diff --git a/crates/buzz-media/src/bin/buzz-media-layout-delete-legacy.rs b/crates/buzz-media/src/bin/buzz-media-layout-delete-legacy.rs index 3b0de954c..a8cff6152 100644 --- a/crates/buzz-media/src/bin/buzz-media-layout-delete-legacy.rs +++ b/crates/buzz-media/src/bin/buzz-media-layout-delete-legacy.rs @@ -8,7 +8,29 @@ use buzz_media::MediaStorage; use clap::Parser; mod media_layout_common; -use media_layout_common::{verify_destination, CommonArgs}; +use media_layout_common::{verify_destination, CommonArgs, DestinationVerification}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum CleanupVerification { + SkipAbsentSource, + VerifyLegacyKey, +} + +fn decide_cleanup_verification( + verification: DestinationVerification, + sharded_key: &str, + checkpoint: &str, +) -> Result { + match verification { + DestinationVerification::LegacySourceAbsent => Ok(CleanupVerification::SkipAbsentSource), + DestinationVerification::Verified => Ok(CleanupVerification::VerifyLegacyKey), + DestinationVerification::DestinationMissing | DestinationVerification::ByteMismatch => { + anyhow::bail!( + "refusing deletion: sharded destination bytes do not match legacy source: {sharded_key} (checkpoint: {checkpoint})" + ) + } + } +} const CONFIRMATION: &str = "delete-verified-legacy-media"; @@ -59,14 +81,16 @@ async fn verify_selected_destinations( .context("read media sidecar")?; let meta = serde_json::from_slice(&bytes).context("parse media sidecar")?; for object in objects_for_sidecar(community, sha, &meta)? { - if !verify_destination(storage, pacer, &object).await? { - anyhow::bail!( - "refusing deletion: sharded destination bytes do not match legacy source: {} (checkpoint: {sidecar_key})", - object.sharded - ); + let verification = verify_destination(storage, pacer, &object).await?; + match decide_cleanup_verification(verification, &object.sharded, &sidecar_key)? { + CleanupVerification::SkipAbsentSource => { + tracing::info!(source = %object.legacy, destination = %object.sharded, "legacy source absent; skipping deletion candidate"); + } + CleanupVerification::VerifyLegacyKey => { + verified += 1; + verified_legacy_keys.insert(object.legacy); + } } - verified += 1; - verified_legacy_keys.insert(object.legacy); } checkpoint = Some(sidecar_key); } @@ -145,7 +169,8 @@ async fn main() -> Result<()> { #[cfg(test)] mod tests { - use super::ensure_no_start_after; + use super::{decide_cleanup_verification, ensure_no_start_after, CleanupVerification}; + use crate::media_layout_common::DestinationVerification; #[test] fn rejects_start_after_flag() { @@ -167,4 +192,57 @@ mod tests { fn allows_full_bucket_scan() { ensure_no_start_after(None, false).expect("full scan must be allowed"); } + #[test] + fn cleanup_skips_absent_legacy_source() { + let decision = decide_cleanup_verification( + DestinationVerification::LegacySourceAbsent, + "media/aa/bb/hash.png", + "_meta/community/hash.json", + ) + .expect("absent legacy source should be skipped"); + + assert_eq!(decision, CleanupVerification::SkipAbsentSource); + } + + #[test] + fn cleanup_verifies_matching_destination_for_deletion() { + let decision = decide_cleanup_verification( + DestinationVerification::Verified, + "media/aa/bb/hash.png", + "_meta/community/hash.json", + ) + .expect("verified destination should enter deletion set"); + + assert_eq!(decision, CleanupVerification::VerifyLegacyKey); + } + + #[test] + fn cleanup_fails_closed_when_destination_is_missing() { + let error = decide_cleanup_verification( + DestinationVerification::DestinationMissing, + "media/aa/bb/hash.png", + "_meta/community/hash.json", + ) + .expect_err("missing sharded destination must block cleanup"); + + assert!( + error.to_string().contains("refusing deletion"), + "unexpected error: {error}" + ); + } + + #[test] + fn cleanup_fails_closed_on_byte_mismatch() { + let error = decide_cleanup_verification( + DestinationVerification::ByteMismatch, + "media/aa/bb/hash.png", + "_meta/community/hash.json", + ) + .expect_err("byte mismatch must block cleanup"); + + assert!( + error.to_string().contains("refusing deletion"), + "unexpected error: {error}" + ); + } } diff --git a/crates/buzz-media/src/bin/media_layout_common/mod.rs b/crates/buzz-media/src/bin/media_layout_common/mod.rs index 43037eb4f..48a4aa2ed 100644 --- a/crates/buzz-media/src/bin/media_layout_common/mod.rs +++ b/crates/buzz-media/src/bin/media_layout_common/mod.rs @@ -28,21 +28,51 @@ pub struct CommonArgs { pub page_size: usize, } -pub async fn verify_destination( - storage: &MediaStorage, - pacer: &mut RequestPacer, - object: &MigrationObject, -) -> Result { - pacer.wait().await; - let source = storage - .get(&object.legacy) - .await - .with_context(|| format!("read legacy source for verification: {}", object.legacy))?; +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum DestinationVerification { + /// The legacy source is already gone (or never existed). There is nothing + /// left for backfill to copy or cleanup to delete for this object. + LegacySourceAbsent, + /// The legacy source exists and the sharded destination exists with the same bytes. + Verified, + /// The legacy source exists, but the sharded destination has not been created yet. + DestinationMissing, + /// Both objects exist, but the sharded destination does not match the legacy bytes. + ByteMismatch, +} - pacer.wait().await; - let destination = match storage.get(&object.sharded).await { +fn classify_destination_verification( + source: Result, MediaError>, + destination: Result, MediaError>, + object: &MigrationObject, +) -> Result { + if matches!(source, Err(MediaError::NotFound)) { + if let Err(error) = destination { + if !matches!(error, MediaError::NotFound) { + return Err(error).with_context(|| { + format!( + "read sharded destination for verification: {}", + object.sharded + ) + }); + } + } + return Ok(DestinationVerification::LegacySourceAbsent); + } + + let source = match source { Ok(bytes) => bytes, - Err(MediaError::NotFound) => return Ok(false), + Err(MediaError::NotFound) => unreachable!("handled above"), + Err(error) => { + return Err(error).with_context(|| { + format!("read legacy source for verification: {}", object.legacy) + }); + } + }; + + let destination = match destination { + Ok(bytes) => bytes, + Err(MediaError::NotFound) => return Ok(DestinationVerification::DestinationMissing), Err(error) => { return Err(error).with_context(|| { format!( @@ -53,7 +83,25 @@ pub async fn verify_destination( } }; - Ok(destination == source) + Ok(if destination == source { + DestinationVerification::Verified + } else { + DestinationVerification::ByteMismatch + }) +} + +pub async fn verify_destination( + storage: &MediaStorage, + pacer: &mut RequestPacer, + object: &MigrationObject, +) -> Result { + pacer.wait().await; + let source = storage.get(&object.legacy).await; + + pacer.wait().await; + let destination = storage.get(&object.sharded).await; + + classify_destination_verification(source, destination, object) } impl CommonArgs { @@ -84,3 +132,114 @@ impl CommonArgs { .context("create media S3 client") } } + +#[cfg(test)] +mod tests { + use super::*; + + fn object() -> MigrationObject { + MigrationObject { + legacy: "abc.png".to_string(), + sharded: "media/ab/cd/abc.png".to_string(), + } + } + + #[test] + fn classifies_matching_destination_as_verified() { + let result = classify_destination_verification( + Ok(b"same".to_vec()), + Ok(b"same".to_vec()), + &object(), + ) + .expect("verification should classify"); + + assert_eq!(result, DestinationVerification::Verified); + } + + #[test] + fn classifies_absent_legacy_source_as_skip() { + let result = classify_destination_verification( + Err(MediaError::NotFound), + Ok(b"sharded-only".to_vec()), + &object(), + ) + .expect("not-found legacy source should be non-fatal"); + + assert_eq!(result, DestinationVerification::LegacySourceAbsent); + } + + #[test] + fn classifies_present_source_with_missing_destination_as_missing() { + let result = classify_destination_verification( + Ok(b"source".to_vec()), + Err(MediaError::NotFound), + &object(), + ) + .expect("missing destination should be classified"); + + assert_eq!(result, DestinationVerification::DestinationMissing); + } + + #[test] + fn classifies_present_source_with_mismatched_destination_as_mismatch() { + let result = classify_destination_verification( + Ok(b"source".to_vec()), + Ok(b"different".to_vec()), + &object(), + ) + .expect("mismatched destination should be classified"); + + assert_eq!(result, DestinationVerification::ByteMismatch); + } + + #[test] + fn keeps_non_not_found_source_errors_fatal() { + let error = classify_destination_verification( + Err(MediaError::StorageError("source boom".to_string())), + Ok(b"destination".to_vec()), + &object(), + ) + .expect_err("source storage errors must remain fatal"); + + assert!( + error + .to_string() + .contains("read legacy source for verification"), + "unexpected error: {error}" + ); + } + + #[test] + fn keeps_non_not_found_destination_errors_fatal_with_present_source() { + let error = classify_destination_verification( + Ok(b"source".to_vec()), + Err(MediaError::StorageError("destination boom".to_string())), + &object(), + ) + .expect_err("destination storage errors must remain fatal"); + + assert!( + error + .to_string() + .contains("read sharded destination for verification"), + "unexpected error: {error}" + ); + } + + #[test] + fn keeps_non_not_found_destination_errors_fatal_with_absent_source() { + let error = classify_destination_verification( + Err(MediaError::NotFound), + Err(MediaError::StorageError("destination boom".to_string())), + &object(), + ) + .expect_err("destination storage errors must remain fatal"); + + assert!( + error + .to_string() + .contains("read sharded destination for verification"), + "unexpected error: {error}" + ); + } +}