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 3ff68d977..2238dcd49 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 @@ -1,5 +1,7 @@ //! Delete verified legacy media payloads after migration reconciliation. +use std::collections::BTreeSet; + use anyhow::{Context, Result}; use buzz_media::migration::{objects_for_sidecar, parse_sidecar_key, RequestPacer}; use buzz_media::MediaStorage; @@ -25,18 +27,17 @@ struct Args { /// Scan all selected sidecars before deleting anything. This two-pass design is /// important because one flat legacy CAS key can serve multiple communities: -/// every community destination must exist before the shared source is removed. -async fn scan( +/// every community destination must be byte-verified while the shared source is +/// still present before the shared source is removed. +async fn verify_selected_destinations( storage: &MediaStorage, common: &CommonArgs, pacer: &mut RequestPacer, - delete: bool, - dry_run: bool, -) -> Result<(u64, u64, Option)> { +) -> Result<(u64, BTreeSet, Option)> { let mut continuation = None; let mut start_after = common.start_after.clone(); - let mut changed = 0_u64; - let mut skipped = 0_u64; + let mut verified = 0_u64; + let mut verified_legacy_keys = BTreeSet::new(); let mut checkpoint = start_after.clone(); loop { pacer.wait().await; @@ -61,34 +62,14 @@ async fn scan( .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 !delete { - pacer.wait().await; - if !storage.head(&object.sharded).await? { - anyhow::bail!( - "refusing deletion: sharded destination missing: {} (checkpoint: {sidecar_key})", - object.sharded - ); - } - continue; - } - pacer.wait().await; - if !storage.head(&object.legacy).await? { - skipped += 1; - continue; - } - if dry_run { - tracing::info!(key = %object.legacy, "would delete verified legacy object"); - } else { - if !verify_destination(storage, pacer, &object).await? { - anyhow::bail!( - "refusing deletion: sharded destination bytes do not match legacy source: {} (checkpoint: {sidecar_key})", - object.sharded - ); - } - pacer.wait().await; - storage.delete(&object.legacy).await?; - changed += 1; + if !verify_destination(storage, pacer, &object).await? { + anyhow::bail!( + "refusing deletion: sharded destination bytes do not match legacy source: {} (checkpoint: {sidecar_key})", + object.sharded + ); } + verified += 1; + verified_legacy_keys.insert(object.legacy); } checkpoint = Some(sidecar_key); } @@ -100,7 +81,32 @@ async fn scan( anyhow::bail!("truncated S3 listing returned no continuation token"); } } - Ok((changed, skipped, checkpoint)) + Ok((verified, verified_legacy_keys, checkpoint)) +} + +async fn delete_verified_legacy_objects( + storage: &MediaStorage, + pacer: &mut RequestPacer, + verified_legacy_keys: BTreeSet, + dry_run: bool, +) -> Result<(u64, u64)> { + let mut changed = 0_u64; + let mut skipped = 0_u64; + for legacy_key in verified_legacy_keys { + if dry_run { + tracing::info!(key = %legacy_key, "would delete verified legacy object"); + continue; + } + pacer.wait().await; + if !storage.head(&legacy_key).await? { + skipped += 1; + continue; + } + pacer.wait().await; + storage.delete(&legacy_key).await?; + changed += 1; + } + Ok((changed, skipped)) } #[tokio::main] @@ -112,11 +118,12 @@ async fn main() -> Result<()> { } let storage = args.common.storage()?; let mut pacer = RequestPacer::new(args.common.requests_per_second); - let (_, _, verified_checkpoint) = - scan(&storage, &args.common, &mut pacer, false, args.dry_run).await?; - tracing::info!(checkpoint = ?verified_checkpoint, "all selected sharded destinations exist; beginning deletion pass"); - let (deleted, skipped, checkpoint) = - scan(&storage, &args.common, &mut pacer, true, args.dry_run).await?; - tracing::info!(deleted, skipped, dry_run = args.dry_run, checkpoint = ?checkpoint, "legacy deletion complete"); + let (verified, verified_legacy_keys, verified_checkpoint) = + verify_selected_destinations(&storage, &args.common, &mut pacer).await?; + tracing::info!(verified, checkpoint = ?verified_checkpoint, "all selected sharded destinations match legacy sources; beginning deletion pass"); + let (deleted, skipped) = + delete_verified_legacy_objects(&storage, &mut pacer, verified_legacy_keys, args.dry_run) + .await?; + tracing::info!(deleted, skipped, dry_run = args.dry_run, checkpoint = ?verified_checkpoint, "legacy deletion complete"); Ok(()) }