Skip to main content

amalgam/
memory.rs

1//! The L1 (in-memory) store, backed by [`moka`].
2//!
3//! We use moka purely as a concurrent, evicting storage layer with per-entry
4//! physical expiration; the single-flight, fail-safe and timeout logic lives in
5//! [`Cache`](crate::Cache) so we keep full control over the flow (moka's own
6//! `get_with` coalescing can't express "return a stale value now and finish in
7//! the background").
8
9use std::sync::Arc;
10use std::time::{Duration, Instant};
11
12use moka::Expiry;
13use moka::future::Cache as MokaCache;
14use moka::notification::RemovalCause;
15
16use crate::entry::Entry;
17use crate::events::{CacheEvent, Events};
18
19/// Per-entry expiry policy: each entry carries the physical TTL the backend
20/// should honour ([`Entry::backend_ttl`]).
21struct EntryExpiry;
22
23impl<V: Send + Sync + 'static> Expiry<Arc<str>, Entry<V>> for EntryExpiry {
24    fn expire_after_create(
25        &self,
26        _key: &Arc<str>,
27        value: &Entry<V>,
28        _created_at: Instant,
29    ) -> Option<Duration> {
30        Some(value.backend_ttl())
31    }
32
33    fn expire_after_update(
34        &self,
35        _key: &Arc<str>,
36        value: &Entry<V>,
37        _updated_at: Instant,
38        _duration_until_expiry: Option<Duration>,
39    ) -> Option<Duration> {
40        Some(value.backend_ttl())
41    }
42
43    // `expire_after_read` keeps the default (no sliding): expiration is absolute,
44    // matching FusionCache's behaviour.
45}
46
47/// The L1 memory store.
48#[derive(Clone)]
49pub struct MemoryStore<V: Clone + Send + Sync + 'static> {
50    cache: MokaCache<Arc<str>, Entry<V>>,
51}
52
53impl<V: Clone + Send + Sync + 'static> MemoryStore<V> {
54    /// Creates a store with an optional maximum entry capacity (`None` =
55    /// unbounded). `events` receives an [`Eviction`](CacheEvent::Eviction) event
56    /// whenever the backend evicts an entry for expiry or capacity (not for our
57    /// own explicit removes/updates).
58    #[must_use]
59    pub fn new(max_capacity: Option<u64>, events: Events) -> Self {
60        let listener = move |key: Arc<Arc<str>>, _value: Entry<V>, cause: RemovalCause| {
61            if matches!(cause, RemovalCause::Expired | RemovalCause::Size) {
62                events.emit(CacheEvent::Eviction {
63                    key: key.as_ref().clone(),
64                });
65            }
66        };
67        let mut builder = MokaCache::builder()
68            .expire_after(EntryExpiry)
69            .eviction_listener(listener);
70        if let Some(capacity) = max_capacity {
71            builder = builder.max_capacity(capacity);
72        }
73        Self {
74            cache: builder.build(),
75        }
76    }
77
78    /// Reads an entry, if present and not physically expired.
79    pub async fn get(&self, key: &str) -> Option<Entry<V>> {
80        self.cache.get(key).await
81    }
82
83    /// Writes an entry.
84    pub async fn insert(&self, key: Arc<str>, entry: Entry<V>) {
85        self.cache.insert(key, entry).await;
86    }
87
88    /// Removes an entry, returning the previous value if any.
89    pub async fn remove(&self, key: &str) -> Option<Entry<V>> {
90        self.cache.remove(key).await
91    }
92
93    /// Invalidates every entry (lazily).
94    pub fn invalidate_all(&self) {
95        self.cache.invalidate_all();
96    }
97
98    /// Forces pending maintenance to run — primarily for deterministic tests.
99    pub async fn run_pending_tasks(&self) {
100        self.cache.run_pending_tasks().await;
101    }
102}