mod build;
mod recover;
mod util;
pub use build::build_snapshot;
pub use recover::recover_snapshot;
pub use util::get_current_snapshot;
#[cfg(test)]
mod tests {
use super::util;
use super::*;
use openraft::Membership;
use openraft::StoredMembership;
use rocksdb::DB;
use std::path::PathBuf;
use std::sync::Arc;
use tempfile::tempdir;
use tokio::time::{Duration, sleep};
use crate::raft::store::keys::SM_DATA_FAMILY;
#[tokio::test]
async fn test_build_and_recover_snapshot_with_data() {
let source_db = create_test_db_with_sample_data();
let snapshot_dir = tempdir().unwrap();
let snapshot_dir_path = snapshot_dir.path().to_path_buf();
let cf_handle = source_db
.cf_handle(SM_DATA_FAMILY)
.expect("column family not found");
source_db
.put_cf(&cf_handle, b"user:1", b"John Doe")
.unwrap();
source_db
.put_cf(&cf_handle, b"user:2", b"Jane Smith")
.unwrap();
source_db
.put_cf(&cf_handle, b"user:3", b"Bob Johnson")
.unwrap();
source_db
.put_cf(&cf_handle, b"config:theme", b"dark")
.unwrap();
let last_applied_log_id = None;
let last_membership = StoredMembership::new(None, Membership::default());
let snapshot = build_snapshot(
&source_db,
&snapshot_dir_path,
last_applied_log_id,
last_membership,
)
.await
.expect("Failed to build snapshot");
assert!(!snapshot.meta.snapshot_id.is_empty());
println!("Snapshot created with ID: {}", snapshot.meta.snapshot_id);
let target_db = create_empty_test_db();
let cf_handle_target = target_db
.cf_handle(SM_DATA_FAMILY)
.expect("column family not found");
let mut count = 0u64;
let iter = target_db.iterator_cf(&cf_handle_target, rocksdb::IteratorMode::Start);
for _ in iter {
count += 1;
}
assert_eq!(count, 0, "Target database should be empty initially");
recover_snapshot(&target_db, snapshot)
.await
.expect("Failed to recover snapshot");
sleep(Duration::from_millis(100)).await;
let mut recovered_count = 0u64;
let iter = target_db.iterator_cf(&cf_handle_target, rocksdb::IteratorMode::Start);
for _ in iter {
recovered_count += 1;
}
assert_eq!(
recovered_count, 4,
"Expected 4 entries after recovery, got {}",
recovered_count
);
let value1 = target_db
.get_cf(&cf_handle_target, b"user:1")
.unwrap()
.expect("user:1 should exist after recovery");
assert_eq!(
String::from_utf8(value1).unwrap(),
"John Doe",
"user:1 value mismatch"
);
let value2 = target_db
.get_cf(&cf_handle_target, b"user:2")
.unwrap()
.expect("user:2 should exist after recovery");
assert_eq!(
String::from_utf8(value2).unwrap(),
"Jane Smith",
"user:2 value mismatch"
);
let value3 = target_db
.get_cf(&cf_handle_target, b"user:3")
.unwrap()
.expect("user:3 should exist after recovery");
assert_eq!(
String::from_utf8(value3).unwrap(),
"Bob Johnson",
"user:3 value mismatch"
);
let value4 = target_db
.get_cf(&cf_handle_target, b"config:theme")
.unwrap()
.expect("config:theme should exist after recovery");
assert_eq!(
String::from_utf8(value4).unwrap(),
"dark",
"config:theme value mismatch"
);
println!("Integration test passed: Build and recover snapshot with data");
}
#[tokio::test]
async fn test_build_and_recover_empty_snapshot() {
let source_db = create_test_db_with_sample_data();
let snapshot_dir = tempdir().unwrap();
let snapshot_dir_path = snapshot_dir.path().to_path_buf();
let last_applied_log_id = None;
let last_membership = StoredMembership::new(None, Membership::default());
let snapshot = build_snapshot(
&source_db,
&snapshot_dir_path,
last_applied_log_id,
last_membership,
)
.await
.expect("Failed to build empty snapshot");
println!(
"Empty snapshot created with ID: {}",
snapshot.meta.snapshot_id
);
let target_db = create_empty_test_db();
let cf_handle_target = target_db
.cf_handle(SM_DATA_FAMILY)
.expect("column family not found");
recover_snapshot(&target_db, snapshot)
.await
.expect("Failed to recover empty snapshot");
sleep(Duration::from_millis(100)).await;
let mut count = 0u64;
let iter = target_db.iterator_cf(&cf_handle_target, rocksdb::IteratorMode::Start);
for _ in iter {
count += 1;
}
assert_eq!(
count, 0,
"Database should be empty after recovering empty snapshot"
);
println!("Integration test passed: Build and recover empty snapshot");
}
#[tokio::test]
async fn test_build_and_recover_large_snapshot() {
let source_db = create_test_db_with_sample_data();
let snapshot_dir = tempdir().unwrap();
let snapshot_dir_path = snapshot_dir.path().to_path_buf();
let cf_handle = source_db
.cf_handle(SM_DATA_FAMILY)
.expect("column family not found");
let data_count = 1000;
for i in 0..data_count {
let key = format!("entry:{}", i);
let value = format!("data-{}-value", i);
source_db
.put_cf(&cf_handle, key.as_bytes(), value.as_bytes())
.unwrap();
}
println!("Inserted {} entries into source database", data_count);
let last_applied_log_id = None;
let last_membership = StoredMembership::new(None, Membership::default());
let snapshot = build_snapshot(
&source_db,
&snapshot_dir_path,
last_applied_log_id,
last_membership,
)
.await
.expect("Failed to build large snapshot");
println!(
"Large snapshot created with ID: {}",
snapshot.meta.snapshot_id
);
let target_db = create_empty_test_db();
let cf_handle_target = target_db
.cf_handle(SM_DATA_FAMILY)
.expect("column family not found");
recover_snapshot(&target_db, snapshot)
.await
.expect("Failed to recover large snapshot");
sleep(Duration::from_millis(200)).await;
let mut recovered_count = 0u64;
let iter = target_db.iterator_cf(&cf_handle_target, rocksdb::IteratorMode::Start);
for _ in iter {
recovered_count += 1;
}
assert_eq!(
recovered_count, data_count as u64,
"Expected {} entries after recovery, got {}",
data_count, recovered_count
);
let sample_keys = vec![0, 100, 500, 999];
for &key_num in &sample_keys {
let key = format!("entry:{}", key_num);
let expected_value = format!("data-{}-value", key_num);
let actual_value = target_db
.get_cf(&cf_handle_target, key.as_bytes())
.unwrap()
.expect(format!("Entry {} should exist", key).as_str());
assert_eq!(
String::from_utf8(actual_value).unwrap(),
expected_value,
"Value mismatch for {}",
key
);
}
println!(
"Integration test passed: Build and recover large snapshot with {} entries",
data_count
);
}
#[tokio::test]
async fn test_build_and_recover_binary_data() {
let source_db = create_test_db_with_sample_data();
let snapshot_dir = tempdir().unwrap();
let snapshot_dir_path = snapshot_dir.path().to_path_buf();
let cf_handle = source_db
.cf_handle(SM_DATA_FAMILY)
.expect("column family not found");
let binary_data = vec![
(b"bin:key1".to_vec(), vec![0x00, 0x01, 0x02, 0x03, 0xFF]),
(b"bin:key2".to_vec(), vec![0xDE, 0xAD, 0xBE, 0xEF]),
(
b"bin:key3".to_vec(),
vec![0x80, 0x90, 0xA0, 0xB0, 0xC0, 0xD0, 0xE0, 0xF0],
),
];
for (key, value) in &binary_data {
source_db
.put_cf(&cf_handle, key.as_slice(), value.as_slice())
.unwrap();
}
let last_applied_log_id = None;
let last_membership = StoredMembership::new(None, Membership::default());
let snapshot = build_snapshot(
&source_db,
&snapshot_dir_path,
last_applied_log_id,
last_membership,
)
.await
.expect("Failed to build snapshot with binary data");
let target_db = create_empty_test_db();
let cf_handle_target = target_db
.cf_handle(SM_DATA_FAMILY)
.expect("column family not found");
recover_snapshot(&target_db, snapshot)
.await
.expect("Failed to recover snapshot with binary data");
sleep(Duration::from_millis(100)).await;
for (key, expected_value) in &binary_data {
let actual_value = target_db
.get_cf(&cf_handle_target, key.as_slice())
.unwrap()
.expect(format!("Binary entry should exist: {:?}", key).as_str());
assert_eq!(
actual_value, *expected_value,
"Binary data mismatch for key {:?}",
key
);
}
println!("Integration test passed: Build and recover snapshot with binary data");
}
#[tokio::test]
async fn test_build_and_recover_snapshot_with_metadata() {
let source_db = create_test_db_with_sample_data();
let snapshot_dir = tempdir().unwrap();
let snapshot_dir_path = snapshot_dir.path().to_path_buf();
let last_applied_log_id = None; let last_membership = StoredMembership::new(None, Membership::default());
let snapshot = build_snapshot(
&source_db,
&snapshot_dir_path,
last_applied_log_id,
last_membership,
)
.await
.expect("Failed to build snapshot");
let snapshot_id = snapshot.meta.snapshot_id.clone();
println!("Snapshot created with ID: {}", snapshot_id);
let current_snapshot = get_current_snapshot(&snapshot_dir_path)
.await
.expect("Failed to get current snapshot");
assert!(
current_snapshot.is_some(),
"Current snapshot should exist after building"
);
let retrieved_snapshot = current_snapshot.unwrap();
assert_eq!(
retrieved_snapshot.meta.snapshot_id, snapshot_id,
"Retrieved snapshot ID should match built snapshot ID"
);
let snapshot_id_dir = util::snapshot_id_dir(&snapshot_dir_path, &snapshot_id);
let meta_file = util::snapshot_meta_file(&snapshot_id_dir);
assert!(
PathBuf::from(&meta_file).exists(),
"Metadata file should exist"
);
let data_file = util::snapshot_data_file(&snapshot_id_dir);
assert!(PathBuf::from(&data_file).exists(), "Data file should exist");
println!("Integration test passed: Build and recover snapshot with metadata");
}
fn create_test_db_with_sample_data() -> Arc<DB> {
let temp_dir = tempdir().unwrap();
let mut opts = rocksdb::Options::default();
opts.create_if_missing(true);
opts.create_missing_column_families(true);
let db = DB::open_cf(&opts, temp_dir.path(), vec![SM_DATA_FAMILY]).unwrap();
std::mem::forget(temp_dir);
Arc::new(db)
}
fn create_empty_test_db() -> Arc<DB> {
let temp_dir = tempdir().unwrap();
let mut opts = rocksdb::Options::default();
opts.create_if_missing(true);
opts.create_missing_column_families(true);
let db = DB::open_cf(&opts, temp_dir.path(), vec![SM_DATA_FAMILY]).unwrap();
std::mem::forget(temp_dir);
Arc::new(db)
}
}