use std::path::PathBuf;
use web_time::Duration;
use serde::{Deserialize, Serialize};
use crate::types::Children;
#[derive(Serialize, Deserialize, Default, Debug, Clone, PartialEq)]
pub(crate) struct NodeRecord {
pub(crate) node_id: String,
pub(crate) children: Children,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum Backend {
Redb,
Persy,
}
impl Backend {
pub fn as_str(&self) -> &'static str {
match self {
Backend::Redb => "redb",
Backend::Persy => "persy",
}
}
}
impl std::fmt::Display for Backend {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone)]
pub struct MigrateOpts {
pub from: Backend,
pub to: Backend,
pub source_path: PathBuf,
pub target_path: PathBuf,
pub batch_size: usize,
pub force: bool,
pub dry_run: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MigrationReport {
pub records_migrated: usize,
pub source_count: usize,
pub target_count_after: usize,
pub elapsed: Duration,
pub dry_run: bool,
}
#[derive(thiserror::Error, Debug)]
pub enum MigrateError {
#[error("redb error at {path}: {source}")]
Redb {
path: PathBuf,
#[source]
source: redb::Error,
},
#[error("redb transaction error at {path}: {source}")]
RedbTx {
path: PathBuf,
#[source]
source: redb::TransactionError,
},
#[error("redb table error at {path}: {source}")]
RedbTable {
path: PathBuf,
#[source]
source: redb::TableError,
},
#[error("redb commit error at {path}: {source}")]
RedbCommit {
path: PathBuf,
#[source]
source: redb::CommitError,
},
#[error("persy error: {0}")]
Persy(String),
#[error("postcard error: {0}")]
Postcard(#[from] postcard::Error),
#[error("io error: {0}")]
Io(#[from] std::io::Error),
#[error("json error: {0}")]
Json(#[from] serde_json::Error),
#[error("target already exists at {0} (use --force to overwrite)")]
TargetExists(PathBuf),
#[error("unsupported migration: {from} -> {to} (use --from and --to with different backends)")]
Unsupported { from: Backend, to: Backend },
#[error("invalid backend string: {0} (expected 'redb' or 'persy')")]
InvalidBackend(String),
}
impl MigrateError {
pub fn parse_backend(s: &str) -> Result<Backend, MigrateError> {
match s.to_lowercase().as_str() {
"redb" => Ok(Backend::Redb),
"persy" => Ok(Backend::Persy),
_ => Err(MigrateError::InvalidBackend(s.to_string())),
}
}
}
pub fn redb_to_persy_payload(key: &str, value: &[u8]) -> Result<Vec<u8>, MigrateError> {
let children: Children = postcard::from_bytes(value)?;
let record = NodeRecord {
node_id: key.to_string(),
children,
};
Ok(postcard::to_allocvec(&record)?)
}
pub fn persy_to_redb_record(payload: &[u8]) -> Result<(String, Vec<u8>), MigrateError> {
let record: NodeRecord = postcard::from_bytes(payload)?;
let children_bytes = postcard::to_allocvec(&record.children)?;
Ok((record.node_id, children_bytes))
}
#[cfg(feature = "persy")]
pub(crate) mod io {
use super::*;
use web_time::Instant;
use redb::{Database, ReadableDatabase, ReadableTable, TableDefinition};
use crate::adapters::persy_storage::BEAM_NODES as PERSY_BEAM_NODES;
const REDB_BEAM_NODES: TableDefinition<&str, &[u8]> = TableDefinition::new("beam_nodes_v1");
pub fn migrate(opts: &MigrateOpts) -> Result<MigrationReport, MigrateError> {
if opts.from == opts.to {
return Err(MigrateError::Unsupported {
from: opts.from,
to: opts.to,
});
}
if !opts.dry_run && opts.target_path.exists() && !opts.force {
return Err(MigrateError::TargetExists(opts.target_path.clone()));
}
let start = Instant::now();
let (source_count, migrated) = match (opts.from, opts.to) {
(Backend::Redb, Backend::Persy) => migrate_redb_to_persy(opts)?,
(Backend::Persy, Backend::Redb) => migrate_persy_to_redb(opts)?,
(Backend::Redb, Backend::Redb) | (Backend::Persy, Backend::Persy) => {
unreachable!("from == to caught above")
}
};
Ok(MigrationReport {
records_migrated: if opts.dry_run { source_count } else { migrated },
source_count,
target_count_after: if opts.dry_run { 0 } else { migrated },
elapsed: start.elapsed(),
dry_run: opts.dry_run,
})
}
fn migrate_redb_to_persy(opts: &MigrateOpts) -> Result<(usize, usize), MigrateError> {
let src_db = Database::open(&opts.source_path).map_err(|source| MigrateError::Redb {
path: opts.source_path.clone(),
source: source.into(),
})?;
let src_tx = src_db.begin_read().map_err(|source| MigrateError::RedbTx {
path: opts.source_path.clone(),
source,
})?;
let src_table = match src_tx.open_table(REDB_BEAM_NODES) {
Ok(t) => t,
Err(e) => {
use redb::TableError;
if matches!(e, TableError::TableDoesNotExist { .. }) {
return Ok((0, 0));
}
return Err(MigrateError::RedbTable {
path: opts.source_path.clone(),
source: e,
});
}
};
let mut payloads: Vec<Vec<u8>> = Vec::new();
{
let iter = src_table.iter().map_err(|source| MigrateError::RedbTable {
path: opts.source_path.clone(),
source: redb::TableError::Storage(source),
})?;
for entry in iter {
let (key_guard, value_guard) = entry.map_err(|source| MigrateError::RedbTable {
path: opts.source_path.clone(),
source: redb::TableError::Storage(source),
})?;
let key = key_guard.value();
let value = value_guard.value();
payloads.push(redb_to_persy_payload(key, value)?);
}
}
let source_count = payloads.len();
drop(src_table);
drop(src_tx);
drop(src_db);
if opts.dry_run {
return Ok((source_count, 0));
}
let target_db = persy::Persy::open_or_create_with(
opts.target_path.to_string_lossy().as_ref(),
persy::Config::new(),
|persy_db| -> Result<(), Box<dyn std::error::Error>> {
let mut create_tx = persy_db
.begin()
.map_err(|e| Box::new(e) as Box<dyn std::error::Error>)?;
create_tx
.create_segment(PERSY_BEAM_NODES)
.map_err(|e| Box::new(e) as Box<dyn std::error::Error>)?;
create_tx
.prepare()
.map_err(Box::<dyn std::error::Error>::from)?
.commit()
.map_err(Box::<dyn std::error::Error>::from)?;
Ok(())
},
)
.map_err(|source| {
MigrateError::Persy(format!(
"{}: {}\nhelp: ensure the target directory is writable",
opts.target_path.display(),
source
))
})?;
let target_seg = target_db.solve_segment_id(PERSY_BEAM_NODES).map_err(|e| {
MigrateError::Persy(format!(
"{}: solve_segment_id: {}",
opts.target_path.display(),
e
))
})?;
let mut tx = target_db.begin().map_err(|source| {
MigrateError::Persy(format!("{}: begin: {}", opts.target_path.display(), source))
})?;
for payload in &payloads {
tx.insert(target_seg, payload.as_slice())
.map_err(|source| {
MigrateError::Persy(format!(
"{}: insert: {}",
opts.target_path.display(),
source
))
})?;
}
tx.prepare()
.map_err(|source| {
MigrateError::Persy(format!(
"{}: prepare: {}",
opts.target_path.display(),
source
))
})?
.commit()
.map_err(|source| {
MigrateError::Persy(format!(
"{}: commit: {}",
opts.target_path.display(),
source
))
})?;
drop(target_db);
Ok((source_count, payloads.len()))
}
fn migrate_persy_to_redb(opts: &MigrateOpts) -> Result<(usize, usize), MigrateError> {
let src_db = persy::Persy::open(
opts.source_path.to_string_lossy().as_ref(),
persy::Config::new(),
)
.map_err(|source| {
MigrateError::Persy(format!(
"{}: {}",
opts.source_path.clone().display(),
source
))
})?;
let src_seg = src_db.solve_segment_id(PERSY_BEAM_NODES).map_err(|e| {
MigrateError::Persy(format!(
"{}: solve_segment_id failed: {}",
opts.source_path.display(),
e
))
})?;
let target_db =
Database::create(&opts.target_path).map_err(|source| MigrateError::Redb {
path: opts.target_path.clone(),
source: source.into(),
})?;
let target_tx = target_db
.begin_write()
.map_err(|source| MigrateError::RedbTx {
path: opts.target_path.clone(),
source,
})?;
let mut target_table =
target_tx
.open_table(REDB_BEAM_NODES)
.map_err(|source| MigrateError::RedbTable {
path: opts.target_path.clone(),
source,
})?;
let scan = src_db.scan(src_seg).map_err(|source| {
MigrateError::Persy(format!(
"{}: {}",
opts.source_path.clone().display(),
source
))
})?;
let mut source_count = 0;
let mut migrated = 0;
let mut batch: Vec<(String, Vec<u8>)> = Vec::with_capacity(opts.batch_size);
for entry in scan {
let (_id, bytes) = entry;
let (key, value_bytes) = persy_to_redb_record(&bytes)?;
source_count += 1;
if !opts.dry_run {
batch.push((key, value_bytes));
if batch.len() >= opts.batch_size {
for (k, v) in batch.drain(..) {
target_table
.insert(k.as_str(), v.as_slice())
.map_err(|source| MigrateError::RedbTable {
path: opts.target_path.clone(),
source: redb::TableError::Storage(source),
})?;
migrated += 1;
}
}
}
}
if !batch.is_empty() && !opts.dry_run {
for (k, v) in batch.drain(..) {
target_table
.insert(k.as_str(), v.as_slice())
.map_err(|source| MigrateError::RedbTable {
path: opts.target_path.clone(),
source: redb::TableError::Storage(source),
})?;
migrated += 1;
}
}
drop(target_table);
target_tx
.commit()
.map_err(|source| MigrateError::RedbCommit {
path: opts.target_path.clone(),
source,
})?;
Ok((source_count, migrated))
}
}
#[cfg(feature = "persy")]
pub use io::migrate;
#[cfg(test)]
mod tests {
use super::*;
use crate::types::{NodeData, Value};
use std::collections::BTreeMap;
fn make_test_children() -> Children {
let mut children = BTreeMap::new();
children.insert(
"greeting".to_string(),
NodeData {
value: Value::Text("hello".to_string()),
updated_at: 12345.0,
},
);
children.insert(
"count".to_string(),
NodeData {
value: Value::Number(42.0),
updated_at: 67890.0,
},
);
children.insert(
"flag".to_string(),
NodeData {
value: Value::Bit(true),
updated_at: 11111.0,
},
);
children
}
#[test]
fn redb_to_persy_roundtrips_children() {
let children = make_test_children();
let original_bytes = postcard::to_allocvec(&children).unwrap();
let translated = redb_to_persy_payload("test-node", &original_bytes).unwrap();
let record: NodeRecord = postcard::from_bytes(&translated).unwrap();
assert_eq!(record.node_id, "test-node");
assert_eq!(record.children, children);
}
#[test]
fn persy_to_redb_roundtrips_children() {
let children = make_test_children();
let record = NodeRecord {
node_id: "test-node".to_string(),
children: children.clone(),
};
let payload = postcard::to_allocvec(&record).unwrap();
let (key, value_bytes) = persy_to_redb_record(&payload).unwrap();
assert_eq!(key, "test-node");
let recovered: Children = postcard::from_bytes(&value_bytes).unwrap();
assert_eq!(recovered, children);
}
#[test]
fn translation_is_pure_and_deterministic() {
let children = make_test_children();
let bytes = postcard::to_allocvec(&children).unwrap();
let result1 = redb_to_persy_payload("k", &bytes).unwrap();
let result2 = redb_to_persy_payload("k", &bytes).unwrap();
assert_eq!(result1, result2, "same input must produce same output");
}
#[test]
fn empty_children_translates_cleanly() {
let empty: Children = BTreeMap::new();
let bytes = postcard::to_allocvec(&empty).unwrap();
let translated = redb_to_persy_payload("empty-node", &bytes).unwrap();
let record: NodeRecord = postcard::from_bytes(&translated).unwrap();
assert_eq!(record.node_id, "empty-node");
assert!(record.children.is_empty());
}
#[test]
fn all_value_variants_preserved() {
let mut children = BTreeMap::new();
children.insert(
"null".to_string(),
NodeData {
value: Value::Null,
updated_at: 1.0,
},
);
children.insert(
"bit".to_string(),
NodeData {
value: Value::Bit(false),
updated_at: 2.0,
},
);
children.insert(
"num".to_string(),
NodeData {
value: Value::Number(-3.15),
updated_at: 3.0,
},
);
children.insert(
"text".to_string(),
NodeData {
value: Value::Text("unicode: ☃ snowman".to_string()),
updated_at: 4.0,
},
);
children.insert(
"link".to_string(),
NodeData {
value: Value::Link("node/abc".to_string()),
updated_at: 5.0,
},
);
let bytes = postcard::to_allocvec(&children).unwrap();
let translated = redb_to_persy_payload("root", &bytes).unwrap();
let record: NodeRecord = postcard::from_bytes(&translated).unwrap();
assert_eq!(record.children, children);
if let Value::Link(ref s) = record.children.get("link").unwrap().value {
assert_eq!(s, "node/abc");
} else {
panic!("link value not preserved");
}
}
#[test]
fn backend_parse_accepts_lowercase() {
assert_eq!(MigrateError::parse_backend("redb").unwrap(), Backend::Redb);
assert_eq!(
MigrateError::parse_backend("persy").unwrap(),
Backend::Persy
);
}
#[test]
fn backend_parse_accepts_mixed_case() {
assert_eq!(MigrateError::parse_backend("Redb").unwrap(), Backend::Redb);
assert_eq!(
MigrateError::parse_backend("PERSY").unwrap(),
Backend::Persy
);
}
#[test]
fn backend_parse_rejects_unknown() {
assert!(matches!(
MigrateError::parse_backend("sqlite"),
Err(MigrateError::InvalidBackend(_))
));
}
#[test]
fn backend_as_str_roundtrips() {
assert_eq!(Backend::Redb.as_str(), "redb");
assert_eq!(Backend::Persy.as_str(), "persy");
}
#[test]
fn unsorted_keys_preserved_after_roundtrip() {
let mut children = BTreeMap::new();
children.insert(
"z".to_string(),
NodeData {
value: Value::Text("last".to_string()),
updated_at: 1.0,
},
);
children.insert(
"a".to_string(),
NodeData {
value: Value::Text("first".to_string()),
updated_at: 2.0,
},
);
children.insert(
"m".to_string(),
NodeData {
value: Value::Text("middle".to_string()),
updated_at: 3.0,
},
);
let bytes = postcard::to_allocvec(&children).unwrap();
let translated = redb_to_persy_payload("k", &bytes).unwrap();
let record: NodeRecord = postcard::from_bytes(&translated).unwrap();
let keys: Vec<&String> = record.children.keys().collect();
assert_eq!(keys, vec!["a", "m", "z"]); }
}