#[cfg(feature = "persist")]
use crdt_lite::persist::{PersistedCRDT, PersistConfig};
#[cfg(feature = "persist")]
use std::path::PathBuf;
#[cfg(feature = "persist")]
use std::sync::{Arc, Mutex};
#[cfg(feature = "persist")]
fn main() -> Result<(), Box<dyn std::error::Error>> {
println!("=== CRDT with R2 Backup Hooks ===\n");
let uploaded_snapshots = Arc::new(Mutex::new(Vec::new()));
let uploaded_wal_segments = Arc::new(Mutex::new(Vec::new()));
let data_dir = PathBuf::from("./r2_backup_example_data");
let config = PersistConfig {
snapshot_threshold: 3, auto_cleanup_snapshots: None, max_batch_size: Some(10000), ..Default::default() };
let mut pcrdt = PersistedCRDT::<String, String, String>::open(
data_dir.clone(),
1, config,
)?;
let snapshots_clone = uploaded_snapshots.clone();
pcrdt.add_snapshot_hook(Box::new(move |snapshot_path: &std::path::Path, db_version: u64| {
println!("๐ธ Snapshot created: {:?} (db_version: {})", snapshot_path.file_name().unwrap(), db_version);
println!(" โโ Uploading to R2...");
snapshots_clone.lock().unwrap().push(snapshot_path.to_path_buf());
println!(" โโ โ Uploaded to R2");
}));
let wal_clone = uploaded_wal_segments.clone();
pcrdt.add_wal_segment_hook(Box::new(move |segment_path: &std::path::Path| {
println!("๐ WAL segment sealed: {:?}", segment_path.file_name().unwrap());
println!(" โโ Uploading to R2...");
wal_clone.lock().unwrap().push(segment_path.to_path_buf());
println!(" โโ โ Uploaded to R2");
}));
println!("โ CRDT opened with backup hooks\n");
println!("Inserting records...\n");
for i in 1..=10 {
pcrdt.insert_or_update(
&format!("user{}", i),
[
("name".to_string(), format!("User {}", i)),
("email".to_string(), format!("user{}@example.com", i)),
]
.into_iter(),
)?;
println!(" โ Inserted user{} ({} changes since snapshot)",
i, pcrdt.changes_since_snapshot());
std::thread::sleep(std::time::Duration::from_millis(100));
}
println!("\nโ
All operations complete!\n");
let snapshots = uploaded_snapshots.lock().unwrap();
let wal_segments = uploaded_wal_segments.lock().unwrap();
println!("๐ฆ Backup Summary:");
println!(" Snapshots uploaded: {}", snapshots.len());
for (i, path) in snapshots.iter().enumerate() {
println!(" {}. {:?}", i + 1, path.file_name().unwrap());
}
println!("\n WAL segments uploaded: {}", wal_segments.len());
for (i, path) in wal_segments.iter().enumerate() {
println!(" {}. {:?}", i + 1, path.file_name().unwrap());
}
println!("\n๐งน Cleanup demo:");
println!(" Before cleanup:");
let snapshot_count_before = std::fs::read_dir(&data_dir)?
.filter_map(|e| e.ok())
.filter(|e| e.file_name().to_string_lossy().starts_with("snapshot_"))
.count();
println!(" Local snapshots: {}", snapshot_count_before);
pcrdt.cleanup_old_snapshots(2, false)?;
let snapshot_count_after = std::fs::read_dir(&data_dir)?
.filter_map(|e| e.ok())
.filter(|e| e.file_name().to_string_lossy().starts_with("snapshot_"))
.count();
println!(" After cleanup (keep 2):");
println!(" Local snapshots: {}", snapshot_count_after);
println!(" (Old snapshots safely backed up in R2)");
println!("\n๐ก Recovery flow:");
println!(" 1. Container starts โ Download latest snapshot from R2");
println!(" 2. Download WAL segments since snapshot");
println!(" 3. PersistedCRDT::open() recovers from local files");
println!(" 4. Continue operations with hooks uploading new changes");
println!("\nโจ Example complete! Data in: {:?}", data_dir);
Ok(())
}
#[cfg(not(feature = "persist"))]
fn main() {
eprintln!("This example requires the 'persist' feature.");
eprintln!("Run with: cargo run --example r2_backup_example --features persist");
std::process::exit(1);
}