use std::{
fs, io,
path::{Path, PathBuf},
};
use eyre::{Context, Result, ensure};
use future_form::Sendable;
use sedimentree_fs_storage::FsStorage;
use subduction_core::storage::traits::Storage;
use subduction_redb_storage::RedbStorage;
use crate::keyhive::{ARCHIVES_SUBDIR, KEYHIVE_DIR, OPS_SUBDIR};
#[derive(Debug, clap::Parser)]
pub(crate) struct MigrateArgs {
#[arg(long)]
pub(crate) from: PathBuf,
#[arg(long)]
pub(crate) to: PathBuf,
#[arg(long, default_value_t = 1000)]
pub(crate) progress_every: usize,
#[arg(long, default_value_t = false)]
pub(crate) dry_run: bool,
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
struct MigrationStats {
total: usize,
migrated: usize,
skipped: usize,
commits: usize,
fragments: usize,
keyhive_archives: usize,
keyhive_ops: usize,
}
pub(crate) async fn run(args: MigrateArgs) -> Result<()> {
ensure!(
args.from != args.to,
"--from and --to must differ: pick a fresh destination directory so the \
source filesystem store is left untouched (the migration is non-destructive)"
);
ensure!(
args.from.exists(),
"source directory does not exist: {}",
args.from.display()
);
let source = FsStorage::new(args.from.clone()).wrap_err("open source filesystem store")?;
tracing::info!(
from = %args.from.display(),
to = %args.to.display(),
dry_run = args.dry_run,
"starting filesystem → redb migration"
);
let dest = if args.dry_run {
None
} else {
Some(RedbStorage::new(&args.to).wrap_err("open destination redb store")?)
};
let mut stats = migrate_all(&source, dest.as_ref(), args.progress_every).await?;
let (archives, ops) =
migrate_keyhive(&args.from, &args.to, args.dry_run).wrap_err("copy keyhive state")?;
stats.keyhive_archives = archives;
stats.keyhive_ops = ops;
tracing::info!(
total = stats.total,
migrated = stats.migrated,
skipped = stats.skipped,
commits = stats.commits,
fragments = stats.fragments,
keyhive_archives = stats.keyhive_archives,
keyhive_ops = stats.keyhive_ops,
dry_run = args.dry_run,
"filesystem → redb migration complete"
);
let keyhive = stats.keyhive_archives + stats.keyhive_ops;
if args.dry_run {
println!(
"[dry run] Would migrate {} trees ({} commits, {} fragments) and {keyhive} keyhive \
files ({} archives, {} ops). Source totals; destination not inspected.",
stats.total, stats.commits, stats.fragments, stats.keyhive_archives, stats.keyhive_ops,
);
} else {
println!(
"Migrated {} trees ({} commits, {} fragments) and {keyhive} keyhive files \
({} archives, {} ops); skipped {} trees already present; {} trees total.",
stats.migrated,
stats.commits,
stats.fragments,
stats.keyhive_archives,
stats.keyhive_ops,
stats.skipped,
stats.total,
);
}
Ok(())
}
async fn migrate_all(
source: &FsStorage,
dest: Option<&RedbStorage>,
progress_every: usize,
) -> Result<MigrationStats> {
let progress_every = progress_every.max(1);
let ids = Storage::<Sendable>::load_all_sedimentree_ids(source)
.await
.wrap_err("enumerate source sedimentrees")?;
let mut stats = MigrationStats {
total: ids.len(),
..MigrationStats::default()
};
for (processed, id) in ids.into_iter().enumerate() {
let already_present = match dest {
Some(dest) => Storage::<Sendable>::contains_sedimentree_id(dest, id)
.await
.wrap_err("check destination for existing tree")?,
None => false,
};
if already_present {
stats.skipped += 1;
} else if let Some(dest) = dest {
let commits = Storage::<Sendable>::load_loose_commits(source, id)
.await
.wrap_err_with(|| format!("load commits for {id:?}"))?;
let fragments = Storage::<Sendable>::load_fragments(source, id)
.await
.wrap_err_with(|| format!("load fragments for {id:?}"))?;
stats.commits += commits.len();
stats.fragments += fragments.len();
Storage::<Sendable>::save_batch(dest, id, commits, fragments)
.await
.wrap_err_with(|| format!("write tree {id:?} into redb"))?;
stats.migrated += 1;
} else {
stats.commits += Storage::<Sendable>::load_loose_commit_metas(source, id)
.await
.wrap_err_with(|| format!("load commit metas for {id:?}"))?
.len();
stats.fragments += Storage::<Sendable>::load_fragment_metas(source, id)
.await
.wrap_err_with(|| format!("load fragment metas for {id:?}"))?
.len();
}
if (processed + 1) % progress_every == 0 {
tracing::info!(
processed = processed + 1,
total = stats.total,
migrated = stats.migrated,
skipped = stats.skipped,
"migration progress"
);
}
}
Ok(stats)
}
fn migrate_keyhive(from: &Path, to: &Path, dry_run: bool) -> Result<(usize, usize)> {
let from_root = from.join(KEYHIVE_DIR);
let to_root = to.join(KEYHIVE_DIR);
let archives = copy_keyhive_subdir(
&from_root.join(ARCHIVES_SUBDIR),
&to_root.join(ARCHIVES_SUBDIR),
dry_run,
)
.wrap_err("copy keyhive archives")?;
let ops = copy_keyhive_subdir(
&from_root.join(OPS_SUBDIR),
&to_root.join(OPS_SUBDIR),
dry_run,
)
.wrap_err("copy keyhive ops")?;
Ok((archives, ops))
}
fn copy_keyhive_subdir(src: &Path, dst: &Path, dry_run: bool) -> io::Result<usize> {
let entries = match fs::read_dir(src) {
Ok(entries) => entries,
Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(0),
Err(e) => return Err(e),
};
let mut copied = 0;
let mut created_dst = false;
for entry in entries {
let path = entry?.path();
if path.extension().and_then(|e| e.to_str()) != Some("bin") {
continue;
}
let Some(name) = path.file_name() else {
continue;
};
let target = dst.join(name);
if target.try_exists()? {
continue; }
copied += 1;
if dry_run {
continue;
}
if !created_dst {
fs::create_dir_all(dst)?;
created_dst = true;
}
durable_copy(&path, &target)?;
}
if created_dst {
fsync_dir(dst)?;
}
Ok(copied)
}
fn durable_copy(src: &Path, dst: &Path) -> io::Result<()> {
let tmp = dst.with_extension("bin.tmp");
let staged = (|| {
fs::copy(src, &tmp)?;
fs::File::open(&tmp)?.sync_all()?;
fs::rename(&tmp, dst)
})();
if staged.is_err() {
fs::remove_file(&tmp).ok();
}
staged
}
fn fsync_dir(dir: &Path) -> io::Result<()> {
fs::File::open(dir)?.sync_all()?;
Ok(())
}
#[cfg(test)]
mod tests {
use std::collections::BTreeSet;
use sedimentree_core::{
blob::{Blob, verified::VerifiedBlobMeta},
crypto::digest::Digest,
fragment::Fragment,
id::SedimentreeId,
loose_commit::{LooseCommit, id::CommitId},
};
use subduction_crypto::{signer::memory::MemorySigner, verified_meta::VerifiedMeta};
use super::*;
fn signer() -> MemorySigner {
MemorySigner::from_bytes(&[7u8; 32])
}
async fn commit(
s: &MemorySigner,
id: SedimentreeId,
head: u8,
blob_len: usize,
) -> VerifiedMeta<LooseCommit> {
VerifiedMeta::seal::<Sendable, _>(
s,
(id, CommitId::new([head; 32]), BTreeSet::new()),
VerifiedBlobMeta::new(Blob::new(vec![head; blob_len])),
)
.await
}
async fn fragment(
s: &MemorySigner,
id: SedimentreeId,
head: u8,
blob_len: usize,
) -> VerifiedMeta<Fragment> {
VerifiedMeta::seal::<Sendable, _>(
s,
(
id,
CommitId::new([head; 32]),
BTreeSet::from([CommitId::new([0xF0; 32])]),
vec![CommitId::new([0xF1; 32])],
),
VerifiedBlobMeta::new(Blob::new(vec![head; blob_len])),
)
.await
}
async fn commit_digests<S>(s: &S, id: SedimentreeId) -> Result<BTreeSet<Digest<LooseCommit>>>
where
S: Storage<Sendable>,
S::Error: Send + Sync + 'static,
{
Ok(Storage::<Sendable>::load_loose_commits(s, id)
.await
.wrap_err("load commits")?
.iter()
.map(|vm| Digest::hash(vm.payload()))
.collect())
}
async fn fragment_digests<S>(s: &S, id: SedimentreeId) -> Result<BTreeSet<Digest<Fragment>>>
where
S: Storage<Sendable>,
S::Error: Send + Sync + 'static,
{
Ok(Storage::<Sendable>::load_fragments(s, id)
.await
.wrap_err("load fragments")?
.iter()
.map(|vm| Digest::hash(vm.payload()))
.collect())
}
#[tokio::test]
async fn migrates_all_trees_and_is_resumable() -> Result<()> {
let src_dir = tempfile::tempdir()?;
let dst_dir = tempfile::tempdir()?;
let source = FsStorage::new(src_dir.path().to_path_buf())?;
let dest = RedbStorage::new(dst_dir.path())?;
let s = signer();
let tree_a = SedimentreeId::new([0xA1; 32]);
let tree_b = SedimentreeId::new([0xB2; 32]);
let a_commits = vec![
commit(&s, tree_a, 0x01, 16).await,
commit(&s, tree_a, 0x02, 20 * 1024).await,
];
let a_fragments = vec![fragment(&s, tree_a, 0x03, 32).await];
Storage::<Sendable>::save_batch(&source, tree_a, a_commits, a_fragments).await?;
let b_commits = vec![commit(&s, tree_b, 0x10, 64).await];
Storage::<Sendable>::save_batch(&source, tree_b, b_commits, Vec::new()).await?;
let stats = migrate_all(&source, Some(&dest), 1).await?;
assert_eq!(
stats,
MigrationStats {
total: 2,
migrated: 2,
skipped: 0,
commits: 3,
fragments: 1,
keyhive_archives: 0,
keyhive_ops: 0,
}
);
for id in [tree_a, tree_b] {
assert!(
Storage::<Sendable>::contains_sedimentree_id(&dest, id).await?,
"migrated tree {id:?} must be registered in the destination"
);
assert_eq!(
commit_digests(&source, id).await?,
commit_digests(&dest, id).await?,
"commit content must survive migration for {id:?}"
);
assert_eq!(
fragment_digests(&source, id).await?,
fragment_digests(&dest, id).await?,
"fragment content must survive migration for {id:?}"
);
}
let again = migrate_all(&source, Some(&dest), 1).await?;
assert_eq!(
again,
MigrationStats {
total: 2,
migrated: 0,
skipped: 2,
commits: 0,
fragments: 0,
keyhive_archives: 0,
keyhive_ops: 0,
}
);
Ok(())
}
fn seed_keyhive(root: &Path, archives: &[&str], ops: &[&str]) -> Result<()> {
for (sub, stems) in [(ARCHIVES_SUBDIR, archives), (OPS_SUBDIR, ops)] {
let dir = root.join(KEYHIVE_DIR).join(sub);
std::fs::create_dir_all(&dir)?;
for stem in stems {
std::fs::write(
dir.join(format!("{stem}.bin")),
format!("data-{stem}").as_bytes(),
)?;
}
}
let tmp = root.join(KEYHIVE_DIR).join("tmp");
std::fs::create_dir_all(&tmp)?;
std::fs::write(tmp.join("stale.bin.tmp"), b"junk")?;
Ok(())
}
#[test]
fn migrate_keyhive_copies_archives_ops_and_is_idempotent() -> Result<()> {
let from = tempfile::tempdir()?;
let to = tempfile::tempdir()?;
seed_keyhive(from.path(), &["aa", "bb"], &["cc"])?;
let (archives, ops) = migrate_keyhive(from.path(), to.path(), false)?;
assert_eq!((archives, ops), (2, 1));
let to_kh = to.path().join(KEYHIVE_DIR);
for (sub, stem) in [
(ARCHIVES_SUBDIR, "aa"),
(ARCHIVES_SUBDIR, "bb"),
(OPS_SUBDIR, "cc"),
] {
let copied = to_kh.join(sub).join(format!("{stem}.bin"));
assert_eq!(
std::fs::read(&copied)?,
format!("data-{stem}").into_bytes(),
"{} must be copied with identical content",
copied.display()
);
}
assert!(
!to_kh.join("tmp").join("stale.bin.tmp").exists(),
"tmp/ must not be migrated"
);
let (archives, ops) = migrate_keyhive(from.path(), to.path(), false)?;
assert_eq!((archives, ops), (0, 0));
Ok(())
}
#[tokio::test]
async fn dry_run_writes_nothing() -> Result<()> {
let src_dir = tempfile::tempdir()?;
let dst_dir = tempfile::tempdir()?;
let source = FsStorage::new(src_dir.path().to_path_buf())?;
let s = signer();
let tree = SedimentreeId::new([0xC3; 32]);
Storage::<Sendable>::save_batch(
&source,
tree,
vec![commit(&s, tree, 0x01, 16).await],
vec![fragment(&s, tree, 0x02, 32).await],
)
.await?;
seed_keyhive(src_dir.path(), &["aa"], &["bb", "cc"])?;
let dst = dst_dir.path().join("redb-out");
let stats = migrate_all(&source, None, 1).await?;
let (archives, ops) = migrate_keyhive(src_dir.path(), &dst, true)?;
assert_eq!(stats.total, 1);
assert_eq!(stats.commits, 1);
assert_eq!(stats.fragments, 1);
assert_eq!(stats.migrated, 0, "dry run must migrate nothing");
assert_eq!(stats.skipped, 0, "dry run does not inspect the destination");
assert_eq!((archives, ops), (1, 2));
assert!(
!dst.exists(),
"dry run must not create the destination directory"
);
Ok(())
}
#[tokio::test]
async fn run_rejects_same_from_and_to() -> Result<()> {
let dir = tempfile::tempdir()?;
let args = MigrateArgs {
from: dir.path().to_path_buf(),
to: dir.path().to_path_buf(),
progress_every: 1000,
dry_run: false,
};
let result = run(args).await;
assert!(
matches!(&result, Err(e) if e.to_string().contains("must differ")),
"expected the distinct-directory guard, got: {result:?}"
);
Ok(())
}
#[tokio::test]
async fn run_rejects_missing_source() -> Result<()> {
let dir = tempfile::tempdir()?;
let args = MigrateArgs {
from: dir.path().join("does-not-exist"),
to: dir.path().join("dest"),
progress_every: 1000,
dry_run: false,
};
let result = run(args).await;
assert!(
matches!(&result, Err(e) if e.to_string().contains("does not exist")),
"expected the missing-source guard, got: {result:?}"
);
Ok(())
}
#[test]
#[allow(clippy::expect_used, clippy::cast_possible_truncation)]
fn prop_migrate_preserves_content() {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("current-thread runtime");
let s = signer();
bolero::check!()
.with_generator(((0u8..=3, 0u8..=2), (0u8..=3, 0u8..=2), (0u8..=3, 0u8..=2)))
.for_each(|&(a, b, c)| {
let corpus = [a, b, c];
rt.block_on(async {
let src_dir = tempfile::tempdir().expect("src tempdir");
let dst_dir = tempfile::tempdir().expect("dst tempdir");
let source = FsStorage::new(src_dir.path().to_path_buf()).expect("fs source");
let dest = RedbStorage::new(dst_dir.path()).expect("redb dest");
let mut tree_ids = Vec::new();
for (i, &(n_c_u, n_f_u)) in corpus.iter().enumerate() {
let n_c = usize::from(n_c_u);
let n_f = usize::from(n_f_u);
if n_c + n_f == 0 {
continue;
}
let mut id_bytes = [0u8; 32];
id_bytes[0] = i as u8 + 1;
let id = SedimentreeId::new(id_bytes);
tree_ids.push(id);
let mut commits = Vec::new();
for j in 0..n_c {
let len = if j % 2 == 0 { 32 } else { 300 };
commits.push(commit(&s, id, j as u8, len).await);
}
let mut frags = Vec::new();
for j in 0..n_f {
frags.push(fragment(&s, id, 0xF0 + j as u8, 64).await);
}
Storage::<Sendable>::save_batch(&source, id, commits, frags)
.await
.expect("populate source");
}
let stats = migrate_all(&source, Some(&dest), 1).await.expect("migrate");
assert_eq!(stats.migrated, tree_ids.len(), "every tree migrated");
for &id in &tree_ids {
assert_eq!(
commit_digests(&source, id).await.expect("src commits"),
commit_digests(&dest, id).await.expect("dst commits"),
"commit content preserved for {id:?}"
);
assert_eq!(
fragment_digests(&source, id).await.expect("src fragments"),
fragment_digests(&dest, id).await.expect("dst fragments"),
"fragment content preserved for {id:?}"
);
}
let src_ids: BTreeSet<_> =
Storage::<Sendable>::load_all_sedimentree_ids(&source)
.await
.expect("src ids")
.into_iter()
.collect();
let dst_ids: BTreeSet<_> = Storage::<Sendable>::load_all_sedimentree_ids(&dest)
.await
.expect("dst ids")
.into_iter()
.collect();
assert_eq!(src_ids, dst_ids, "tree-id set preserved");
let again = migrate_all(&source, Some(&dest), 1)
.await
.expect("re-migrate");
assert_eq!(again.migrated, 0, "second pass is a no-op");
});
});
}
}