From e13ba7f4c752e7f2e98cceb2a7db1d1d5a461c9f Mon Sep 17 00:00:00 2001 From: Lewis Date: Wed, 10 Jun 2026 10:04:23 +0300 Subject: [PATCH] ripple: anti-entropy gossip sync Lewis: May this revision serve well! --- crates/tranquil-ripple/src/gossip.rs | 26 ++++++++++++++++++++++++++ 1 file changed, 26 insertions(+) diff --git a/crates/tranquil-ripple/src/gossip.rs b/crates/tranquil-ripple/src/gossip.rs index 2149051..548f39b 100644 --- a/crates/tranquil-ripple/src/gossip.rs +++ b/crates/tranquil-ripple/src/gossip.rs @@ -4,6 +4,7 @@ use crate::crdt::lww_map::LwwDelta; use crate::metrics; use crate::transport::{ChannelTag, IncomingFrame, Transport}; use foca::{Config, Foca, Notification, Runtime, Timer}; +use rand::Rng; use rand::SeedableRng; use rand::rngs::StdRng; use std::collections::HashSet; @@ -187,6 +188,7 @@ impl GossipEngine { let (timer_tx, mut timer_rx) = mpsc::channel::<(Timer, Duration)>(256); const WATERMARK_STALE_SECS: u64 = 30; + const ANTI_ENTROPY_SECS: u64 = 30; tokio::spawn(async move { let mut runtime = BufferedRuntime::new(); @@ -216,6 +218,11 @@ impl GossipEngine { let mut gossip_tick = tokio::time::interval(Duration::from_millis(gossip_interval_ms)); gossip_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + let mut anti_entropy_tick = + tokio::time::interval(Duration::from_secs(ANTI_ENTROPY_SECS)); + anti_entropy_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + let mut ae_rng = StdRng::from_os_rng(); + loop { tokio::select! { _ = shutdown.cancelled() => { @@ -340,6 +347,25 @@ impl GossipEngine { "gossip health check" ); } + _ = anti_entropy_tick.tick() => { + let peers: Vec = members.active_peers().collect(); + if !peers.is_empty() { + let snapshot = store.peek_full_state(); + if !snapshot.is_empty() { + let peer = peers[ae_rng.random_range(0..peers.len())]; + chunk_and_serialize(&snapshot).into_iter().for_each(|chunk| { + let t = transport.clone(); + let c = shutdown.clone(); + tokio::spawn(async move { + tokio::select! { + _ = c.cancelled() => {} + _ = t.send(peer, ChannelTag::CrdtSync, &chunk) => {} + } + }); + }); + } + } + } } } })