use std::net::IpAddr;
use std::sync::Arc;
use std::time::Duration;
use crate::ReplicatedMap;
use super::ephemeral_config;
#[tokio::test]
async fn start_reconciliation_actually_drives_a_round() {
use std::net::SocketAddr;
use crate::transport::InMemoryNetwork;
let net = InMemoryNetwork::new();
let port = 5101u16;
let a_ip: IpAddr = "127.0.3.5".parse().unwrap();
let b_ip: IpAddr = "127.0.3.6".parse().unwrap();
let cfg = |ip: IpAddr| {
ephemeral_config()
.with_listen_addr(ip)
.with_port(port)
.with_reconcile_interval(Duration::from_secs(3600))
};
let a = ReplicatedMap::<i32, i32>::new_with_transport(
cfg(a_ip),
Arc::new(net.bind(SocketAddr::new(a_ip, port))),
);
let b = ReplicatedMap::<i32, i32>::new_with_transport(
cfg(b_ip),
Arc::new(net.bind(SocketAddr::new(b_ip, port))),
);
a.insert(99, 42);
let task_a = tokio::spawn(a.clone().run());
let task_b = tokio::spawn(b.clone().run());
tokio::time::sleep(Duration::from_millis(150)).await;
a.engine
.peers
.write()
.insert(b_ip, std::time::Instant::now());
let mut converged = false;
for _ in 0..300 {
if b.get(&99).as_deref() == Some(&42) {
converged = true;
break;
}
a.start_reconciliation().await;
tokio::time::sleep(Duration::from_millis(20)).await;
}
task_a.abort();
task_b.abort();
assert!(
converged,
"explicit start_reconciliation calls never converged the peer"
);
}
#[tokio::test]
async fn peers_map_len_reflects_the_engine_peers_map() {
let store = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.unwrap();
assert_eq!(store.peers_map_len(), 0);
store
.engine
.peers
.write()
.insert("127.0.0.222".parse().unwrap(), std::time::Instant::now());
assert_eq!(store.peers_map_len(), 1);
}
#[tokio::test]
async fn tombstone_acks_len_reflects_the_engine_tombstone_acks_map() {
let store = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.unwrap();
assert_eq!(store.tombstone_acks_len(), 0);
store
.engine
.tombstone_acks
.write()
.insert(1, std::collections::HashMap::new());
assert_eq!(store.tombstone_acks_len(), 1);
}
#[tokio::test]
async fn replay_filter_len_reflects_the_engine_replay_filter() {
use std::net::SocketAddr;
use crate::transport::InMemoryNetwork;
let net = InMemoryNetwork::new();
let port = 5104u16;
let a_ip: IpAddr = "127.0.5.5".parse().unwrap();
let b_ip: IpAddr = "127.0.5.6".parse().unwrap();
let key = gossip::auth::ClusterKey::new([7u8; 32]);
let cfg = |ip: IpAddr| {
ephemeral_config()
.with_listen_addr(ip)
.with_port(port)
.with_cluster_key(key.clone())
.with_reconcile_interval(Duration::from_millis(5))
};
let a = ReplicatedMap::<i32, i32>::new_with_transport(
cfg(a_ip),
Arc::new(net.bind(SocketAddr::new(a_ip, port))),
);
let b = ReplicatedMap::<i32, i32>::new_with_transport(
cfg(b_ip),
Arc::new(net.bind(SocketAddr::new(b_ip, port))),
);
a.engine
.peers
.write()
.insert(b_ip, std::time::Instant::now());
b.engine
.peers
.write()
.insert(a_ip, std::time::Instant::now());
assert_eq!(b.replay_filter_len(), 0);
a.insert(1, 1);
let task_a = tokio::spawn(a.clone().run());
let task_b = tokio::spawn(b.clone().run());
let mut seen = false;
for _ in 0..300 {
tokio::time::sleep(Duration::from_millis(20)).await;
if b.replay_filter_len() >= 1 {
seen = true;
break;
}
}
task_a.abort();
task_b.abort();
assert!(
seen,
"receiving an authenticated datagram from A must register A in B's replay filter"
);
}
#[tokio::test]
async fn bulk_dumps_in_flight_count_reflects_a_dump_actually_in_progress() {
use std::net::SocketAddr;
use crate::transport::InMemoryNetwork;
let net = InMemoryNetwork::new();
let port = 5105u16;
let a_ip: IpAddr = "127.0.6.5".parse().unwrap();
let b_ip: IpAddr = "127.0.6.6".parse().unwrap();
let cfg = |ip: IpAddr| {
ephemeral_config()
.with_listen_addr(ip)
.with_port(port)
.with_reconcile_interval(Duration::from_millis(20))
.with_bulk_send_rate(super::super::MIN_BULK_SEND_RATE)
};
let a = ReplicatedMap::<i32, Vec<u8>>::new_with_transport(
cfg(a_ip),
Arc::new(net.bind(SocketAddr::new(a_ip, port))),
);
let b = ReplicatedMap::<i32, Vec<u8>>::new_with_transport(
cfg(b_ip),
Arc::new(net.bind(SocketAddr::new(b_ip, port))),
);
a.engine
.peers
.write()
.insert(b_ip, std::time::Instant::now());
b.engine
.peers
.write()
.insert(a_ip, std::time::Instant::now());
let payload = vec![0u8; 200];
let entries: Vec<(i32, Vec<u8>)> = (0..8000).map(|k| (k, payload.clone())).collect();
a.just_insert_bulk(&entries);
assert_eq!(a.bulk_dumps_in_flight_count(), 0);
let task_a = tokio::spawn(a.clone().run());
let task_b = tokio::spawn(b.clone().run());
let mut seen_in_flight = false;
for _ in 0..400 {
tokio::time::sleep(Duration::from_millis(20)).await;
if a.bulk_dumps_in_flight_count() >= 1 {
seen_in_flight = true;
break;
}
}
task_a.abort();
task_b.abort();
assert!(
seen_in_flight,
"sending a cold-sync bulk dump must register as in-flight"
);
}
#[tokio::test]
async fn set_remote_interval_actually_retunes_the_cross_network_cadence() {
use std::net::SocketAddr;
use ipnet::IpNet;
use crate::transport::InMemoryNetwork;
let net = InMemoryNetwork::new();
let port = 5102u16;
let net_a: IpNet = "127.1.0.0/24".parse().unwrap();
let a_ip: IpAddr = "127.1.0.5".parse().unwrap();
let b_ip: IpAddr = "127.2.0.5".parse().unwrap();
let a = ReplicatedMap::<i32, i32>::new_with_transport(
ephemeral_config()
.with_listen_addr(a_ip)
.with_port(port)
.with_net(net_a)
.with_reconcile_interval(Duration::from_millis(5)),
Arc::new(net.bind(SocketAddr::new(a_ip, port))),
);
let b = ReplicatedMap::<i32, i32>::new_with_transport(
ephemeral_config()
.with_listen_addr(b_ip)
.with_port(port)
.with_reconcile_interval(Duration::from_millis(5)),
Arc::new(net.bind(SocketAddr::new(b_ip, port))),
);
a.insert(7, 42);
let task_a = tokio::spawn(a.clone().run());
let task_b = tokio::spawn(b.clone().run());
tokio::time::sleep(Duration::from_millis(300)).await;
a.engine
.peers
.write()
.insert(b_ip, std::time::Instant::now());
a.set_remote_interval(100_000);
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
b.get(&7).is_none(),
"remote_interval=100000 must starve cross-network contact"
);
a.set_remote_interval(1);
let mut converged = false;
for _ in 0..300 {
tokio::time::sleep(Duration::from_millis(20)).await;
if b.get(&7).as_deref() == Some(&42) {
converged = true;
break;
}
}
task_a.abort();
task_b.abort();
assert!(
converged,
"retuning remote_interval down must let cross-network contact resume"
);
}
#[tokio::test]
async fn set_remote_fanout_actually_retunes_the_cross_network_sample_size() {
use std::net::SocketAddr;
use ipnet::IpNet;
use crate::transport::InMemoryNetwork;
let net = InMemoryNetwork::new();
let port = 5106u16;
let net_a: IpNet = "127.1.1.0/24".parse().unwrap();
let a_ip: IpAddr = "127.1.1.5".parse().unwrap();
let b_ip: IpAddr = "127.2.1.5".parse().unwrap();
let a = ReplicatedMap::<i32, i32>::new_with_transport(
ephemeral_config()
.with_listen_addr(a_ip)
.with_port(port)
.with_net(net_a)
.with_reconcile_interval(Duration::from_millis(5)),
Arc::new(net.bind(SocketAddr::new(a_ip, port))),
);
let b = ReplicatedMap::<i32, i32>::new_with_transport(
ephemeral_config()
.with_listen_addr(b_ip)
.with_port(port)
.with_reconcile_interval(Duration::from_millis(5)),
Arc::new(net.bind(SocketAddr::new(b_ip, port))),
);
a.insert(7, 42);
a.engine
.peers
.write()
.insert(b_ip, std::time::Instant::now());
a.set_remote_interval(1); a.set_remote_fanout(0);
let task_a = tokio::spawn(a.clone().run());
let task_b = tokio::spawn(b.clone().run());
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
b.get(&7).is_none(),
"remote_fanout=0 must starve cross-network contact"
);
a.set_remote_fanout(1);
let mut converged = false;
for _ in 0..300 {
tokio::time::sleep(Duration::from_millis(20)).await;
if b.get(&7).as_deref() == Some(&42) {
converged = true;
break;
}
}
task_a.abort();
task_b.abort();
assert!(
converged,
"retuning remote_fanout up must let cross-network contact resume"
);
}
#[tokio::test]
async fn set_reconcile_interval_actually_retunes_the_round_cadence() {
use std::net::SocketAddr;
use ipnet::IpNet;
use crate::transport::InMemoryNetwork;
let net_fabric = InMemoryNetwork::new();
let port = 5103u16;
let shared_net: IpNet = "127.0.4.0/24".parse().unwrap();
let a_ip: IpAddr = "127.0.4.5".parse().unwrap();
let b_ip: IpAddr = "127.0.4.6".parse().unwrap();
let cfg = |ip: IpAddr| {
ephemeral_config()
.with_listen_addr(ip)
.with_port(port)
.with_net(shared_net)
.with_reconcile_interval(Duration::from_secs(3600))
};
let a = ReplicatedMap::<i32, i32>::new_with_transport(
cfg(a_ip),
Arc::new(net_fabric.bind(SocketAddr::new(a_ip, port))),
);
let b = ReplicatedMap::<i32, i32>::new_with_transport(
cfg(b_ip),
Arc::new(net_fabric.bind(SocketAddr::new(b_ip, port))),
);
a.set_reconcile_interval(Duration::from_millis(5));
a.insert(7, 42);
let task_a = tokio::spawn(a.clone().run());
let task_b = tokio::spawn(b.clone().run());
tokio::time::sleep(Duration::from_millis(150)).await;
a.engine
.peers
.write()
.insert(b_ip, std::time::Instant::now());
let mut converged = false;
for _ in 0..300 {
tokio::time::sleep(Duration::from_millis(20)).await;
if b.get(&7).as_deref() == Some(&42) {
converged = true;
break;
}
}
task_a.abort();
task_b.abort();
assert!(
converged,
"retuning reconcile_interval before the loop's first wait must make it converge \
quickly, not wait out the original 3600s interval"
);
}