use std::time::Duration;
use edgestore::{EdgestoreConfig, EdgestoreError};
use edgestore_repl::ReplicatedEngine;
use tempfile::TempDir;
fn main() -> Result<(), Box<dyn std::error::Error>> {
let primary_dir = TempDir::new()?;
let replica_dir = TempDir::new()?;
println!("=== ReplicatedEngine demo ===\n");
let primary = ReplicatedEngine::open_primary(
EdgestoreConfig::new(primary_dir.path()),
"127.0.0.1:0", )?;
let port = primary.bound_port().unwrap();
let primary_url = format!("http://127.0.0.1:{}", port);
println!("Primary listening on {}", primary_url);
{
let engine = primary.engine().lock().unwrap();
drop(engine);
}
{
let mut engine = primary.engine().lock().unwrap();
engine.put(b"catalog", b"item1", b"Widget A")?;
engine.put(b"catalog", b"item2", b"Widget B")?;
engine.put(b"catalog", b"item3", b"Gadget X")?;
engine.flush_to_segments()?;
println!("Primary: wrote 3 items and flushed to segment.");
}
let mut replica_cfg = EdgestoreConfig::new(replica_dir.path());
replica_cfg.readonly = true; let replica =
ReplicatedEngine::open_replica(EdgestoreConfig::new(replica_dir.path()), &primary_url)?;
println!("Replica connected to {}", primary_url);
match replica
.engine()
.lock()
.unwrap()
.put(b"catalog", b"rogue", b"write")
{
Err(EdgestoreError::ReadOnly) => println!("Replica: write correctly rejected (ReadOnly)."),
other => panic!("Expected ReadOnly, got {:?}", other),
}
println!("Waiting for anti-entropy pull cycle...");
std::thread::sleep(Duration::from_secs(35));
let items = replica
.engine()
.lock()
.unwrap()
.range(b"catalog", b"", b"\xff")?;
println!(
"Replica: {} items visible after sync (expected 3).",
items.len()
);
for (k, v) in &items {
println!(
" {} = {}",
String::from_utf8_lossy(k),
String::from_utf8_lossy(v)
);
}
println!("\nDone.");
Ok(())
}