use std::sync::Arc;
use std::time::{Duration, Instant};
use moka::Expiry;
use moka::future::Cache as MokaCache;
use moka::notification::RemovalCause;
use crate::entry::Entry;
use crate::events::{CacheEvent, Events};
struct EntryExpiry;
impl<V: Send + Sync + 'static> Expiry<Arc<str>, Entry<V>> for EntryExpiry {
fn expire_after_create(
&self,
_key: &Arc<str>,
value: &Entry<V>,
_created_at: Instant,
) -> Option<Duration> {
Some(value.backend_ttl())
}
fn expire_after_update(
&self,
_key: &Arc<str>,
value: &Entry<V>,
_updated_at: Instant,
_duration_until_expiry: Option<Duration>,
) -> Option<Duration> {
Some(value.backend_ttl())
}
}
#[derive(Clone)]
pub struct MemoryStore<V: Clone + Send + Sync + 'static> {
cache: MokaCache<Arc<str>, Entry<V>>,
}
impl<V: Clone + Send + Sync + 'static> MemoryStore<V> {
#[must_use]
pub fn new(max_capacity: Option<u64>, events: Events) -> Self {
let listener = move |key: Arc<Arc<str>>, _value: Entry<V>, cause: RemovalCause| {
if matches!(cause, RemovalCause::Expired | RemovalCause::Size) {
events.emit(CacheEvent::Eviction {
key: key.as_ref().clone(),
});
}
};
let mut builder = MokaCache::builder()
.expire_after(EntryExpiry)
.eviction_listener(listener);
if let Some(capacity) = max_capacity {
builder = builder.max_capacity(capacity);
}
Self {
cache: builder.build(),
}
}
pub async fn get(&self, key: &str) -> Option<Entry<V>> {
self.cache.get(key).await
}
pub async fn insert(&self, key: Arc<str>, entry: Entry<V>) {
self.cache.insert(key, entry).await;
}
pub async fn remove(&self, key: &str) -> Option<Entry<V>> {
self.cache.remove(key).await
}
pub fn invalidate_all(&self) {
self.cache.invalidate_all();
}
pub async fn run_pending_tasks(&self) {
self.cache.run_pending_tasks().await;
}
}