#![cfg(feature = "persy")]
use beam::migration::{Backend, MigrateOpts, migrate};
use beam::types::{NodeData, Value};
use std::collections::BTreeMap;
use std::env;
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{SystemTime, UNIX_EPOCH};
static TEST_COUNTER: AtomicU64 = AtomicU64::new(0);
fn temp_path(name: &str, ext: &str) -> PathBuf {
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_nanos();
let counter = TEST_COUNTER.fetch_add(1, Ordering::SeqCst);
env::temp_dir().join(format!(
"beam-migrate-{}-{}-{}-{}{}",
name,
std::process::id(),
nanos,
counter,
ext
))
}
fn write_redb_records(path: &std::path::Path, count: usize) -> Result<usize, String> {
use redb::{Database, TableDefinition};
const BEAM_NODES: TableDefinition<&str, &[u8]> = TableDefinition::new("beam_nodes_v1");
let db = Database::create(path).map_err(|e| format!("create: {:?}", e))?;
let txn = db
.begin_write()
.map_err(|e| format!("begin_write: {:?}", e))?;
{
let mut table = txn
.open_table(BEAM_NODES)
.map_err(|e| format!("open_table: {:?}", e))?;
for i in 0..count {
let mut children: BTreeMap<String, NodeData> = BTreeMap::new();
children.insert(
format!("leaf_{:04}", i),
NodeData {
value: Value::Text(format!("test-{}", i)),
updated_at: 1_700_000_000.0 + i as f64,
},
);
let bytes =
postcard::to_allocvec(&children).map_err(|e| format!("postcard: {:?}", e))?;
let key = format!("node_{:04}", i);
table
.insert(key.as_str(), bytes.as_slice())
.map_err(|e| format!("insert {}: {:?}", key, e))?;
}
}
txn.commit().map_err(|e| format!("commit: {:?}", e))?;
drop(db);
Ok(count)
}
fn count_persy_records(path: &std::path::Path) -> Result<usize, String> {
use persy::Persy;
let db = Persy::open(path, persy::Config::default()).map_err(|e| format!("open: {:?}", e))?;
let segment_id = db
.solve_segment_id("beam_nodes_v1")
.map_err(|e| format!("solve_segment_id: {:?}", e))?;
let iter = db.scan(segment_id).map_err(|e| format!("scan: {:?}", e))?;
Ok(iter.count())
}
fn count_redb_records(path: &std::path::Path) -> Result<usize, String> {
use redb::{Database, ReadableDatabase, ReadableTable, TableDefinition};
const BEAM_NODES: TableDefinition<&str, &[u8]> = TableDefinition::new("beam_nodes_v1");
let db = Database::open(path).map_err(|e| format!("open: {:?}", e))?;
let txn = db
.begin_read()
.map_err(|e| format!("begin_read: {:?}", e))?;
let table = match txn.open_table(BEAM_NODES) {
Ok(t) => t,
Err(_) => return Ok(0), };
let count = table.iter().map_err(|e| format!("iter: {:?}", e))?.count();
drop(table);
drop(txn);
drop(db);
Ok(count)
}
fn cleanup(path: &PathBuf) {
let _ = std::fs::remove_file(path);
}
#[tokio::test]
async fn e2e_redb_to_persy_basic() {
let source = temp_path("redb-source", ".redb");
let target = temp_path("persy-target", ".persy");
cleanup(&target);
let written = write_redb_records(&source, 100).expect("write redb");
assert_eq!(written, 100);
let opts = MigrateOpts {
from: Backend::Redb,
to: Backend::Persy,
source_path: source.clone(),
target_path: target.clone(),
batch_size: 50,
force: false,
dry_run: false,
};
let report = migrate(&opts).expect("migrate redb→persy");
assert_eq!(report.records_migrated, 100);
assert_eq!(report.source_count, 100);
drop(opts);
let target_count = count_persy_records(&target).expect("count persy");
assert_eq!(target_count, 100, "Persy target should have 100 records");
cleanup(&source);
cleanup(&target);
}
#[tokio::test]
async fn e2e_persy_to_redb_basic() {
let source = temp_path("persy-source", ".persy");
let target = temp_path("redb-target", ".redb");
cleanup(&target);
write_redb_records(&source, 50).expect("write initial redb");
let intermediate = temp_path("persy-intermediate", ".persy");
cleanup(&intermediate);
let _ = migrate(&MigrateOpts {
from: Backend::Redb,
to: Backend::Persy,
source_path: source.clone(),
target_path: intermediate.clone(),
batch_size: 25,
force: true,
dry_run: false,
})
.expect("redb→persy intermediate");
let report = migrate(&MigrateOpts {
from: Backend::Persy,
to: Backend::Redb,
source_path: intermediate.clone(),
target_path: target.clone(),
batch_size: 25,
force: false,
dry_run: false,
})
.expect("persy→redb");
assert_eq!(report.records_migrated, 50);
let target_count = count_redb_records(&target).expect("count redb target");
assert_eq!(target_count, 50);
cleanup(&source);
cleanup(&intermediate);
cleanup(&target);
}
#[tokio::test]
async fn e2e_migration_preserves_children() {
let source = temp_path("redb-nested-source", ".redb");
let target = temp_path("persy-nested-target", ".persy");
cleanup(&target);
use redb::{Database, TableDefinition};
const BEAM_NODES: TableDefinition<&str, &[u8]> = TableDefinition::new("beam_nodes_v1");
let db = Database::create(&source).expect("create");
let txn = db.begin_write().expect("begin_write");
{
let mut table = txn.open_table(BEAM_NODES).expect("open_table");
let mut root_children: BTreeMap<String, NodeData> = BTreeMap::new();
root_children.insert(
"level1_a".to_string(),
NodeData {
value: Value::Null,
updated_at: 1_700_000_000.0,
},
);
let bytes = postcard::to_allocvec(&root_children).expect("postcard root");
table.insert("root", bytes.as_slice()).expect("insert root");
let mut l1_children: BTreeMap<String, NodeData> = BTreeMap::new();
l1_children.insert(
"level2_a".to_string(),
NodeData {
value: Value::Null,
updated_at: 1_700_000_001.0,
},
);
let bytes = postcard::to_allocvec(&l1_children).expect("postcard l1");
table
.insert("level1_a", bytes.as_slice())
.expect("insert l1");
let mut l2_children: BTreeMap<String, NodeData> = BTreeMap::new();
l2_children.insert(
"leaf1".to_string(),
NodeData {
value: Value::Text("deep value 1".into()),
updated_at: 1_700_000_002.0,
},
);
l2_children.insert(
"leaf2".to_string(),
NodeData {
value: Value::Text("deep value 2".into()),
updated_at: 1_700_000_003.0,
},
);
let bytes = postcard::to_allocvec(&l2_children).expect("postcard l2");
table
.insert("level2_a", bytes.as_slice())
.expect("insert l2");
}
txn.commit().expect("commit");
drop(db);
let report = migrate(&MigrateOpts {
from: Backend::Redb,
to: Backend::Persy,
source_path: source.clone(),
target_path: target.clone(),
batch_size: 100,
force: false,
dry_run: false,
})
.expect("migrate nested");
assert_eq!(report.records_migrated, 3, "3 top-level records");
let target_count = count_persy_records(&target).expect("count");
assert_eq!(target_count, 3);
cleanup(&source);
cleanup(&target);
}
#[tokio::test]
async fn e2e_migration_empty_dataset() {
let source = temp_path("redb-empty-source", ".redb");
let target = temp_path("persy-empty-target", ".persy");
cleanup(&target);
use redb::Database;
let db = Database::create(&source).expect("create empty");
drop(db);
let report = migrate(&MigrateOpts {
from: Backend::Redb,
to: Backend::Persy,
source_path: source.clone(),
target_path: target.clone(),
batch_size: 100,
force: false,
dry_run: false,
})
.expect("migrate empty");
assert_eq!(report.records_migrated, 0);
assert_eq!(report.source_count, 0);
cleanup(&source);
cleanup(&target);
}
#[tokio::test]
async fn e2e_migration_dry_run_no_write() {
let source = temp_path("redb-dry-source", ".redb");
let target = temp_path("persy-dry-target", ".persy");
cleanup(&target);
write_redb_records(&source, 10).expect("write");
let report = migrate(&MigrateOpts {
from: Backend::Redb,
to: Backend::Persy,
source_path: source.clone(),
target_path: target.clone(),
batch_size: 100,
force: false,
dry_run: true,
})
.expect("dry run migrate");
assert_eq!(report.records_migrated, 10);
assert!(report.dry_run);
assert!(!target.exists(), "dry run should not create target file");
cleanup(&source);
}
#[tokio::test]
async fn diag_persy_count_basic() {
use redb::{Database, TableDefinition};
const BEAM_NODES: TableDefinition<&str, &[u8]> = TableDefinition::new("beam_nodes_v1");
let source = temp_path("diag-source", ".redb");
let target = temp_path("diag-target", ".persy");
cleanup(&source);
cleanup(&target);
let db = Database::create(&source).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(BEAM_NODES).unwrap();
for i in 0..100 {
let mut children: BTreeMap<String, NodeData> = BTreeMap::new();
children.insert(
format!("leaf_{:04}", i),
NodeData {
value: Value::Text(format!("test-{}", i)),
updated_at: 1_700_000_000.0 + i as f64,
},
);
let bytes = postcard::to_allocvec(&children).unwrap();
let key = format!("node_{:04}", i);
table.insert(key.as_str(), bytes.as_slice()).unwrap();
}
}
txn.commit().unwrap();
drop(db);
let report = migrate(&MigrateOpts {
from: Backend::Redb,
to: Backend::Persy,
source_path: source.clone(),
target_path: target.clone(),
batch_size: 50,
force: false,
dry_run: false,
})
.expect("migrate");
eprintln!(
"DIAG: migration reports records_migrated={}",
report.records_migrated
);
eprintln!("DIAG: source_count={}", report.source_count);
let target_count = count_persy_records(&target).expect("count persy");
assert_eq!(
target_count, 100,
"all 100 records must persist across batches (post-fix checkpoint)"
);
assert_eq!(
report.records_migrated, 100,
"migration report must match actual persisted count"
);
cleanup(&source);
cleanup(&target);
}