fraiseql_core/cache/result.rs
1//! Query result caching with W-TinyLFU eviction and per-entry TTL.
2//!
3//! This module provides a `moka::sync::Cache`-backed store for GraphQL query results.
4//! Moka uses Concurrent W-TinyLFU policy with lock-free reads — cache hits do NOT
5//! acquire any shared lock, eliminating the hot-key serialisation bottleneck present
6//! in the old 64-shard `parking_lot::Mutex<LruCache>` design.
7//!
8//! ## Performance characteristics
9//!
10//! - **`get()` hot path** (cache hit): lock-free frequency-counter update (thread-local ring
11//! buffer, drained lazily on writes), `Arc` clone (single atomic increment), one atomic counter
12//! bump.
13//! - **`put()` path**: early-exit guards (disabled / list / size) before touching the store.
14//! Reverse-index updates use `DashMap` (fine-grained sharding, no global lock).
15//! - **`metrics()`**: reads `store.entry_count()` directly — no shard scan.
16//! - **`invalidate_views()` / `invalidate_by_entity()`**: O(k) where k = matching entries (via
17//! reverse indexes), not O(total entries).
18//!
19//! ## Reverse indexes
20//!
21//! Because `moka` does not support arbitrary iteration, view-based and entity-based
22//! invalidation rely on two `DashMap` reverse indexes maintained alongside the store:
23//!
24//! ```text
25//! view_index: DashMap<view_name, DashSet<(cache_key, epoch)>>
26//! entity_index: DashMap<entity_type, DashMap<entity_id, DashSet<(cache_key, epoch)>>>
27//! ```
28//!
29//! Indexes are populated in `put()` and pruned via moka's eviction listener.
30//! `clear()` resets all indexes synchronously.
31//!
32//! ### Why the indexes key on `(cache_key, epoch)` and not on `cache_key` alone (#740)
33//!
34//! `put_arc()` registers the key in the reverse indexes *before* `store.insert()`,
35//! deliberately: an `invalidate_views()` racing the insert must not miss the key.
36//! The consequence is that moka fires the eviction listener for the entry the insert
37//! just displaced — while the replacement is already live under the same key. A
38//! listener that pruned by `cache_key` would delete the registrations the *live*
39//! entry depends on, detaching it from every invalidation path: it would then be
40//! served until TTL expiry, or forever for the `cache_ttl_seconds = 0` entries
41//! documented as "mutation-invalidated only".
42//!
43//! Each cached entry therefore carries a process-unique [`CachedResult::epoch`] and
44//! registers/deregisters itself under `(cache_key, epoch)`. Registration and removal
45//! are symmetric per *entry instance* rather than per key, so the listener needs no
46//! knowledge of moka's [`moka::notification::RemovalCause`] taxonomy — `Replaced`,
47//! `Expired`, `Size` and `Explicit` are all handled by the same arithmetic. Keep it
48//! that way: a cause-inspecting listener is only as correct as the exact set of
49//! causes the current moka version reports for a displaced entry.
50
51use std::{
52 collections::HashSet,
53 sync::{
54 Arc,
55 atomic::{AtomicU64, AtomicUsize, Ordering},
56 },
57 time::Duration,
58};
59
60use dashmap::{DashMap, DashSet};
61use fraiseql_db::ViewName;
62use moka::sync::Cache as MokaCache;
63use serde::{Deserialize, Serialize};
64
65use super::config::CacheConfig;
66use crate::{db::types::JsonbValue, error::Result};
67
68/// Cached query result with metadata.
69///
70/// Stores the query result along with tracking information for
71/// TTL expiry, view-based invalidation, and monitoring.
72#[derive(Debug, Clone)]
73pub struct CachedResult {
74 /// The actual query result (JSONB array from database).
75 ///
76 /// Wrapped in `Arc` for cheap cloning on cache hits (zero-copy).
77 pub result: Arc<Vec<JsonbValue>>,
78
79 /// Which views/tables this query accesses.
80 ///
81 /// Format: `[ViewName::from("v_user"), ViewName::from("v_post")]`
82 ///
83 /// Stored as a boxed slice of [`ViewName`] (each backed by `Arc<str>`)
84 /// so cloning a name into the reverse index is a cheap atomic ref-count
85 /// bump rather than a fresh heap allocation. Views are fixed at `put()`
86 /// time and never modified.
87 pub accessed_views: Box<[ViewName]>,
88
89 /// When this entry was cached (Unix timestamp in seconds).
90 ///
91 /// Wall-clock timestamp for debugging. TTL enforcement is handled by moka
92 /// internally via `CacheEntryExpiry`.
93 pub cached_at: u64,
94
95 /// Per-entry TTL in seconds.
96 ///
97 /// Overrides `CacheConfig::ttl_seconds` when set via `put(..., Some(ttl))`.
98 /// Read by `CacheEntryExpiry::expire_after_create` to tell moka the expiry.
99 pub ttl_seconds: u64,
100
101 /// Entity references for selective entity-level invalidation.
102 ///
103 /// Contains one `(entity_type, entity_id)` pair per row in `result` that has
104 /// a valid string in its `"id"` field. Empty for queries with no `id` column
105 /// or when `put()` is called without an `entity_type`.
106 /// Used by the eviction listener to clean up `entity_index` on eviction.
107 pub entity_refs: Box<[(String, String)]>,
108
109 /// Process-unique identity of this cached *entry instance*.
110 ///
111 /// Distinguishes an entry from its own replacement under the same cache key.
112 /// The reverse indexes register `(cache_key, epoch)` pairs so the eviction
113 /// listener can deregister exactly the instance being evicted and never the
114 /// live one that displaced it — see the module docs (#740).
115 pub epoch: u64,
116}
117
118/// A reverse-index registration: the cache key plus the epoch of the entry
119/// instance that registered it.
120type IndexRef = (u64, u64);
121
122/// Moka `Expiry` implementation: reads TTL from `CachedResult.ttl_seconds`.
123struct CacheEntryExpiry;
124
125impl moka::Expiry<u64, Arc<CachedResult>> for CacheEntryExpiry {
126 fn expire_after_create(
127 &self,
128 _key: &u64,
129 value: &Arc<CachedResult>,
130 _created_at: std::time::Instant,
131 ) -> Option<Duration> {
132 if value.ttl_seconds == 0 {
133 // TTL=0 means "no time-based expiry" — entry lives until explicitly
134 // invalidated by a mutation. Return None so moka never schedules
135 // a timer-wheel eviction for this entry.
136 None
137 } else {
138 Some(Duration::from_secs(value.ttl_seconds))
139 }
140 }
141
142 // `expire_after_read` is intentionally NOT overridden.
143 //
144 // Moka's default returns `None` (no change to the timer) which skips the
145 // internal timer-wheel reschedule on every get(). Overriding it to return
146 // `duration_until_expiry` — even though the value is semantically unchanged —
147 // forces moka to acquire its timer-wheel lock on every cache hit. Under 40
148 // concurrent workers reading the same key, that lock becomes the new hot-key
149 // bottleneck, serialising reads and degrading list-query throughput ~3×.
150 //
151 // Entries expire at creation_time + ttl_seconds regardless of read frequency,
152 // which is the correct fixed-TTL semantics for query result caching.
153}
154
155/// Thread-safe W-TinyLFU cache for query results.
156///
157/// Backed by [`moka::sync::Cache`] which provides lock-free reads via
158/// Concurrent `TinyLFU`. Reverse `DashMap` indexes enable O(k) invalidation.
159///
160/// # Thread Safety
161///
162/// `moka::sync::Cache` is `Send + Sync`. All reverse indexes use `DashMap`
163/// (fine-grained shard locking) and `DashSet` (also shard-locked). There is no
164/// global mutex on the read path.
165///
166/// # Example
167///
168/// ```rust
169/// use fraiseql_core::cache::{QueryResultCache, CacheConfig};
170/// use fraiseql_core::db::types::JsonbValue;
171/// use serde_json::json;
172///
173/// let cache = QueryResultCache::new(CacheConfig::default());
174///
175/// // Cache a result
176/// let result = vec![JsonbValue::new(json!({"id": 1, "name": "Alice"}))];
177/// let fence = cache.invalidation_generation(); // snapshot before the read (#1079)
178/// cache.put(
179/// 12345_u64,
180/// result.clone(),
181/// vec!["v_user".to_string()],
182/// None, // use global TTL
183/// None, // no entity type index
184/// Some(fence), // discard the write if an invalidation raced the read
185/// ).unwrap();
186///
187/// // Retrieve from cache
188/// if let Some(cached) = cache.get(12345).unwrap() {
189/// println!("Cache hit! {} results", cached.len());
190/// }
191/// ```
192pub struct QueryResultCache {
193 /// Moka W-TinyLFU store.
194 ///
195 /// `Arc<CachedResult>` rather than `CachedResult` so that `get()` returns in
196 /// one atomic increment instead of deep-cloning the struct (which would copy
197 /// `accessed_views: Box<[String]>` on every cache hit).
198 store: MokaCache<u64, Arc<CachedResult>>,
199
200 /// Configuration (immutable after creation).
201 config: CacheConfig,
202
203 // Metrics counters — `Relaxed` ordering is sufficient: these counters are
204 // used only for monitoring, not for correctness or synchronisation.
205 hits: AtomicU64,
206 misses: AtomicU64,
207 total_cached: AtomicU64,
208 invalidations: AtomicU64,
209
210 /// Estimated total memory in use.
211 ///
212 /// Wrapped in `Arc` so the eviction listener closure (which requires `'static`)
213 /// can hold a clone and decrement on eviction.
214 memory_bytes: Arc<AtomicUsize>,
215
216 /// Reverse index: view name → set of `(cache key, epoch)` accessing that view.
217 ///
218 /// Keys are [`ViewName`] (`Arc<str>` inside) so inserts share the same
219 /// allocation as the names stored in [`CachedResult::accessed_views`].
220 /// Lookup by `&str` still works via the `Borrow<str>` impl on `ViewName`.
221 view_index: Arc<DashMap<ViewName, DashSet<IndexRef>>>,
222
223 /// Reverse index: entity type → entity id → set of `(cache key, epoch)`.
224 entity_index: Arc<DashMap<String, DashMap<String, DashSet<IndexRef>>>>,
225
226 /// Source of [`CachedResult::epoch`] values. Monotonic for the process
227 /// lifetime; wrap-around at 2^64 puts is not reachable.
228 next_epoch: AtomicU64,
229
230 /// Bumped by every invalidation, so a read that straddled one can tell (#1079).
231 ///
232 /// A read is `get` → miss → **await the database** → `put`. An invalidation landing
233 /// inside that await evicts nothing (the key is not in `view_index` yet) and the `put`
234 /// then stores rows fetched *before* the mutation committed — the client that just
235 /// wrote sees its own write vanish on the next read.
236 ///
237 /// The counter is **global**, not per view: a skipped cache write costs one uncached
238 /// read, whereas a stale entry costs correctness, and a per-view map would inherit
239 /// every lifetime question `view_index` already has. If a benchmark ever shows the
240 /// hit-rate loss matters, refine it then — with a measurement, not a guess.
241 ///
242 /// Distinct from [`next_epoch`](Self::next_epoch), which identifies one entry instance
243 /// so the eviction listener does not deregister a replacement (#740).
244 invalidation_generation: AtomicU64,
245}
246
247/// Cache metrics for monitoring.
248///
249/// Exposed via API for observability and debugging.
250#[derive(Debug, Clone, Serialize, Deserialize)]
251pub struct CacheMetrics {
252 /// Number of cache hits (returned cached result).
253 pub hits: u64,
254
255 /// Number of cache misses (executed query).
256 pub misses: u64,
257
258 /// Total entries cached across all time.
259 pub total_cached: u64,
260
261 /// Number of invalidations triggered.
262 pub invalidations: u64,
263
264 /// Current size of cache (number of entries).
265 pub size: usize,
266
267 /// Estimated memory usage in bytes.
268 ///
269 /// This is a rough estimate based on `CachedResult` struct size.
270 /// Actual memory usage may vary based on result sizes.
271 pub memory_bytes: usize,
272}
273
274/// Estimate the per-entry accounting overhead.
275const fn entry_overhead() -> usize {
276 std::mem::size_of::<CachedResult>() + std::mem::size_of::<u64>() * 2
277}
278
279/// Build the moka store, wiring the eviction listener to the reverse indexes
280/// and memory counter.
281fn build_store(
282 config: &CacheConfig,
283 memory_bytes: Arc<AtomicUsize>,
284 view_index: Arc<DashMap<ViewName, DashSet<IndexRef>>>,
285 entity_index: Arc<DashMap<String, DashMap<String, DashSet<IndexRef>>>>,
286) -> MokaCache<u64, Arc<CachedResult>> {
287 let max_cap = config.max_entries as u64;
288 let mb = memory_bytes;
289 let vi = view_index;
290 let ei = entity_index;
291
292 MokaCache::builder()
293 .max_capacity(max_cap)
294 .expire_after(CacheEntryExpiry)
295 .eviction_listener(move |key: Arc<u64>, value: Arc<CachedResult>, _cause| {
296 // Decrement memory budget so put()'s byte-gate stays accurate.
297 mb.fetch_sub(entry_overhead(), Ordering::Relaxed);
298
299 // Deregister exactly THIS entry instance. `_cause` is deliberately
300 // unused: `(key, epoch)` already distinguishes a displaced entry from
301 // the live replacement that displaced it, so `Replaced` needs no
302 // special case and no future moka cause can detach a live entry (#740).
303 let reference: IndexRef = (*key, value.epoch);
304
305 // Remove this instance from the view index.
306 for view in &value.accessed_views {
307 if let Some(keys) = vi.get(view) {
308 keys.remove(&reference);
309 }
310 }
311
312 // Remove ALL entity_refs from entity index.
313 for (et, id) in &*value.entity_refs {
314 if let Some(by_type) = ei.get(et) {
315 if let Some(keys) = by_type.get(id) {
316 keys.remove(&reference);
317 }
318 }
319 }
320 })
321 .build()
322}
323
324impl QueryResultCache {
325 /// Create new cache with configuration.
326 ///
327 /// # Panics
328 ///
329 /// Panics if `config.max_entries` is 0 (invalid configuration).
330 ///
331 /// # Example
332 ///
333 /// ```rust
334 /// use fraiseql_core::cache::{QueryResultCache, CacheConfig};
335 ///
336 /// let cache = QueryResultCache::new(CacheConfig::default());
337 /// ```
338 #[must_use]
339 pub fn new(config: CacheConfig) -> Self {
340 assert!(config.max_entries > 0, "max_entries must be > 0");
341
342 let memory_bytes = Arc::new(AtomicUsize::new(0));
343 let view_index: Arc<DashMap<ViewName, DashSet<IndexRef>>> = Arc::new(DashMap::new());
344 let entity_index: Arc<DashMap<String, DashMap<String, DashSet<IndexRef>>>> =
345 Arc::new(DashMap::new());
346
347 let store = build_store(
348 &config,
349 Arc::clone(&memory_bytes),
350 Arc::clone(&view_index),
351 Arc::clone(&entity_index),
352 );
353
354 Self {
355 store,
356 config,
357 hits: AtomicU64::new(0),
358 misses: AtomicU64::new(0),
359 total_cached: AtomicU64::new(0),
360 invalidations: AtomicU64::new(0),
361 memory_bytes,
362 view_index,
363 entity_index,
364 next_epoch: AtomicU64::new(0),
365 invalidation_generation: AtomicU64::new(0),
366 }
367 }
368
369 /// Snapshot of the invalidation generation, for fencing a cache write (#1079).
370 ///
371 /// Take this **before** the database round trip and hand it to
372 /// [`put_arc`](Self::put_arc). If any invalidation lands while the read is in flight,
373 /// the counter moves and the write is refused — the rows are still returned to the
374 /// caller, they are simply not cached.
375 #[must_use]
376 pub fn invalidation_generation(&self) -> u64 {
377 self.invalidation_generation.load(Ordering::Acquire)
378 }
379
380 /// Returns whether caching is enabled.
381 ///
382 /// Used by `CachedDatabaseAdapter` to short-circuit key generation
383 /// and result clone overhead when caching is disabled.
384 #[must_use]
385 pub const fn is_enabled(&self) -> bool {
386 self.config.enabled
387 }
388
389 /// The configuration this cache was built with (immutable after creation).
390 ///
391 /// The admin surface reports the effective TTL and entry ceiling, which used to be
392 /// hard-coded constants in the handler that happened to match a different cache's
393 /// defaults (#941).
394 #[must_use]
395 pub const fn config(&self) -> &CacheConfig {
396 &self.config
397 }
398
399 /// Look up a cached result by its cache key.
400 ///
401 /// Returns `None` when caching is disabled or the key is not present or expired.
402 /// Moka handles TTL expiry internally — if `get()` returns `Some`, the entry is live.
403 ///
404 /// # Errors
405 ///
406 /// This method is infallible. The `Result` return type is kept for API compatibility.
407 pub fn get(&self, cache_key: u64) -> Result<Option<Arc<Vec<JsonbValue>>>> {
408 if !self.config.enabled {
409 return Ok(None);
410 }
411
412 // moka::sync::Cache::get() is lock-free on the read path.
413 if let Some(cached) = self.store.get(&cache_key) {
414 self.hits.fetch_add(1, Ordering::Relaxed);
415 Ok(Some(Arc::clone(&cached.result)))
416 } else {
417 self.misses.fetch_add(1, Ordering::Relaxed);
418 Ok(None)
419 }
420 }
421
422 /// Store query result in cache, accepting an already-`Arc`-wrapped result.
423 ///
424 /// Preferred over [`put`](Self::put) on the hot miss path: callers that already
425 /// hold an `Arc<Vec<JsonbValue>>` (e.g. `CachedDatabaseAdapter`) can store it
426 /// without an extra `Vec` clone.
427 ///
428 /// # Arguments
429 ///
430 /// * `cache_key` - Cache key (from `generate_cache_key()`)
431 /// * `result` - Arc-wrapped query result to cache
432 /// * `accessed_views` - List of views accessed by this query
433 /// * `ttl_override` - Per-entry TTL in seconds; `None` uses `CacheConfig::ttl_seconds`
434 /// * `entity_type` - Optional GraphQL type name for entity-ID indexing
435 /// * `fence` - Generation snapshot from
436 /// [`invalidation_generation`](Self::invalidation_generation), taken **before** the database
437 /// round trip whose rows are being stored. `None` stores unconditionally, and is only correct
438 /// when the value did not come from a read that could have raced a mutation (a test fixture,
439 /// or a synchronous re-population).
440 ///
441 /// # Errors
442 ///
443 /// This method is infallible. The `Result` return type is kept for API compatibility.
444 /// A fenced-out write is **not** an error — the caller already has its rows; they are
445 /// simply not cached. Turning a benign race into a request failure would be worse than
446 /// the staleness this prevents.
447 pub fn put_arc(
448 &self,
449 cache_key: u64,
450 result: Arc<Vec<JsonbValue>>,
451 accessed_views: Vec<String>,
452 ttl_override: Option<u64>,
453 entity_type: Option<&str>,
454 fence: Option<u64>,
455 ) -> Result<()> {
456 if !self.config.enabled {
457 return Ok(());
458 }
459
460 // Cheap pre-check: if an invalidation already landed, do no work at all. The
461 // authoritative check is after registration, below — this one only saves the
462 // serialisation and index churn in the common case.
463 if fence.is_some_and(|snapshot| snapshot != self.invalidation_generation()) {
464 return Ok(());
465 }
466
467 let ttl_seconds = ttl_override.unwrap_or(self.config.ttl_seconds);
468
469 // TTL=0 means "no time-based expiry" — store the entry and rely entirely
470 // on mutation-based invalidation. expire_after_create returns None for
471 // these entries so moka never schedules a timer-wheel eviction.
472
473 // Respect cache_list_queries: a result with more than one row is considered a list.
474 if !self.config.cache_list_queries && result.len() > 1 {
475 return Ok(());
476 }
477
478 // Enforce per-entry size limit: estimate entry size from serialized JSON.
479 if let Some(max_entry) = self.config.max_entry_bytes {
480 let estimated = serde_json::to_vec(&*result).map_or(0, |v| v.len());
481 if estimated > max_entry {
482 return Ok(()); // silently skip oversized entries
483 }
484 }
485
486 // Enforce total cache size limit.
487 if let Some(max_total) = self.config.max_total_bytes {
488 if self.memory_bytes.load(Ordering::Relaxed) >= max_total {
489 return Ok(()); // silently skip when budget is exhausted
490 }
491 }
492
493 // Extract entity refs from ALL rows (not just the first).
494 let entity_refs: Box<[(String, String)]> = if let Some(et) = entity_type {
495 result
496 .iter()
497 .filter_map(|row| {
498 row.as_value()
499 .as_object()?
500 .get("id")?
501 .as_str()
502 .map(|id| (et.to_string(), id.to_string()))
503 })
504 .collect::<Vec<_>>()
505 .into_boxed_slice()
506 } else {
507 Box::default()
508 };
509
510 // Promote owned `String` view names into `ViewName(Arc<str>)` exactly
511 // once. The same Arc is then shared by `view_index` and
512 // `accessed_views` (the slice stored on the cached entry).
513 let accessed_views: Box<[ViewName]> =
514 accessed_views.into_iter().map(ViewName::from).collect();
515
516 // Identity of this entry instance, so the eviction listener can
517 // deregister it without touching a replacement under the same key (#740).
518 let epoch = self.next_epoch.fetch_add(1, Ordering::Relaxed);
519 let reference: IndexRef = (cache_key, epoch);
520
521 // Register in view index.
522 for view in &accessed_views {
523 self.view_index.entry(view.clone()).or_default().insert(reference);
524 }
525
526 // Register ALL entity refs in entity index.
527 for (et, id) in &*entity_refs {
528 self.entity_index
529 .entry(et.clone())
530 .or_default()
531 .entry(id.clone())
532 .or_default()
533 .insert(reference);
534 }
535
536 // ── The fence (#1079) ───────────────────────────────────────────────
537 //
538 // Authoritative check, AFTER registration and BEFORE `store.insert`. The
539 // pre-check at the top of this function is only an optimisation; this is the
540 // one that closes the window, because #740's "register before insert" makes the
541 // key *visible* to a concurrent `invalidate_views` but does not stop the insert
542 // that follows. An invalidation landing between the registrations above and the
543 // insert below collects a key that is not in the store yet — invalidating
544 // nothing — and the insert then lands stale.
545 //
546 // Deregistration is symmetric with the eviction listener's: the same
547 // `(cache_key, epoch)` reference, removed from the same two indexes, so no index
548 // row survives pointing at a key that was never stored. No moka call is made
549 // while a DashMap guard is held — the listener runs synchronously on the calling
550 // thread and re-enters `view_index`, which is why that rule exists.
551 if fence.is_some_and(|snapshot| snapshot != self.invalidation_generation()) {
552 for view in &accessed_views {
553 if let Some(keys) = self.view_index.get(view) {
554 keys.remove(&reference);
555 }
556 }
557 for (et, id) in &*entity_refs {
558 if let Some(by_type) = self.entity_index.get(et) {
559 if let Some(keys) = by_type.get(id) {
560 keys.remove(&reference);
561 }
562 }
563 }
564 return Ok(());
565 }
566
567 let cached = CachedResult {
568 result,
569 accessed_views,
570 cached_at: std::time::SystemTime::now()
571 .duration_since(std::time::UNIX_EPOCH)
572 .map_or(0, |d| d.as_secs()),
573 ttl_seconds,
574 entity_refs,
575 epoch,
576 };
577
578 self.memory_bytes.fetch_add(entry_overhead(), Ordering::Relaxed);
579 // Wrap in Arc so moka's get() costs one atomic increment, not a full clone.
580 self.store.insert(cache_key, Arc::new(cached));
581 self.total_cached.fetch_add(1, Ordering::Relaxed);
582 Ok(())
583 }
584
585 /// Store query result in cache.
586 ///
587 /// If caching is disabled, this is a no-op.
588 ///
589 /// Wraps `result` in an `Arc` and delegates to [`put_arc`](Self::put_arc).
590 /// Prefer [`put_arc`](Self::put_arc) when the caller already holds an `Arc`.
591 ///
592 /// # Arguments
593 ///
594 /// * `cache_key` - Cache key (from `generate_cache_key()`)
595 /// * `result` - Query result to cache
596 /// * `accessed_views` - List of views accessed by this query
597 /// * `ttl_override` - Per-entry TTL in seconds; `None` uses `CacheConfig::ttl_seconds`
598 /// * `entity_type` - Optional GraphQL type name (e.g. `"User"`) for entity-ID indexing. When
599 /// provided, each row's `"id"` field is extracted and stored in `entity_index` so that
600 /// `invalidate_by_entity()` can perform selective eviction.
601 ///
602 /// # Errors
603 ///
604 /// This method is infallible. The `Result` return type is kept for API compatibility.
605 ///
606 /// # Example
607 ///
608 /// ```rust
609 /// use fraiseql_core::cache::{QueryResultCache, CacheConfig};
610 /// use fraiseql_core::db::types::JsonbValue;
611 /// use serde_json::json;
612 ///
613 /// let cache = QueryResultCache::new(CacheConfig::default());
614 ///
615 /// let result = vec![JsonbValue::new(json!({"id": "uuid-1"}))];
616 /// let fence = cache.invalidation_generation();
617 /// cache.put(0xabc123, result, vec!["v_user".to_string()], None, Some("User"), Some(fence))?;
618 /// # Ok::<(), fraiseql_core::error::FraiseQLError>(())
619 /// ```
620 pub fn put(
621 &self,
622 cache_key: u64,
623 result: Vec<JsonbValue>,
624 accessed_views: Vec<String>,
625 ttl_override: Option<u64>,
626 entity_type: Option<&str>,
627 fence: Option<u64>,
628 ) -> Result<()> {
629 self.put_arc(cache_key, Arc::new(result), accessed_views, ttl_override, entity_type, fence)
630 }
631
632 /// Invalidate entries accessing specified views.
633 ///
634 /// Uses the `view_index` for O(k) lookup instead of O(n) full-cache scan.
635 /// Keys accessing multiple views in `views` are deduplicated before invalidation.
636 ///
637 /// # Arguments
638 ///
639 /// * `views` - List of view/table names modified by mutation
640 ///
641 /// # Returns
642 ///
643 /// Number of cache entries invalidated.
644 ///
645 /// # Errors
646 ///
647 /// This method is infallible. The `Result` return type is kept for API compatibility.
648 ///
649 /// # Example
650 ///
651 /// ```rust
652 /// use fraiseql_core::cache::{QueryResultCache, CacheConfig};
653 /// use fraiseql_db::ViewName;
654 ///
655 /// let cache = QueryResultCache::new(CacheConfig::default());
656 ///
657 /// // After createUser mutation
658 /// let invalidated = cache.invalidate_views(&[ViewName::from("v_user")])?;
659 /// println!("Invalidated {} cache entries", invalidated);
660 /// # Ok::<(), fraiseql_core::error::FraiseQLError>(())
661 /// ```
662 pub fn invalidate_views(&self, views: &[ViewName]) -> Result<u64> {
663 if !self.config.enabled {
664 return Ok(0);
665 }
666
667 // Bump BEFORE collecting keys (#1079). A read whose database round trip
668 // straddles this call must observe the move even if it snapshotted a moment
669 // ago — ordering the bump first means the window can only ever be too wide
670 // (a refused cache write), never too narrow (a stale entry).
671 self.invalidation_generation.fetch_add(1, Ordering::Release);
672
673 // Collect keys first (releases DashMap guards) then invalidate.
674 // Moka's eviction listener fires synchronously on the calling thread, so
675 // we must NOT hold any DashMap shard guard when calling store.invalidate() —
676 // the listener itself calls view_index.get() on the same shard, which
677 // would deadlock on a non-re-entrant parking_lot::RwLock.
678 let mut keys_to_invalidate: HashSet<u64> = HashSet::new();
679 for view in views {
680 // ViewName implements Borrow<str>, so DashMap lookup by &str works
681 // without materialising a fresh ViewName.
682 if let Some(keys) = self.view_index.get(view.as_str()) {
683 // Dedup: a query accessing multiple views in `views` would
684 // otherwise be counted and invalidated once per view, and a
685 // key mid-replacement is registered under two epochs.
686 for reference in keys.iter() {
687 keys_to_invalidate.insert(reference.0);
688 }
689 }
690 // Guard dropped here — safe to proceed
691 }
692
693 #[allow(clippy::cast_possible_truncation)]
694 // Reason: entry count never exceeds u64
695 let count = keys_to_invalidate.len() as u64;
696
697 for key in keys_to_invalidate {
698 self.store.invalidate(&key);
699 // Index cleanup handled by eviction listener.
700 }
701
702 self.invalidations.fetch_add(count, Ordering::Relaxed);
703 Ok(count)
704 }
705
706 /// Evict cache entries that contain a specific entity UUID.
707 ///
708 /// Uses the `entity_index` for O(k) lookup. Entries not referencing this
709 /// entity are left untouched.
710 ///
711 /// # Arguments
712 ///
713 /// * `entity_type` - GraphQL type name (e.g. `"User"`)
714 /// * `entity_id` - UUID string of the mutated entity
715 ///
716 /// # Returns
717 ///
718 /// Number of cache entries evicted.
719 ///
720 /// # Errors
721 ///
722 /// This method is infallible. The `Result` return type is kept for API compatibility.
723 pub fn invalidate_by_entity(&self, entity_type: &str, entity_id: &str) -> Result<u64> {
724 if !self.config.enabled {
725 return Ok(0);
726 }
727
728 // Bump BEFORE collecting keys (#1079) — see invalidate_views.
729 self.invalidation_generation.fetch_add(1, Ordering::Release);
730
731 // Short-circuit: if entity_type has no indexed entries, skip the DashMap
732 // lookup entirely. Covers cold-cache and write-heavy workloads where no
733 // reads are cached yet.
734 if !self.entity_index.contains_key(entity_type) {
735 return Ok(0);
736 }
737
738 // Collect keys first (releases DashMap guards) then invalidate.
739 // Moka's eviction listener fires synchronously on the calling thread, so
740 // we must NOT hold any DashMap shard guard when calling store.invalidate() —
741 // the listener itself calls entity_index.get() on the same shard, which
742 // would deadlock on a non-re-entrant parking_lot::RwLock.
743 let keys_to_invalidate: HashSet<u64> = self
744 .entity_index
745 .get(entity_type)
746 .and_then(|by_type| {
747 by_type.get(entity_id).map(|keys| keys.iter().map(|r| r.0).collect())
748 })
749 .unwrap_or_default();
750
751 #[allow(clippy::cast_possible_truncation)]
752 // Reason: entry count never exceeds u64
753 let count = keys_to_invalidate.len() as u64;
754
755 for key in keys_to_invalidate {
756 self.store.invalidate(&key);
757 // Index cleanup handled by eviction listener.
758 }
759
760 self.invalidations.fetch_add(count, Ordering::Relaxed);
761 Ok(count)
762 }
763
764 /// Get cache metrics snapshot.
765 ///
766 /// Returns a consistent snapshot of current counters. Individual fields may
767 /// be updated independently (atomics), so the snapshot is not a single atomic
768 /// transaction, but is accurate enough for monitoring.
769 ///
770 /// # Errors
771 ///
772 /// This method is infallible. The `Result` return type is kept for API compatibility.
773 ///
774 /// # Example
775 ///
776 /// ```rust
777 /// use fraiseql_core::cache::{QueryResultCache, CacheConfig};
778 ///
779 /// let cache = QueryResultCache::new(CacheConfig::default());
780 /// let metrics = cache.metrics()?;
781 ///
782 /// println!("Hit rate: {:.1}%", metrics.hit_rate() * 100.0);
783 /// println!("Size: {} / {} entries", metrics.size, 10_000);
784 /// # Ok::<(), fraiseql_core::error::FraiseQLError>(())
785 /// ```
786 pub fn metrics(&self) -> Result<CacheMetrics> {
787 Ok(CacheMetrics {
788 hits: self.hits.load(Ordering::Relaxed),
789 misses: self.misses.load(Ordering::Relaxed),
790 total_cached: self.total_cached.load(Ordering::Relaxed),
791 invalidations: self.invalidations.load(Ordering::Relaxed),
792 #[allow(clippy::cast_possible_truncation)]
793 // Reason: entry count fits in usize on any 64-bit target
794 size: self.store.entry_count() as usize,
795 memory_bytes: self.memory_bytes.load(Ordering::Relaxed),
796 })
797 }
798
799 /// Clear all cache entries.
800 ///
801 /// Resets the store, reverse indexes, and `memory_bytes` synchronously.
802 /// The eviction listener will still fire asynchronously for each evicted entry,
803 /// but its index-cleanup operations will be no-ops on the already-cleared maps.
804 ///
805 /// # Errors
806 ///
807 /// This method is infallible. The `Result` return type is kept for API compatibility.
808 ///
809 /// # Example
810 ///
811 /// ```rust
812 /// use fraiseql_core::cache::{QueryResultCache, CacheConfig};
813 ///
814 /// let cache = QueryResultCache::new(CacheConfig::default());
815 /// cache.clear()?;
816 /// # Ok::<(), fraiseql_core::error::FraiseQLError>(())
817 /// ```
818 pub fn clear(&self) -> Result<()> {
819 // Bump first (#1079): a read in flight across a clear must not repopulate it.
820 self.invalidation_generation.fetch_add(1, Ordering::Release);
821 self.store.invalidate_all();
822 // Reset indexes and memory counter synchronously — don't rely on the
823 // async eviction listener to do this.
824 self.view_index.clear();
825 self.entity_index.clear();
826 self.memory_bytes.store(0, Ordering::Relaxed);
827 Ok(())
828 }
829
830 /// Flush pending background tasks in the moka store.
831 ///
832 /// Moka applies writes and evictions on a background schedule, so `entry_count()`
833 /// lags: a `put` followed immediately by `metrics()` reports zero entries. Tests
834 /// need this to synchronise before an assertion, and so does the admin stats
835 /// endpoint — an operator asking how many entries are cached wants the settled
836 /// number, not an estimate that reads as "nothing is cached" (#941).
837 ///
838 /// Not on the query path: this walks the pending write buffer.
839 pub fn run_pending_tasks(&self) {
840 self.store.run_pending_tasks();
841 }
842}
843
844impl CacheMetrics {
845 /// Calculate cache hit rate.
846 ///
847 /// Returns ratio of hits to total requests (0.0 to 1.0).
848 ///
849 /// # Returns
850 ///
851 /// - `1.0` if all requests were hits
852 /// - `0.0` if all requests were misses
853 /// - `0.0` if no requests yet
854 ///
855 /// # Example
856 ///
857 /// ```rust
858 /// use fraiseql_core::cache::CacheMetrics;
859 ///
860 /// let metrics = CacheMetrics {
861 /// hits: 80,
862 /// misses: 20,
863 /// total_cached: 100,
864 /// invalidations: 5,
865 /// size: 95,
866 /// memory_bytes: 1_000_000,
867 /// };
868 ///
869 /// assert_eq!(metrics.hit_rate(), 0.8); // 80% hit rate
870 /// ```
871 #[must_use]
872 pub fn hit_rate(&self) -> f64 {
873 let total = self.hits + self.misses;
874 if total == 0 {
875 return 0.0;
876 }
877 #[allow(clippy::cast_precision_loss)]
878 // Reason: hit-rate is a display metric; f64 precision loss on u64 counters is acceptable
879 {
880 self.hits as f64 / total as f64
881 }
882 }
883
884 /// Check if cache is performing well.
885 ///
886 /// Returns `true` if hit rate is above 60% (reasonable threshold).
887 ///
888 /// # Example
889 ///
890 /// ```rust
891 /// use fraiseql_core::cache::CacheMetrics;
892 ///
893 /// let good_metrics = CacheMetrics {
894 /// hits: 80,
895 /// misses: 20,
896 /// total_cached: 100,
897 /// invalidations: 5,
898 /// size: 95,
899 /// memory_bytes: 1_000_000,
900 /// };
901 ///
902 /// assert!(good_metrics.is_healthy()); // 80% > 60%
903 /// ```
904 #[must_use]
905 pub fn is_healthy(&self) -> bool {
906 self.hit_rate() > 0.6
907 }
908}