mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(media): make layout cleanup idempotent
Co-authored-by: coder 0 <d97ebdbb198c7237c94f84ea8bb8a73583ea067407eebd0062abbb3962527fb1@buzz.block.builderlab.xyz> Signed-off-by: coder 0 <d97ebdbb198c7237c94f84ea8bb8a73583ea067407eebd0062abbb3962527fb1@buzz.block.builderlab.xyz>
This commit is contained in:
@@ -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<BackfillDecision> {
|
||||
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}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<CleanupVerification> {
|
||||
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}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -28,21 +28,51 @@ pub struct CommonArgs {
|
||||
pub page_size: usize,
|
||||
}
|
||||
|
||||
pub async fn verify_destination(
|
||||
storage: &MediaStorage,
|
||||
pacer: &mut RequestPacer,
|
||||
object: &MigrationObject,
|
||||
) -> Result<bool> {
|
||||
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<Vec<u8>, MediaError>,
|
||||
destination: Result<Vec<u8>, MediaError>,
|
||||
object: &MigrationObject,
|
||||
) -> Result<DestinationVerification> {
|
||||
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<DestinationVerification> {
|
||||
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}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user