use core::cmp::Eq;
use core::hash::Hash;
use hashbrown::HashMap;
use std::fmt::Display;
use std::sync::{Arc, RwLock, RwLockWriteGuard};
use std::time::{Duration, SystemTime};
use crate::store::StoreItemPart;
pub(super) trait StoreGeneric {
fn ref_last_used(&self) -> &RwLock<SystemTime>;
}
pub(super) trait StoreGenericPool:
std::ops::Deref<Target = RwLock<HashMap<Self::StoreId, Arc<Self::Store>, Self::HashBuilder>>>
{
type StoreId: Hash + Eq;
type Store: StoreGeneric;
type HashBuilder: std::hash::BuildHasher;
fn kind() -> &'static str;
fn consider_inactive_after_secs(&self) -> u64;
fn access_lock(&self) -> &RwLock<()>;
fn proceed_erase_collection(&self, collection: StoreItemPart) -> Result<u32, ()>;
fn proceed_erase_bucket(
&self,
collection: StoreItemPart,
bucket: StoreItemPart,
) -> Result<u32, ()>;
}
pub(super) trait StoreGenericPoolExt: StoreGenericPool {
fn proceed_acquire_cache(
store_id: Self::StoreId,
store: &Arc<Self::Store>,
) -> Result<Arc<Self::Store>, ()>
where
Self::StoreId: Display,
{
let kind = Self::kind();
tracing::debug!("{kind} store {store_id} acquired from pool");
*store.ref_last_used().write().unwrap() = SystemTime::now();
Ok(Arc::clone(store))
}
#[allow(clippy::type_complexity, reason = "We can’t avoid it")]
fn proceed_acquire_open<'a>(
&'a self,
store_id: Self::StoreId,
build: impl FnOnce(&'a Self, Self::StoreId) -> Result<Self::Store, ()>,
write_guard: Option<
&mut RwLockWriteGuard<'a, HashMap<Self::StoreId, Arc<Self::Store>, Self::HashBuilder>>,
>,
) -> Result<Arc<Self::Store>, ()>
where
Self::StoreId: Display + Copy,
{
let kind = Self::kind();
match build(self, store_id) {
Ok(store) => {
let store_pool_write = match write_guard {
Some(x) => x,
None => &mut self.write().unwrap(),
};
let store_box = Arc::new(store);
store_pool_write.insert(store_id, Arc::clone(&store_box));
tracing::debug!("opened and cached {kind} store {store_id}");
Ok(store_box)
}
Err(error) => {
tracing::error!("failed opening {kind} store {store_id}: {error:?}");
Err(())
}
}
}
fn proceed_janitor(&self, filter: impl Fn(&Self::StoreId) -> bool)
where
Self::StoreId: Display + Copy,
{
let kind = Self::kind();
tracing::debug!("scanning for {kind} store pool items to janitor");
let _access = self.access_lock().write().unwrap();
let mut removal_register: Vec<Self::StoreId> = Vec::new();
let store_pool_read = self.read().unwrap();
for (collection_bucket, store) in store_pool_read.iter().filter(|(key, _)| filter(key)) {
let last_used_elapsed = (store.ref_last_used().read().unwrap())
.elapsed()
.unwrap_or_else(|err| {
tracing::error!(
"store pool item: {} last used duration clock issue, zeroing: {}",
collection_bucket,
err
);
Duration::ZERO
});
if last_used_elapsed.as_secs() >= self.consider_inactive_after_secs() {
tracing::debug!(
"found expired {kind} store pool item: {}; elapsed time: {last_used_elapsed:.1?}",
collection_bucket
);
removal_register.push(*collection_bucket);
} else {
tracing::debug!(
"found non-expired {kind} store pool item: {}; elapsed time: {last_used_elapsed:.1?}",
collection_bucket
);
}
}
let store_pool_read = if removal_register.is_empty() {
store_pool_read
} else {
drop(store_pool_read);
let mut store_pool_write = self.write().unwrap();
for collection_bucket in removal_register.iter() {
store_pool_write.remove(collection_bucket);
}
RwLockWriteGuard::downgrade(store_pool_write)
};
tracing::info!(
"done scanning for {kind} store pool items to janitor, expired {} items, now has {} items",
removal_register.len(),
store_pool_read.len(),
);
}
fn dispatch_erase(
&self,
collection: StoreItemPart,
bucket: Option<StoreItemPart>,
) -> Result<u32, ()> {
let kind = Self::kind();
tracing::info!("{kind} erase requested on collection: {collection}");
if let Some(bucket) = bucket {
self.proceed_erase_bucket(collection, bucket)
} else {
self.proceed_erase_collection(collection)
}
}
}
impl<Pool: StoreGenericPool> StoreGenericPoolExt for Pool {}
#[inline]
pub(super) fn u32_from_hex(hex: &str) -> Result<u32, std::io::Error> {
u32::from_str_radix(hex, 16).map_err(std::io::Error::other)
}