use std::net::{IpAddr, SocketAddr};
use std::sync::Arc;
use std::time::{Duration, Instant};
use crate::clock::{ManualClock, NodeId};
use crate::entry::Entry;
use crate::replica::Replica;
use crate::replicated_map::Config;
use crate::transport::InMemoryNetwork;
#[tokio::test]
async fn insert_broadcasts_immediately_without_a_reconciliation_round() {
let net = InMemoryNetwork::new();
let port = 5006u16;
let a_ip: IpAddr = "127.0.4.10".parse().unwrap();
let b_ip: IpAddr = "127.0.4.11".parse().unwrap();
let cfg = |ip: IpAddr| {
Config::default()
.with_listen_addr(ip)
.with_port(port)
.with_reconcile_interval(Duration::from_secs(3600))
.with_insecure_no_key()
};
let a: Replica<u32, u32> = Replica::new_with_transport(
cfg(a_ip),
Arc::new(net.bind(SocketAddr::new(a_ip, port))),
Arc::new(ManualClock::new(NodeId::new(1))),
);
let b: Replica<u32, u32> = Replica::new_with_transport(
cfg(b_ip),
Arc::new(net.bind(SocketAddr::new(b_ip, port))),
Arc::new(ManualClock::new(NodeId::new(2))),
);
let ta = tokio::spawn(a.clone().run());
let tb = tokio::spawn(b.clone().run());
tokio::time::sleep(Duration::from_millis(100)).await;
a.peers.write().insert(b_ip, Instant::now());
a.insert(7, Entry::present(a.clock_now(), 42));
let deadline = Instant::now() + Duration::from_secs(5);
let mut received = None;
while Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(20)).await;
if let Some(entry) = b.map.read().get(&7) {
received = entry.value().copied();
if received.is_some() {
break;
}
}
}
ta.abort();
tb.abort();
assert_eq!(
received,
Some(42),
"B never received A's insert via the immediate broadcast path — both engines' \
reconcile_interval is far beyond this test's deadline, and B was never told about A, \
so only insert's own broadcast call could have delivered it"
);
}