use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
use boatramp_core::{ByteStream, PutMeta, Storage, StorageError};
use futures::StreamExt;
#[cfg(feature = "blob-migrate-gate-mutation")]
pub(crate) mod gate_mutation {
pub(crate) fn env_on(name: &str) -> bool {
std::env::var(name)
.map(|v| !v.is_empty() && v != "0")
.unwrap_or(false)
}
}
#[derive(Debug, Clone)]
pub struct MigrateOptions {
pub concurrency: usize,
pub verify: bool,
pub dry_run: bool,
pub prefix: String,
}
impl Default for MigrateOptions {
fn default() -> Self {
Self {
concurrency: 8,
verify: true,
dry_run: false,
prefix: String::new(),
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct MigrateReport {
pub total_objects: u64,
pub copied_objects: u64,
pub skipped_objects: u64,
pub copied_bytes: u64,
pub verified: bool,
}
#[derive(Debug, thiserror::Error)]
pub enum MigrateError {
#[error("blob storage: {0}")]
Storage(#[from] StorageError),
#[error(
"verification FAILED: {missing} source object(s) absent from the destination \
(first {shown} shown): {keys:?}"
)]
VerifyMissing {
missing: usize,
shown: usize,
keys: Vec<String>,
},
}
pub const VERIFY_REPORT_CAP: usize = 50;
enum PlanItem {
Copy {
key: String,
size: u64,
content_type: Option<String>,
},
Skip,
}
struct ItemOutcome {
copied: bool,
bytes: u64,
}
pub async fn migrate(
source: Arc<dyn Storage>,
dest: Arc<dyn Storage>,
opts: &MigrateOptions,
) -> Result<MigrateReport, MigrateError> {
let concurrency = opts.concurrency.max(1);
#[cfg_attr(not(feature = "blob-migrate-gate-mutation"), allow(unused_mut))]
let mut objects = source.list(&opts.prefix).await?;
let total_objects = objects.len() as u64;
#[cfg(feature = "blob-migrate-gate-mutation")]
if gate_mutation::env_on("BOATRAMP_BLOBMIG_MUTATE_DROP_LAST") {
objects.sort_by(|a, b| a.key.cmp(&b.key));
objects.pop();
}
tracing::info!(
total = total_objects,
prefix = %opts.prefix,
dry_run = opts.dry_run,
concurrency,
"blob migrate: enumerated source objects"
);
let copied_objects = Arc::new(AtomicU64::new(0));
let skipped_objects = Arc::new(AtomicU64::new(0));
let copied_bytes = Arc::new(AtomicU64::new(0));
let done_objects = Arc::new(AtomicU64::new(0));
let progress = Arc::new(std::sync::Mutex::new(Progress::new(total_objects)));
let results: Vec<Result<ItemOutcome, MigrateError>> = futures::stream::iter(objects)
.map(|meta| {
let source = source.clone();
let dest = dest.clone();
let dry_run = opts.dry_run;
async move { copy_one(&source, &dest, meta, dry_run).await }
})
.buffer_unordered(concurrency)
.map(|outcome| {
if let Ok(ref o) = outcome {
if o.copied {
copied_objects.fetch_add(1, Ordering::Relaxed);
copied_bytes.fetch_add(o.bytes, Ordering::Relaxed);
} else {
skipped_objects.fetch_add(1, Ordering::Relaxed);
}
}
let done = done_objects.fetch_add(1, Ordering::Relaxed) + 1;
let (c, s, b) = (
copied_objects.load(Ordering::Relaxed),
skipped_objects.load(Ordering::Relaxed),
copied_bytes.load(Ordering::Relaxed),
);
if let Ok(mut p) = progress.lock() {
p.maybe_log(done, c, s, b, opts.dry_run);
}
outcome
})
.collect()
.await;
for r in results {
r?;
}
let copied_objects = copied_objects.load(Ordering::Relaxed);
let skipped_objects = skipped_objects.load(Ordering::Relaxed);
let copied_bytes = copied_bytes.load(Ordering::Relaxed);
tracing::info!(
total = total_objects,
copied = copied_objects,
skipped = skipped_objects,
copied_bytes,
dry_run = opts.dry_run,
"blob migrate: copy phase complete"
);
let mut verified = false;
if opts.verify && !opts.dry_run {
let missing = verify(&source, &dest, &opts.prefix, concurrency).await?;
if !missing.is_empty() {
let shown = missing.len().min(VERIFY_REPORT_CAP);
return Err(MigrateError::VerifyMissing {
missing: missing.len(),
shown,
keys: missing.into_iter().take(VERIFY_REPORT_CAP).collect(),
});
}
verified = true;
tracing::info!(
objects = total_objects,
"VERIFY OK: all source objects present in destination"
);
}
Ok(MigrateReport {
total_objects,
copied_objects,
skipped_objects,
copied_bytes,
verified,
})
}
async fn plan_one(
dest: &Arc<dyn Storage>,
key: &str,
source_size: u64,
content_type: Option<String>,
) -> Result<PlanItem, MigrateError> {
#[cfg(feature = "blob-migrate-gate-mutation")]
if gate_mutation::env_on("BOATRAMP_BLOBMIG_MUTATE_SKIP_ALWAYS") {
return Ok(PlanItem::Skip);
}
match dest.head(key).await {
Ok(meta) if meta.size == Some(source_size) => Ok(PlanItem::Skip),
Ok(_) => Ok(PlanItem::Copy {
key: key.to_string(),
size: source_size,
content_type,
}),
Err(StorageError::NotFound(_)) => Ok(PlanItem::Copy {
key: key.to_string(),
size: source_size,
content_type,
}),
Err(e) => Err(MigrateError::Storage(e)),
}
}
async fn copy_one(
source: &Arc<dyn Storage>,
dest: &Arc<dyn Storage>,
meta: boatramp_core::ObjectMeta,
dry_run: bool,
) -> Result<ItemOutcome, MigrateError> {
let source_size = meta.size.unwrap_or(0);
let plan = plan_one(dest, &meta.key, source_size, meta.content_type.clone()).await?;
match plan {
PlanItem::Skip => Ok(ItemOutcome {
copied: false,
bytes: 0,
}),
PlanItem::Copy {
key,
size,
content_type,
} => {
if dry_run {
return Ok(ItemOutcome {
copied: true,
bytes: 0,
});
}
let got = source.get(&key).await?;
let ct = got.meta.content_type.or(content_type);
let dest_key = dest_key_for(&key);
let body: ByteStream = got.body;
let written = dest
.put(&dest_key, body, PutMeta { content_type: ct })
.await?;
#[cfg(feature = "blob-migrate-gate-mutation")]
if gate_mutation::env_on("BOATRAMP_BLOBMIG_MUTATE_DELETE_SOURCE") {
source.delete(&key).await?;
}
Ok(ItemOutcome {
copied: true,
bytes: written.size.unwrap_or(size),
})
}
}
}
fn dest_key_for(source_key: &str) -> String {
#[cfg(feature = "blob-migrate-gate-mutation")]
if gate_mutation::env_on("BOATRAMP_BLOBMIG_MUTATE_REWRITE_KEY") {
return format!("MANGLED/{source_key}");
}
source_key.to_string()
}
async fn verify(
source: &Arc<dyn Storage>,
dest: &Arc<dyn Storage>,
prefix: &str,
concurrency: usize,
) -> Result<Vec<String>, MigrateError> {
let source_keys = source.list(prefix).await?;
let checks: Vec<Result<Option<String>, MigrateError>> = futures::stream::iter(source_keys)
.map(|meta| {
let dest = dest.clone();
async move {
match dest.head(&meta.key).await {
Ok(_) => Ok(None),
Err(StorageError::NotFound(_)) => Ok(Some(meta.key)),
Err(e) => Err(MigrateError::Storage(e)),
}
}
})
.buffer_unordered(concurrency)
.collect()
.await;
let mut missing = Vec::new();
for c in checks {
if let Some(key) = c? {
missing.push(key);
}
}
missing.sort();
Ok(missing)
}
struct Progress {
total: u64,
last: Instant,
}
const PROGRESS_INTERVAL: Duration = Duration::from_secs(2);
const PROGRESS_EVERY_N: u64 = 500;
impl Progress {
fn new(total: u64) -> Self {
Self {
total,
last: Instant::now(),
}
}
fn maybe_log(&mut self, done: u64, copied: u64, skipped: u64, bytes: u64, dry_run: bool) {
let elapsed = self.last.elapsed();
if elapsed >= PROGRESS_INTERVAL
|| done.is_multiple_of(PROGRESS_EVERY_N)
|| done == self.total
{
self.last = Instant::now();
tracing::info!(
done,
total = self.total,
copied,
skipped,
copied_bytes = bytes,
dry_run,
"blob migrate: progress"
);
}
}
}
#[cfg(all(test, feature = "blob-migrate-gate-mutation"))]
mod gate {
use super::*;
use boatramp_storage::FsStorage;
fn seed_keys() -> Vec<(&'static str, &'static [u8])> {
vec![
(
"ab/0000000000000000000000000000000000000000000000000000000000000000",
b"content-addressed-immutable-blob",
),
("manifests/site-alpha", b"{\"deployment\":\"d1\"}"),
("hblob/proj~site/uploads/report.json", b"{\"ok\":true}"),
]
}
async fn put(store: &Arc<dyn Storage>, key: &str, bytes: &[u8]) {
let owned = bytes.to_vec();
let body: ByteStream =
futures::stream::once(async move { Ok(bytes::Bytes::from(owned)) }).boxed();
store
.put(key, body, PutMeta::default())
.await
.expect("seed put");
}
async fn seeded_backends(
tmp: &std::path::Path,
) -> (Arc<dyn Storage>, Arc<dyn Storage>, Vec<String>) {
let source: Arc<dyn Storage> = Arc::new(FsStorage::new(tmp.join("src")));
let dest: Arc<dyn Storage> = Arc::new(FsStorage::new(tmp.join("dst")));
let mut keys = Vec::new();
for (k, v) in seed_keys() {
put(&source, k, v).await;
keys.push(k.to_string());
}
keys.sort();
(source, dest, keys)
}
async fn keys_of(store: &Arc<dyn Storage>) -> Vec<String> {
let mut ks: Vec<String> = store
.list("")
.await
.expect("list")
.into_iter()
.map(|m| m.key)
.collect();
ks.sort();
ks
}
async fn invariant_1_completeness(dest: &Arc<dyn Storage>, source_keys: &[String]) {
for key in source_keys {
dest.head(key).await.unwrap_or_else(|_| {
panic!("I1 completeness: source key {key:?} is absent from the destination")
});
}
}
async fn invariant_2_head_skip_sound(
source: &Arc<dyn Storage>,
dest: &Arc<dyn Storage>,
source_keys: &[String],
) {
let report = migrate(source.clone(), dest.clone(), &MigrateOptions::default())
.await
.expect("second migrate (idempotent) succeeds");
assert_eq!(
report.copied_objects, 0,
"I2: a re-run over a fully-populated destination must copy nothing (all head-skipped)"
);
assert_eq!(
report.skipped_objects,
source_keys.len() as u64,
"I2: every object must be head-skipped on the idempotent re-run"
);
invariant_1_completeness(dest, source_keys).await;
}
async fn invariant_3_key_fidelity(source: &Arc<dyn Storage>, dest: &Arc<dyn Storage>) {
let src = keys_of(source).await;
let dst = keys_of(dest).await;
assert_eq!(
src, dst,
"I3 key fidelity: destination keys must equal source keys byte-exact (no mangling)"
);
}
async fn invariant_4_source_read_only(source: &Arc<dyn Storage>, before: &[String]) {
let after = keys_of(source).await;
assert_eq!(
before,
after.as_slice(),
"I4 source read-only: the source object set must be unchanged after migrate"
);
}
#[tokio::test]
async fn blob_migrate_complete_gate() {
let tmp = tempfile::tempdir().expect("tempdir");
let (source, dest, source_keys) = seeded_backends(tmp.path()).await;
let before_source = keys_of(&source).await;
assert_eq!(before_source, source_keys, "sanity: seeded source keys");
let report = migrate(source.clone(), dest.clone(), &MigrateOptions::default())
.await
.expect("migrate succeeds on a clean run");
assert_eq!(report.total_objects, source_keys.len() as u64);
assert!(report.verified, "verify must have run + passed");
invariant_1_completeness(&dest, &source_keys).await;
invariant_4_source_read_only(&source, &before_source).await;
invariant_3_key_fidelity(&source, &dest).await;
invariant_2_head_skip_sound(&source, &dest, &source_keys).await;
println!(
"BLOB MIGRATE COMPLETE OK: the offline blob-backend migration copies every source \
object to the destination (completeness), skips only a matching-size head-present key \
(idempotent/resumable), preserves keys byte-exact (no tenant/hblob-boundary collapse), \
and never mutates the source (read-only). Mutation-verified: \
BOATRAMP_BLOBMIG_MUTATE_{{DROP_LAST,SKIP_ALWAYS,REWRITE_KEY,DELETE_SOURCE}}=1 each FAIL \
this gate."
);
}
}