use crate::{
messaging::data::StorageLevel,
routing::{Prefix, XorName},
};
use itertools::Itertools;
use std::{
collections::{BTreeMap, BTreeSet},
sync::Arc,
};
use tokio::sync::RwLock;
pub(crate) const CHUNK_COPY_COUNT: usize = 4;
pub(crate) const MIN_LEVEL_WHEN_FULL: u8 = 9;
#[derive(Clone)]
pub(crate) struct Capacity {
adult_levels: Arc<RwLock<BTreeMap<XorName, Arc<RwLock<StorageLevel>>>>>,
}
impl Capacity {
pub(super) fn new(adult_levels: BTreeMap<XorName, StorageLevel>) -> Self {
let adult_levels = adult_levels
.into_iter()
.map(|(adult, level)| (adult, Arc::new(RwLock::new(level))))
.collect();
Self {
adult_levels: Arc::new(RwLock::new(adult_levels)),
}
}
pub(super) async fn is_full(&self, adult: &XorName) -> Option<bool> {
let adult_levels = self.adult_levels.read().await;
let level = adult_levels.get(adult)?.read().await.value();
Some(level >= MIN_LEVEL_WHEN_FULL)
}
pub(super) async fn avg_usage(&self) -> u8 {
let mut total = 0_usize;
let levels = self.adult_levels.read().await;
let levels = levels.values().collect_vec();
let num_adults = levels.len();
if num_adults == 0 {
return 0; }
for v in levels {
total += v.read().await.value() as usize;
}
(total / num_adults) as u8
}
pub(super) async fn levels(&self) -> BTreeMap<XorName, StorageLevel> {
let mut map = BTreeMap::new();
for (name, level) in self.adult_levels.read().await.iter() {
let _prev = map.insert(*name, *level.read().await);
}
map
}
pub(super) async fn levels_matching(&self, prefix: Prefix) -> BTreeMap<XorName, StorageLevel> {
self.levels()
.await
.iter()
.filter(|(name, _)| prefix.matches(name))
.map(|(name, level)| (*name, *level))
.collect()
}
pub(super) async fn full_adults(&self) -> BTreeSet<XorName> {
let mut set = BTreeSet::new();
for (name, level) in self.adult_levels.read().await.iter() {
if level.read().await.value() >= MIN_LEVEL_WHEN_FULL {
let _changed = set.insert(*name);
}
}
set
}
pub(super) async fn set_adult_levels(&self, levels: BTreeMap<XorName, StorageLevel>) {
for (name, level) in levels {
let _changed = self.set_adult_level(name, level).await;
}
}
pub(super) async fn set_adult_level(&self, adult: XorName, new_level: StorageLevel) -> bool {
{
let all_levels = self.adult_levels.read().await;
if let Some(level) = all_levels.get(&adult) {
let current_level = { level.read().await.value() };
info!("Current level: {}", current_level);
if new_level.value() > current_level {
*level.write().await = new_level;
info!("Old value overwritten.");
return true; }
return false; }
}
info!("No current level, aqcuiring top level write lock..");
let mut all_levels = self.adult_levels.write().await;
info!("Top level write lock aqcuired.");
if let Some(level) = all_levels.get(&adult) {
info!("Oh wait, a value was just recorded..");
let current_level = { level.read().await.value() };
info!("Current level: {}", current_level);
if new_level.value() > current_level {
*level.write().await = new_level;
info!("Old value overwritten.");
return true; }
false } else {
let _level = all_levels.insert(adult, Arc::new(RwLock::new(new_level)));
info!("New value inserted.");
true }
}
pub(super) async fn retain_members_only(&self, members: &BTreeSet<XorName>) {
let mut adult_levels = self.adult_levels.write().await;
let absent_adults: Vec<_> = adult_levels
.iter()
.filter(|(key, _)| !members.contains(key))
.map(|(key, _)| *key)
.collect();
for adult in &absent_adults {
let _level = adult_levels.remove(adult);
}
}
}