diff --git a/src/engine/db.rs b/src/engine/db.rs index f2147d4..3455c49 100644 --- a/src/engine/db.rs +++ b/src/engine/db.rs @@ -157,6 +157,46 @@ impl EventLog { let mut f = std::fs::OpenOptions::new().create(true).append(true).open(&self.path)?; writeln!(f, "{}", serde_json::to_string(entry).unwrap_or_default()) } + + // How many entries the log currently holds. + fn len(&self) -> usize { + self.entries.len() + } + + // Rewrite the log to a minimal snapshot: one AccountRegistered per account, + // authored under our origin at fresh sequence numbers. Peers re-converge + // because AccountRegistered overwrites and cert replay is idempotent. The + // version vector resets to our origin; peers re-advertise theirs on the next + // sync. Written to a temp file and renamed, so a crash leaves the old log. + fn compact(&mut self, accounts: &HashMap) -> std::io::Result<()> { + let mut seq = self.next_seq(); + let mut snapshot = Vec::with_capacity(accounts.len()); + for account in accounts.values() { + self.lamport += 1; + snapshot.push(LogEntry { + origin: self.origin.clone(), + seq, + lamport: self.lamport, + event: Event::AccountRegistered(account.clone()), + }); + seq += 1; + } + let tmp = self.path.with_extension("compact"); + { + let mut f = std::fs::File::create(&tmp)?; + for entry in &snapshot { + writeln!(f, "{}", serde_json::to_string(entry).unwrap_or_default())?; + } + f.sync_all()?; + } + std::fs::rename(&tmp, &self.path)?; + self.versions = HashMap::new(); + if let Some(last) = seq.checked_sub(1) { + self.versions.insert(self.origin.clone(), last); + } + self.entries = snapshot; + Ok(()) + } } #[derive(Debug)] @@ -224,6 +264,20 @@ impl Db { Ok(()) } + /// Rewrite the log to one entry per account, reclaiming churn. + pub fn compact(&mut self) -> std::io::Result<()> { + let before = self.log.len(); + self.log.compact(&self.accounts)?; + tracing::info!(before, after = self.log.len(), "compacted event log"); + Ok(()) + } + + /// Whether the log has grown enough past the live account count to be worth + /// compacting. + pub fn should_compact(&self) -> bool { + self.log.len() > self.accounts.len() * 3 + 64 + } + /// Wire the log to a broadcast channel; each new entry is pushed to peers. pub fn set_outbound(&mut self, tx: broadcast::Sender) { self.log.set_outbound(tx); @@ -353,7 +407,9 @@ fn apply(accounts: &mut HashMap, event: Event) { } Event::CertAdded { account, fp } => { if let Some(a) = accounts.get_mut(&key(&account)) { - a.certfps.push(fp); + if !a.certfps.contains(&fp) { + a.certfps.push(fp); // idempotent: safe to replay over a snapshot + } } } Event::CertRemoved { account, fp } => { @@ -444,4 +500,51 @@ mod tests { assert!(db.exists("bob"), "peer's account is present locally"); assert!(db.exists("alice"), "local account still there"); } + + // Compaction collapses churn to one entry per account and survives a reload. + #[test] + fn compact_preserves_state_and_shrinks() { + let p = tmp("compact"); + let mut db = Db::open(&p, "local"); + db.scram_iterations = 4096; + db.register("alice", "pw", None).unwrap(); + db.register("bob", "pw", None).unwrap(); + for i in 0..10u32 { + let fp = format!("{i:032x}"); + db.certfp_add("alice", &fp).unwrap(); + db.certfp_del("alice", &fp).unwrap(); + } + let kept = "bb".repeat(16); + db.certfp_add("bob", &kept).unwrap(); + + let before = db.log.len(); + db.compact().unwrap(); + assert!(db.log.len() < before, "log shrank: {before} -> {}", db.log.len()); + assert_eq!(db.log.len(), 2, "one entry per account"); + + drop(db); + let db = Db::open(&p, "local"); + assert!(db.authenticate("alice", "pw").is_some(), "password survives compaction"); + assert!(db.certfps("alice").is_empty(), "churned certs are gone"); + assert_eq!(db.certfps("bob"), &[kept][..], "kept cert survives"); + } + + // A peer with an empty log converges from a node that has already compacted. + #[test] + fn compacted_node_still_converges() { + let mut a = Db::open(tmp("compA"), "A"); + a.scram_iterations = 4096; + a.register("alice", "pw", None).unwrap(); + let fp = "cc".repeat(16); + a.certfp_add("alice", &fp).unwrap(); + a.compact().unwrap(); + + let mut b = Db::open(tmp("compB"), "B"); + b.scram_iterations = 4096; + for entry in a.missing_for(&b.version_vector()) { + b.ingest(entry).unwrap(); + } + assert!(b.exists("alice"), "B converges from a compacted node"); + assert_eq!(b.certfps("alice"), &[fp][..], "certs arrive folded into the snapshot"); + } } diff --git a/src/engine/mod.rs b/src/engine/mod.rs index d9807ac..4704f54 100644 --- a/src/engine/mod.rs +++ b/src/engine/mod.rs @@ -104,6 +104,16 @@ impl Engine { self.db.ingest(entry) } + // Compact the log if it has grown past the live account count. Returns + // whether it rewrote anything. + pub fn maybe_compact(&mut self) -> std::io::Result { + if self.db.should_compact() { + self.db.compact()?; + return Ok(true); + } + Ok(false) + } + #[cfg(test)] pub(crate) fn test_register(&mut self, name: &str) { self.db.register(name, "pw", None).unwrap(); diff --git a/src/main.rs b/src/main.rs index 1f710c5..be43d46 100644 --- a/src/main.rs +++ b/src/main.rs @@ -51,6 +51,19 @@ async fn main() -> Result<()> { tokio::spawn(gossip::run(engine.clone(), gossip, cfg.peer.clone(), cfg.server.sid.clone(), gossip_tx)); } + // Periodically fold log churn into a snapshot when it grows past the accounts. + { + let engine = engine.clone(); + tokio::spawn(async move { + loop { + tokio::time::sleep(std::time::Duration::from_secs(1800)).await; + if let Err(e) = engine.lock().await.maybe_compact() { + tracing::warn!(%e, "compaction failed"); + } + } + }); + } + let addr = format!("{}:{}", cfg.uplink.host, cfg.uplink.port); tracing::info!(server = %cfg.server.name, %addr, "linking to uplink"); link::run(proto, engine, &addr).await