coder
2026-08-11 17:31:21 -04:00
parent a78411559e
commit 1ebf45812e
4 changed files with 58 additions and 15 deletions
@@ -12,6 +12,9 @@ use media_layout_common::{verify_destination, CommonArgs};
struct Args {
#[command(flatten)]
common: CommonArgs,
/// Resume after this sidecar key. The final log line prints the next value.
#[arg(long, env = "BUZZ_MEDIA_MIGRATION_START_AFTER")]
start_after: Option<String>,
/// Report actions without copying objects.
#[arg(long, env = "BUZZ_MEDIA_MIGRATION_DRY_RUN", default_value_t = false)]
dry_run: bool,
@@ -24,7 +27,7 @@ async fn main() -> Result<()> {
let storage = args.common.storage()?;
let mut pacer = RequestPacer::new(args.common.requests_per_second);
let mut continuation = None;
let mut start_after = args.common.start_after.clone();
let mut start_after = args.start_after.clone();
let mut processed = 0_u64;
let mut copied = 0_u64;
let mut skipped = 0_u64;
@@ -1,6 +1,6 @@
//! Delete verified legacy media payloads after migration reconciliation.
use std::collections::BTreeSet;
use std::{collections::BTreeSet, env};
use anyhow::{Context, Result};
use buzz_media::migration::{objects_for_sidecar, parse_sidecar_key, RequestPacer};
@@ -17,6 +17,9 @@ const CONFIRMATION: &str = "delete-verified-legacy-media";
struct Args {
#[command(flatten)]
common: CommonArgs,
/// Rejected for destructive cleanup because partial scans can strand shared legacy CAS keys.
#[arg(long, env = "BUZZ_MEDIA_MIGRATION_START_AFTER")]
start_after: Option<String>,
/// Preview is the safe default. Set false only after reconciliation.
#[arg(long, env = "BUZZ_MEDIA_MIGRATION_DRY_RUN", default_value_t = true, action = clap::ArgAction::Set)]
dry_run: bool,
@@ -35,19 +38,13 @@ async fn verify_selected_destinations(
pacer: &mut RequestPacer,
) -> Result<(u64, BTreeSet<String>, Option<String>)> {
let mut continuation = None;
let mut start_after = common.start_after.clone();
let mut verified = 0_u64;
let mut verified_legacy_keys = BTreeSet::new();
let mut checkpoint = start_after.clone();
let mut checkpoint = None;
loop {
pacer.wait().await;
let page = storage
.list_prefix_page(
"_meta/",
continuation.take(),
start_after.take(),
common.page_size,
)
.list_prefix_page("_meta/", continuation.take(), None, common.page_size)
.await
.context("list media sidecars")?;
for (sidecar_key, _) in page.objects {
@@ -109,10 +106,28 @@ async fn delete_verified_legacy_objects(
Ok((changed, skipped))
}
fn ensure_no_start_after(start_after: Option<&str>, env_start_after_is_set: bool) -> Result<()> {
if start_after.is_some() {
anyhow::bail!(
"--start-after/BUZZ_MEDIA_MIGRATION_START_AFTER is unsafe for legacy deletion; rerun from the beginning instead"
);
}
if env_start_after_is_set {
anyhow::bail!(
"BUZZ_MEDIA_MIGRATION_START_AFTER is unsafe for legacy deletion; rerun from the beginning instead"
);
}
Ok(())
}
#[tokio::main]
async fn main() -> Result<()> {
tracing_subscriber::fmt().with_target(false).init();
let args = Args::parse();
ensure_no_start_after(
args.start_after.as_deref(),
env::var_os("BUZZ_MEDIA_MIGRATION_START_AFTER").is_some(),
)?;
if !args.dry_run && args.confirm.as_deref() != Some(CONFIRMATION) {
anyhow::bail!("destructive mode requires --confirm={CONFIRMATION}");
}
@@ -127,3 +142,29 @@ async fn main() -> Result<()> {
tracing::info!(deleted, skipped, dry_run = args.dry_run, checkpoint = ?verified_checkpoint, "legacy deletion complete");
Ok(())
}
#[cfg(test)]
mod tests {
use super::ensure_no_start_after;
#[test]
fn rejects_start_after_flag() {
let error = ensure_no_start_after(Some("_meta/community/sha.json"), false)
.expect_err("start-after must be unsafe for legacy deletion");
assert!(error.to_string().contains("unsafe for legacy deletion"));
}
#[test]
fn rejects_start_after_env() {
let error = ensure_no_start_after(None, true)
.expect_err("start-after env must be unsafe for legacy deletion");
assert!(error.to_string().contains("unsafe for legacy deletion"));
}
#[test]
fn allows_full_bucket_scan() {
ensure_no_start_after(None, false).expect("full scan must be allowed");
}
}
@@ -24,9 +24,6 @@ pub struct CommonArgs {
default_value_t = 25
)]
pub requests_per_second: u32,
/// Resume after this sidecar key. The final log line prints the next value.
#[arg(long, env = "BUZZ_MEDIA_MIGRATION_START_AFTER")]
pub start_after: Option<String>,
#[arg(long, env = "BUZZ_MEDIA_MIGRATION_PAGE_SIZE", default_value_t = 100)]
pub page_size: usize,
}
+4 -2
View File
@@ -299,8 +299,10 @@ are in [`deploy/kubernetes/examples/`](../../kubernetes/examples/). The tools:
- list canonical `_meta/<community>/<sha>.json` sidecars in bounded pages;
- default to 25 total S3 requests/second, configurable with
`BUZZ_MEDIA_MIGRATION_REQUESTS_PER_SECOND`;
- print the last completed sidecar as `checkpoint`; restart with
`BUZZ_MEDIA_MIGRATION_START_AFTER=<checkpoint>` if a Job fails;
- print the last completed sidecar as `checkpoint`; the backfill job can be
restarted with `BUZZ_MEDIA_MIGRATION_START_AFTER=<checkpoint>` if it fails;
the destructive cleanup job always scans from the beginning so shared legacy
keys are verified for every community before deletion;
- are idempotent: backfill skips an existing destination, and deletion skips an
absent legacy source;
- fail closed on malformed metadata or a missing source/destination.