use std::net::SocketAddr;
use bincode::{DefaultOptions, Serializer};
use serde::Serialize;
use super::super::Message;
use crate::clock::{Hlc, LogicalCounter, NodeId, PhysicalTime, Timestamp};
use crate::entry::{Entry, State};
use crate::replica::{version_hash, Replica};
use crate::replicated_map::Config;
use gossip::auth;
type Tombstoned = Entry<Timestamp, i32>;
async fn engine(addr: &str) -> Replica<i32, i32> {
let config = Config::default()
.with_port(0)
.with_listen_addr(addr.parse().unwrap())
.with_insecure_no_key();
Replica::new(config).await.expect("bind failed")
}
fn ack_bytes(key: i32, version: u64) -> Vec<u8> {
let msg = Message::Ack::<i32, Tombstoned, State<i32>>((key, version));
let mut buf = vec![gossip::auth::WIRE_VERSION];
msg.serialize(&mut Serializer::new(&mut buf, DefaultOptions::new()))
.unwrap();
buf
}
#[tokio::test]
async fn ack_for_unknown_key_does_not_grow_tombstone_acks() {
let eng = engine("127.0.0.93").await;
let peer: SocketAddr = "127.0.0.94:9000".parse().unwrap();
let bytes = ack_bytes(42, 999);
let payload = auth::Authenticator::new(None, false)
.open(&bytes)
.expect("unauthenticated open")
.check_version()
.expect("ack_bytes stamps the current wire version");
let payload = payload
.verify_replay(&eng.replay_filter, peer.ip())
.expect("unauthenticated mode is exempt from the replay check");
let mut send_buf = Vec::new();
eng.handle_messages(payload, peer, &mut send_buf).await;
assert_eq!(
eng.tombstone_acks_len(),
0,
"ack for unknown key must not insert into tombstone_acks"
);
}
#[tokio::test]
async fn ack_for_live_key_does_not_grow_tombstone_acks() {
let eng = engine("127.0.0.95").await;
let key = 10;
eng.just_insert(
key,
Entry::present(
Timestamp::new(
Hlc::new(PhysicalTime::from_millis(1), LogicalCounter::new(0)),
NodeId::new(0),
),
42,
),
);
let peer: SocketAddr = "127.0.0.96:9000".parse().unwrap();
let bytes = ack_bytes(key, 123);
let payload = auth::Authenticator::new(None, false)
.open(&bytes)
.expect("unauthenticated open")
.check_version()
.expect("ack_bytes stamps the current wire version");
let payload = payload
.verify_replay(&eng.replay_filter, peer.ip())
.expect("unauthenticated mode is exempt from the replay check");
let mut send_buf = Vec::new();
eng.handle_messages(payload, peer, &mut send_buf).await;
assert_eq!(
eng.tombstone_acks_len(),
0,
"ack for a live (non-tombstone) key must not insert into tombstone_acks"
);
}
#[tokio::test]
async fn ack_for_local_tombstone_is_recorded() {
let eng = engine("127.0.0.97").await;
let key = 20;
let tombstone: Tombstoned = Entry::tombstone(Timestamp::new(
Hlc::new(PhysicalTime::from_millis(2), LogicalCounter::new(0)),
NodeId::new(0),
));
let version = version_hash(&tombstone);
eng.just_insert(key, tombstone);
let peer: SocketAddr = "127.0.0.98:9000".parse().unwrap();
let bytes = ack_bytes(key, version);
let payload = auth::Authenticator::new(None, false)
.open(&bytes)
.expect("unauthenticated open")
.check_version()
.expect("ack_bytes stamps the current wire version");
let payload = payload
.verify_replay(&eng.replay_filter, peer.ip())
.expect("unauthenticated mode is exempt from the replay check");
let mut send_buf = Vec::new();
eng.handle_messages(payload, peer, &mut send_buf).await;
assert_eq!(
eng.tombstone_acks_len(),
1,
"ack for a local tombstone must be recorded in tombstone_acks"
);
eng.members.write().insert(peer.ip());
assert!(
eng.is_tombstone_stable(&key, version),
"tombstone should be stable after the only member has acked"
);
eng.forget_tombstone(&key);
assert_eq!(
eng.tombstone_acks_len(),
0,
"forget_tombstone must clear tombstone_acks for the key"
);
}