use std::hash::Hash;
use std::io;
use std::net::{IpAddr, SocketAddr};
use std::sync::Arc;
use std::time::{Duration, Instant};
use crate::bounds::{Key, Value};
use crate::clock::NodeId;
use crate::discovery::Discovery;
use crate::persistence::{InMemoryPersistence, Persistence};
use crate::replica::Replica;
use crate::timeout_wheel::TimeoutWheel;
use crate::transport::Transport;
mod config;
mod discovery;
mod membership;
mod mutate;
mod persistence;
mod read;
mod write;
pub(crate) use config::MIN_BULK_SEND_RATE;
pub use config::{Config, MAX_NETS};
#[cfg(test)]
pub(crate) use discovery::MemberPresence;
const DEFAULT_DISCOVERY_INTERVAL: Duration = Duration::from_secs(5);
const DEFAULT_DISCOVERY_MISS_THRESHOLD: u32 = 3;
const DEFAULT_DISCOVERY_DECOMMISSION_FLOOR: Duration = Duration::from_secs(600);
pub struct ReplicatedMap<K, V>
where
K: Clone + Hash + std::cmp::Eq + Send + Sync,
{
engine: Replica<K, V>,
tombstones: TimeoutWheel<K>,
persistence: Arc<dyn Persistence<K, V>>,
discovery: Option<Arc<dyn Discovery>>,
discovery_interval: Duration,
discovery_miss_threshold: u32,
discovery_decommission_floor: Duration,
}
impl<K, V> Clone for ReplicatedMap<K, V>
where
K: Clone + Hash + std::cmp::Eq + Send + Sync,
{
fn clone(&self) -> Self {
ReplicatedMap {
engine: self.engine.clone(),
tombstones: self.tombstones.clone(),
persistence: self.persistence.clone(),
discovery: self.discovery.clone(),
discovery_interval: self.discovery_interval,
discovery_miss_threshold: self.discovery_miss_threshold,
discovery_decommission_floor: self.discovery_decommission_floor,
}
}
}
impl<K: Key + Hash, V: Value> ReplicatedMap<K, V> {
pub async fn new(config: Config) -> io::Result<Self> {
Ok(Self::from_engine(Replica::<K, V>::new(config).await?))
}
pub fn new_with_transport(
config: Config,
transport: Arc<dyn Transport<Addr = SocketAddr>>,
) -> Self {
Self::from_engine(Replica::<K, V>::with_transport(config, transport))
}
pub async fn new_with_clock(
config: Config,
clock: Arc<dyn crate::clock::Clock>,
) -> io::Result<Self> {
Ok(Self::from_engine(
Replica::<K, V>::new_with_clock(config, clock).await?,
))
}
fn from_engine(engine: Replica<K, V>) -> Self {
let svc = ReplicatedMap {
engine,
tombstones: TimeoutWheel::new(),
persistence: Arc::new(InMemoryPersistence::default()),
discovery: None,
discovery_interval: DEFAULT_DISCOVERY_INTERVAL,
discovery_miss_threshold: DEFAULT_DISCOVERY_MISS_THRESHOLD,
discovery_decommission_floor: DEFAULT_DISCOVERY_DECOMMISSION_FLOOR,
};
svc.set_pre_insert(|_, _| {});
svc
}
pub fn node_id(&self) -> NodeId {
self.engine.node_id()
}
pub fn with_seed(self, peer: IpAddr) -> Self {
let now = Instant::now();
self.engine.peers.write().insert(peer, now);
self
}
pub fn seed_peer(&self, peer: IpAddr) {
self.engine.seed_peer(peer);
}
}
#[cfg(test)]
mod tests;