mod collections;
use collections::{first, Map, Prependable, Set};
use fnv::FnvBuildHasher;
use std::{
borrow::Borrow,
cmp,
hash::{BuildHasher, Hash, Hasher},
iter::{self, FromIterator},
ops::{Index, RangeInclusive},
};
#[derive(Clone)]
pub struct RingBuilder<T, S = FnvBuildHasher> {
hasher: S,
vnodes: usize,
nodes: Vec<T>,
weighted_nodes: Vec<(T, usize)>,
}
impl<T: Hash + Eq + Clone> Default for RingBuilder<T> {
fn default() -> Self {
RingBuilder::new(Default::default())
}
}
impl<T: Hash + Eq + Clone, S: BuildHasher> RingBuilder<T, S> {
pub fn new(hasher: S) -> Self {
RingBuilder {
hasher,
vnodes: 10,
nodes: vec![],
weighted_nodes: vec![],
}
}
pub fn vnodes(mut self, vnodes: usize) -> Self {
self.vnodes = cmp::max(1, vnodes);
self
}
pub fn weighted_node(mut self, node: T, vnodes: usize) -> Self {
self.weighted_nodes.push((node, vnodes));
self
}
pub fn weighted_nodes(mut self, weighted_nodes: &[(T, usize)]) -> Self {
self.weighted_nodes.extend_from_slice(weighted_nodes);
self
}
pub fn weighted_nodes_iter<I>(mut self, weighted_nodes: I) -> Self
where
I: Iterator<Item = (T, usize)>,
{
weighted_nodes.for_each(|w_node| self.weighted_nodes.push(w_node));
self
}
pub fn node(mut self, node: T) -> Self {
self.nodes.push(node);
self
}
pub fn nodes(mut self, nodes: &[T]) -> Self {
self.nodes.extend_from_slice(nodes);
self
}
pub fn nodes_iter<I>(mut self, nodes: I) -> Self
where
I: IntoIterator<Item = T>,
{
self.nodes.extend(nodes);
self
}
pub fn build(self) -> Ring<T, S> {
let mut ring = Ring {
n_vnodes: self.vnodes,
hasher: self.hasher,
vnodes: Vec::with_capacity(self.vnodes * self.nodes.len()),
unique: Vec::with_capacity(self.nodes.len() + self.weighted_nodes.len()),
};
let vnodes = self.vnodes;
self.nodes
.into_iter()
.map(|n| (n, vnodes))
.chain(self.weighted_nodes)
.for_each(|(n, v)| {
ring.insert_weight(n, v);
});
ring
}
}
#[derive(Clone)]
pub struct Ring<T, S = FnvBuildHasher> {
n_vnodes: usize, hasher: S,
vnodes: Vec<(u64, (T, u64))>,
unique: Vec<(u64, usize)>,
}
impl<T: Hash + Eq + Clone> Default for Ring<T> {
fn default() -> Self {
RingBuilder::default().build()
}
}
impl<T: Hash + Eq + Clone> FromIterator<T> for Ring<T> {
fn from_iter<I: IntoIterator<Item = T>>(iter: I) -> Self {
RingBuilder::default().nodes_iter(iter).build()
}
}
impl<K: Hash, T: Hash + Eq + Clone, S: BuildHasher> Index<K> for Ring<T, S> {
type Output = T;
fn index(&self, index: K) -> &Self::Output {
self.get(index)
}
}
impl<T: Hash + Eq + Clone, S: BuildHasher> Extend<T> for Ring<T, S> {
fn extend<I: IntoIterator<Item = T>>(&mut self, iter: I) {
let iter = iter.into_iter();
let (min, max) = iter.size_hint();
let n = max.unwrap_or(min);
self.vnodes.reserve(n * self.n_vnodes);
self.unique.reserve(n);
for node in iter {
self.insert(node);
}
}
}
impl<T: Hash + Eq + Clone, S: BuildHasher> Ring<T, S> {
pub fn len(&self) -> usize {
self.unique.len()
}
pub fn is_empty(&self) -> bool {
self.vnodes.is_empty()
}
pub fn vnodes(&self) -> usize {
self.vnodes.len()
}
pub fn weight<Q: ?Sized>(&self, node: &Q) -> Option<usize>
where
T: Borrow<Q>,
Q: Hash + Eq,
{
self.unique.map_lookup(&self.hash(node)).copied()
}
pub fn insert(&mut self, node: T) -> bool {
self.insert_weight(node, self.n_vnodes)
}
pub fn insert_weight(&mut self, node: T, vnodes: usize) -> bool {
let node_hash = self.hash(&node);
let mut hash = node_hash;
for _ in 0..vnodes.saturating_sub(1) {
self.vnodes.map_insert(hash, (node.clone(), node_hash));
hash = self.hash(hash);
}
if vnodes > 0 {
self.vnodes.map_insert(hash, (node, node_hash));
hash = self.hash(hash);
}
while self.vnodes.map_remove(&hash).is_some() {
hash = self.hash(hash);
}
self.unique.map_insert(node_hash, vnodes).is_none()
}
fn hash<K: Hash>(&self, key: K) -> u64 {
let mut digest = self.hasher.build_hasher();
key.hash(&mut digest);
digest.finish()
}
pub fn remove<Q: ?Sized>(&mut self, node: &Q) -> bool
where
T: Borrow<Q>,
Q: Hash + Eq,
{
self.vnodes.retain(|(_, (_node, _))| node != _node.borrow());
self.unique.map_remove(&self.hash(node)).is_some()
}
pub fn clear(&mut self) {
self.vnodes.clear();
self.unique.clear();
}
pub fn try_get<K: Hash>(&self, key: K) -> Option<&T> {
self.vnodes.find_gte(&self.hash(key)).map(first)
}
pub fn get<K: Hash>(&self, key: K) -> &T {
self.try_get(key).unwrap()
}
pub fn replicas<K: Hash>(&self, key: K) -> Candidates<'_, T, S> {
Candidates {
inner: self,
seen: Vec::with_capacity(self.len()),
hash: self.hash(&key),
}
}
unsafe fn get_root_hash(&self, vnode_idx: usize) -> u64 {
(self.vnodes.get_unchecked(vnode_idx).1).1
}
unsafe fn get_node_ref(&self, vnode_idx: usize) -> &T {
&(self.vnodes.get_unchecked(vnode_idx).1).0
}
pub fn resident_ranges(&self) -> impl Iterator<Item = Resident<'_, T>> + '_ {
let mut first = self.vnodes.first().map(|(_, (t, _))| t);
let mut vnodes = self.vnodes.iter();
let mut s = 0;
let mut raw = iter::from_fn(move || match vnodes.next() {
Some((h, (node, _))) => {
let next = Resident { keys: s..=*h, node };
s = h.overflowing_add(1).0;
Some(next)
}
None => Some(Resident {
keys: s..=u64::MAX,
node: first.take()?,
}),
})
.peekable();
iter::from_fn(move || {
let mut elt = raw.next()?;
while let Some(suc) = raw.peek() {
if suc.node != elt.node {
break;
}
let s = *elt.keys.start();
let e = *raw.next().unwrap().keys.end();
elt = Resident { keys: s..=e, ..elt };
}
Some(elt)
})
}
}
pub struct Candidates<'a, T, S = FnvBuildHasher> {
inner: &'a Ring<T, S>, seen: Vec<u64>, hash: u64, }
impl<'a, T: Hash + Eq + Clone, S: BuildHasher> Iterator for Candidates<'a, T, S> {
type Item = &'a T;
fn next(&mut self) -> Option<Self::Item> {
if self.seen.len() >= self.inner.len() {
return None;
}
let checked = |i| match self.inner.vnodes.len() {
n if n == 0 => None,
n if n == i => Some(0),
_ => Some(i),
};
let mut idx = (self.inner.vnodes)
.binary_search_by_key(&&self.hash, first)
.map_or_else(checked, Some)?;
while !self
.seen
.set_insert(unsafe { self.inner.get_root_hash(idx) })
{
if idx < self.inner.vnodes.len() - 1 {
idx += 1;
} else {
idx = 0;
}
}
if self.seen.len() < self.inner.len() {
self.hash = self.inner.hash(self.hash);
}
Some(unsafe { self.inner.get_node_ref(idx) })
}
fn size_hint(&self) -> (usize, Option<usize>) {
let n = self.inner.len() - self.seen.len();
(n, Some(n))
}
}
impl<'a, T: Hash + Eq + Clone, S: BuildHasher> ExactSizeIterator for Candidates<'a, T, S> {
fn len(&self) -> usize {
self.inner.len() - self.seen.len()
}
}
#[derive(Clone, Debug)]
pub struct Resident<'a, T> {
keys: RangeInclusive<u64>,
node: &'a T,
}
impl<'a, T> Resident<'a, T> {
pub fn keys(&self) -> &RangeInclusive<u64> {
&self.keys
}
pub fn node(&self) -> &'a T {
self.node
}
}
pub fn migrated_ranges<'a, T: Hash + Eq + Clone, S: BuildHasher>(
src: &'a Ring<T, S>,
dst: &'a Ring<T, S>,
) -> impl Iterator<Item = Migrated<'a, T>> + 'a {
let mut dst = Prependable::new(dst.resident_ranges());
src.resident_ranges().flat_map(move |src_elt| {
let mut out = vec![];
let mut s = *src_elt.keys.start();
while let Some(dst_elt) = dst.next() {
if src_elt.keys.end() < dst_elt.keys.end() {
dst.push_front(Resident {
keys: src_elt.keys.end().overflowing_add(1).0..=*dst_elt.keys.end(),
node: dst_elt.node,
});
if src_elt.node != dst_elt.node {
out.push(Migrated {
keys: s..=*src_elt.keys.end(),
src: src_elt.node,
dst: dst_elt.node,
});
}
break;
}
if src_elt.node != dst_elt.node {
out.push(Migrated {
keys: s..=*dst_elt.keys.end(),
src: src_elt.node,
dst: dst_elt.node,
});
}
if src_elt.keys.end() > dst_elt.keys.end() {
s = dst_elt.keys.end().overflowing_add(1).0;
continue;
}
break;
}
out
})
}
#[derive(Clone, Debug)]
pub struct Migrated<'a, T> {
keys: RangeInclusive<u64>,
src: &'a T,
dst: &'a T,
}
impl<'a, T> Migrated<'a, T> {
pub fn keys(&self) -> &RangeInclusive<u64> {
&self.keys
}
pub fn src(&self) -> &'a T {
self.src
}
pub fn dst(&self) -> &'a T {
self.dst
}
}