use std::net::IpAddr;
use std::sync::Arc;
use std::sync::Mutex;
use std::time::Duration;
use crate::discovery::{DiscoverFuture, Discovery, DiscoveryKind};
use crate::{
replicated_map::{Config, MemberPresence},
ReplicatedMap,
};
use super::ephemeral_config;
#[derive(Clone)]
struct FakeDiscovery {
resp: Arc<Mutex<FakeResp>>,
}
#[derive(Clone)]
enum FakeResp {
Present(Vec<IpAddr>),
Blip,
}
impl FakeDiscovery {
fn new(initial: FakeResp) -> Self {
FakeDiscovery {
resp: Arc::new(Mutex::new(initial)),
}
}
fn set(&self, resp: FakeResp) {
*self.resp.lock().unwrap() = resp;
}
}
impl Discovery for FakeDiscovery {
fn discover(&self) -> DiscoverFuture<'_> {
let resp = self.resp.lock().unwrap().clone();
Box::pin(async move {
match resp {
FakeResp::Present(addrs) => Ok(addrs),
FakeResp::Blip => Err(std::io::Error::other("blip")),
}
})
}
fn kind(&self) -> DiscoveryKind {
DiscoveryKind::Authoritative
}
}
struct SpeculativeDiscovery;
impl Discovery for SpeculativeDiscovery {
fn discover(&self) -> DiscoverFuture<'_> {
Box::pin(async { Ok(Vec::new()) })
}
fn kind(&self) -> DiscoveryKind {
DiscoveryKind::Speculative
}
}
#[tokio::test]
#[should_panic(expected = "with_discovery expects an authoritative source")]
async fn with_discovery_rejects_a_speculative_source() {
let store = ReplicatedMap::<i32, i32>::new(discovery_config())
.await
.expect("bind failed");
let _ = store.with_discovery(Arc::new(SpeculativeDiscovery));
}
fn discovery_config() -> Config {
Config::default()
.with_port(0)
.with_listen_addr("127.0.0.1".parse().unwrap())
.with_insecure_no_key()
}
#[tokio::test(flavor = "multi_thread")]
async fn discovery_decommissions_vanished_member_but_not_self() {
let own: IpAddr = "127.0.0.1".parse().unwrap();
let member: IpAddr = "127.0.0.200".parse().unwrap();
let fake = FakeDiscovery::new(FakeResp::Present(vec![member]));
let store = ReplicatedMap::<i32, i32>::new(discovery_config())
.await
.expect("bind failed")
.with_discovery(Arc::new(fake.clone()))
.with_discovery_interval(Duration::from_millis(20))
.with_discovery_miss_threshold(3);
store.engine.members.write().insert(own);
store.engine.members.write().insert(member);
let loop_store = store.clone();
let handle = tokio::spawn(async move { loop_store.discover_periodically().await });
tokio::time::sleep(Duration::from_millis(120)).await;
assert!(
store.engine.members.read().contains(&member),
"present member was wrongly decommissioned"
);
fake.set(FakeResp::Present(vec![]));
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
!store.engine.members.read().contains(&member),
"vanished member was not decommissioned after the grace period"
);
assert!(
store.engine.members.read().contains(&own),
"node decommissioned itself"
);
handle.abort();
}
#[tokio::test(flavor = "multi_thread")]
async fn discovery_blip_does_not_decommission() {
let member: IpAddr = "127.0.0.201".parse().unwrap();
let fake = FakeDiscovery::new(FakeResp::Present(vec![member]));
let store = ReplicatedMap::<i32, i32>::new(discovery_config())
.await
.expect("bind failed")
.with_discovery(Arc::new(fake.clone()))
.with_discovery_interval(Duration::from_millis(20))
.with_discovery_miss_threshold(3);
store.engine.members.write().insert(member);
let loop_store = store.clone();
let handle = tokio::spawn(async move { loop_store.discover_periodically().await });
tokio::time::sleep(Duration::from_millis(60)).await;
fake.set(FakeResp::Blip);
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
store.engine.members.read().contains(&member),
"a transient discovery failure wrongly decommissioned a member"
);
fake.set(FakeResp::Present(vec![]));
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
!store.engine.members.read().contains(&member),
"member was not decommissioned once it genuinely vanished"
);
handle.abort();
}
#[test]
fn member_presence_starts_present_and_requires_the_miss_threshold() {
let mut state = MemberPresence::default();
assert!(!state.eligible_for_decommission(1, Duration::ZERO, false));
state.mark_missed();
assert!(state.eligible_for_decommission(1, Duration::ZERO, false));
}
#[test]
fn member_presence_reappearance_resets_the_absence_clock_and_counter() {
let mut state = MemberPresence::default();
state.mark_missed();
state.mark_missed();
state.mark_seen();
state.mark_missed();
assert!(!state.eligible_for_decommission(2, Duration::ZERO, false));
}
#[test]
fn member_presence_pending_tombstone_acks_require_the_wall_time_floor() {
let mut state = MemberPresence::default();
state.mark_missed();
state.mark_missed();
assert!(!state.eligible_for_decommission(3, Duration::ZERO, true));
state.mark_missed();
assert!(!state.eligible_for_decommission(3, Duration::from_secs(3600), true));
assert!(state.eligible_for_decommission(3, Duration::from_secs(3600), false));
assert!(state.eligible_for_decommission(3, Duration::ZERO, true));
}
#[tokio::test(flavor = "multi_thread")]
async fn pending_tombstone_acks_hold_decommission_past_the_miss_threshold() {
let member: IpAddr = "127.0.0.210".parse().unwrap();
let fake = FakeDiscovery::new(FakeResp::Present(vec![member]));
let store = ReplicatedMap::<i32, i32>::new(discovery_config())
.await
.expect("bind failed")
.with_discovery(Arc::new(fake.clone()))
.with_discovery_interval(Duration::from_millis(15))
.with_discovery_miss_threshold(2)
.with_discovery_decommission_floor(Duration::from_millis(300));
store.engine.members.write().insert(member);
store.just_insert(1, 11);
store.just_remove(&1);
let loop_store = store.clone();
let handle = tokio::spawn(async move { loop_store.discover_periodically().await });
tokio::time::sleep(Duration::from_millis(60)).await;
fake.set(FakeResp::Present(vec![]));
tokio::time::sleep(Duration::from_millis(120)).await;
assert!(
store.engine.members.read().contains(&member),
"member with a pending tombstone ack was decommissioned before the wall-time floor \
elapsed"
);
tokio::time::sleep(Duration::from_millis(350)).await;
assert!(
!store.engine.members.read().contains(&member),
"member was never decommissioned even after the wall-time floor elapsed"
);
handle.abort();
}
#[tokio::test(flavor = "multi_thread")]
async fn reappearance_resets_the_floor_for_a_member_with_pending_acks() {
let member: IpAddr = "127.0.0.211".parse().unwrap();
let fake = FakeDiscovery::new(FakeResp::Present(vec![member]));
let store = ReplicatedMap::<i32, i32>::new(discovery_config())
.await
.expect("bind failed")
.with_discovery(Arc::new(fake.clone()))
.with_discovery_interval(Duration::from_millis(15))
.with_discovery_miss_threshold(2)
.with_discovery_decommission_floor(Duration::from_millis(300));
store.engine.members.write().insert(member);
store.just_insert(1, 11);
store.just_remove(&1);
let loop_store = store.clone();
let handle = tokio::spawn(async move { loop_store.discover_periodically().await });
fake.set(FakeResp::Present(vec![]));
tokio::time::sleep(Duration::from_millis(200)).await;
fake.set(FakeResp::Present(vec![member]));
tokio::time::sleep(Duration::from_millis(60)).await;
fake.set(FakeResp::Present(vec![]));
tokio::time::sleep(Duration::from_millis(200)).await;
assert!(
store.engine.members.read().contains(&member),
"reappearance did not reset the wall-time floor's absence clock"
);
handle.abort();
}
#[tokio::test]
async fn a_continuously_present_member_is_never_decommissioned() {
use crate::discovery::{DiscoverFuture, Discovery, DiscoveryKind};
#[derive(Clone)]
struct AlwaysPresent(IpAddr);
impl Discovery for AlwaysPresent {
fn discover(&self) -> DiscoverFuture<'_> {
let addr = self.0;
Box::pin(async move { Ok(vec![addr]) })
}
fn kind(&self) -> DiscoveryKind {
DiscoveryKind::Authoritative
}
}
let peer: IpAddr = "127.0.0.201".parse().unwrap();
let store = ReplicatedMap::<i32, i32>::new(
ephemeral_config().with_listen_addr("127.0.0.200".parse().unwrap()),
)
.await
.expect("bind failed")
.with_discovery_interval(Duration::from_millis(5))
.with_discovery_miss_threshold(1)
.with_discovery(Arc::new(AlwaysPresent(peer)));
store.engine.members.write().insert(peer);
let _ = tokio::time::timeout(Duration::from_millis(80), store.discover_periodically()).await;
assert!(
store.engine.members.read().contains(&peer),
"a continuously-present member must never be decommissioned"
);
}