use std::net::IpAddr;
use std::sync::Arc;
use std::time::Duration;
use crate::clock::NodeId;
use crate::{replicated_map::Config, ReplicatedMap};
use super::ephemeral_config;
#[tokio::test]
async fn node_id_is_readable_and_matches_the_minted_stamp() {
let store =
ReplicatedMap::<i32, i32>::new(ephemeral_config().with_node_id(NodeId::new(0xABCD)))
.await
.unwrap();
assert_eq!(store.node_id(), NodeId::new(0xABCD));
store.insert(1, 1);
let stamp = store.engine.map.read().get(&1).unwrap().stamp;
assert_eq!(stamp.node_id(), store.node_id());
}
#[tokio::test]
async fn stores_converge_over_an_injected_transport() {
use std::net::SocketAddr;
use std::time::Instant;
use crate::transport::InMemoryNetwork;
let net = InMemoryNetwork::new();
let port = 5100u16;
let a_ip: IpAddr = "127.0.0.4".parse().unwrap();
let b_ip: IpAddr = "127.0.0.5".parse().unwrap();
let cfg = |ip: IpAddr, id: u64| {
ephemeral_config()
.with_listen_addr(ip)
.with_port(port)
.with_node_id(NodeId::new(id))
.with_reconcile_interval(Duration::from_millis(5))
};
let a = ReplicatedMap::<i32, i32>::new_with_transport(
cfg(a_ip, 1),
Arc::new(net.bind(SocketAddr::new(a_ip, port))),
);
let b = ReplicatedMap::<i32, i32>::new_with_transport(
cfg(b_ip, 2),
Arc::new(net.bind(SocketAddr::new(b_ip, port))),
);
a.engine.peers.write().insert(b_ip, Instant::now());
b.engine.peers.write().insert(a_ip, Instant::now());
a.insert(7, 42);
let ta = tokio::spawn(a.clone().run());
let tb = tokio::spawn(b.clone().run());
let deadline = Instant::now() + Duration::from_secs(5);
let mut converged = false;
while Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(20)).await;
if b.get(&7).as_deref() == Some(&42) {
converged = true;
break;
}
}
ta.abort();
tb.abort();
assert!(
converged,
"B never learned A's write over the injected transport"
);
}
#[tokio::test]
async fn new_returns_err_on_bind_failure() {
let holder =
std::net::UdpSocket::bind("127.0.0.50:0").expect("pre-condition: a port must be free");
let busy = holder
.local_addr()
.expect("holder must report its bound address");
let config = Config::default()
.with_port(busy.port())
.with_listen_addr(busy.ip())
.with_insecure_no_key();
let result = ReplicatedMap::<i32, i32>::new(config).await;
let err = result
.err()
.expect("expected Err when the bind address is already in use");
assert_eq!(
err.kind(),
std::io::ErrorKind::AddrInUse,
"bind failure should surface as AddrInUse, got {err:?}"
);
}