use std::hash::Hash;
use std::net::{IpAddr, SocketAddr};
use std::sync::Arc;
use crate::bounds::{Key, Value};
use crate::clock::Timestamp;
use crate::entry::{Entry, State};
use crate::observability;
use crate::FingerprintTreeMap;
use super::{send_messages_to, Message, Replica, SendPorts, PEER_EXPIRATION};
impl<K: Key + Hash, V: Value> Replica<K, V> {
pub(super) fn map_insert(
&self,
guard: &mut FingerprintTreeMap<K, Entry<Timestamp, V>>,
key: K,
value: Entry<Timestamp, V>,
) -> Option<Entry<Timestamp, V>> {
{
let mut live_tombstones = self.live_tombstones.write();
if value.is_tombstone() {
live_tombstones.insert(key.clone());
} else {
live_tombstones.remove(&key);
}
}
self.projection.write().insert(key.clone(), value.project());
guard.insert(key, value)
}
pub(super) fn get_peers(&self) -> Vec<IpAddr> {
let mut guard = self.peers.write();
guard.retain(|_, instant| instant.elapsed() < PEER_EXPIRATION);
guard.keys().cloned().collect()
}
pub fn just_insert(&self, key: K, value: Entry<Timestamp, V>) -> Option<Entry<Timestamp, V>> {
(self.pre_insert.read())(&key, &value);
if value.is_tombstone() {
observability::record_remove();
} else {
observability::record_insert();
}
let mut guard = self.map.write();
self.map_insert(&mut guard, key, value)
}
fn broadcast(&self, messages: Vec<Message<K, Entry<Timestamp, V>, State<V>>>) {
let peers = self.get_peers();
let port = self.port;
let transport = Arc::clone(&self.transport);
let authenticator = self.authenticator.clone();
let sender_counter = Arc::clone(&self.sender_counter);
tokio::spawn(async move {
let ports = SendPorts {
transport: &*transport,
authenticator: &authenticator,
sender_counter: &sender_counter,
};
let mut send_buf = Vec::new();
for addr in peers {
let peer = SocketAddr::new(addr, port);
send_messages_to(&messages, &ports, &peer, &mut send_buf).await;
}
});
}
pub fn insert(&self, key: K, value: Entry<Timestamp, V>) -> Option<Entry<Timestamp, V>> {
let ret = self.just_insert(key.clone(), value.clone());
self.broadcast(vec![Message::Update::<K, Entry<Timestamp, V>, State<V>>((
key, value,
))]);
ret
}
pub(crate) fn broadcast_update(&self, key: K, value: Entry<Timestamp, V>) {
self.broadcast(vec![Message::Update::<K, Entry<Timestamp, V>, State<V>>((
key, value,
))]);
}
pub fn just_insert_bulk(&self, key_values: &[(K, Entry<Timestamp, V>)]) {
for (key, value) in key_values {
(self.pre_insert.read())(key, value);
if value.is_tombstone() {
observability::record_remove();
} else {
observability::record_insert();
}
}
let mut guard = self.map.write();
for (key, value) in key_values {
self.map_insert(&mut guard, key.clone(), value.clone());
}
}
pub fn insert_bulk(&self, key_values: &[(K, Entry<Timestamp, V>)]) {
self.just_insert_bulk(key_values);
let messages: Vec<_> = key_values
.iter()
.map(|kv| Message::Update::<K, Entry<Timestamp, V>, State<V>>(kv.clone()))
.collect();
self.broadcast(messages);
}
}