Skip to main content

amalgam/
cache.rs

1//! The [`Cache`] type — the heart of the library.
2//!
3//! This is the Rust counterpart of `IFusionCache`. The marquee method is
4//! [`Cache::get_or_set`]; its flow encodes cache-stampede protection, fail-safe,
5//! soft/hard timeouts with background completion, eager refresh, adaptive
6//! caching and conditional refresh. See `docs/PARITY.md` for the mapping to
7//! FusionCache.
8
9use std::future::Future;
10use std::sync::{Arc, Weak};
11use std::time::Duration;
12
13use tokio::sync::broadcast;
14use tokio::task::JoinHandle;
15
16use async_trait::async_trait;
17
18use crate::backplane::{Backplane, BackplaneAction, BackplaneMessage};
19use crate::circuit::CircuitBreaker;
20use crate::distributed::{DistributedCache, DistributedEntry, DistributedSerializer};
21use crate::distributed_lock::DistributedLocker;
22use crate::entry::Entry;
23use crate::error::{Error, FactoryError, Result};
24use crate::events::{CacheEvent, CircuitComponent, Events};
25use crate::factory::{FactoryContext, FactoryProduct, StaleInfo};
26use crate::locking::{KeyGuard, KeyedLock};
27use crate::maybe::MaybeValue;
28use crate::memory::MemoryStore;
29use crate::options::{EntryOptions, KeyModifierMode, RemoveByTagBehavior};
30use crate::plugins::{Plugin, PluginHost};
31use crate::recovery::{
32    AutoRecoveryService, RecoveryAction, RecoveryConfig, RecoveryExecutor, RecoveryItem,
33};
34use crate::registry::DefaultEntryOptionsProvider;
35use crate::tags::{Tag, TagRegistry, TagVerdict};
36use crate::time::{Clock, SystemClock, Timeout, Timestamp};
37
38/// A robust, multi-level cache.
39///
40/// `Cache<V>` is generic over a single value type `V`. Unlike FusionCache (which
41/// leans on .NET runtime type information to store heterogeneous values in one
42/// instance), the idiomatic Rust model is one value type per cache — this makes
43/// "wrong type for this key" unrepresentable instead of a run-time downcast.
44/// Use several caches, or a sum type / `serde_json::Value`, for heterogeneous
45/// values.
46///
47/// Cloning a `Cache` is cheap (it shares one underlying instance) and is how you
48/// hand it to other tasks.
49pub struct Cache<V: Clone + Send + Sync + 'static> {
50    inner: Arc<CacheInner<V>>,
51}
52
53struct CacheInner<V: Clone + Send + Sync + 'static> {
54    name: Arc<str>,
55    instance_id: Arc<str>,
56    memory: MemoryStore<V>,
57    locks: Arc<KeyedLock>,
58    tags: Arc<TagRegistry>,
59    events: Events,
60    clock: Arc<dyn Clock>,
61    default_options: EntryOptions,
62    key_prefix: Option<Arc<str>>,
63    remove_by_tag_behavior: RemoveByTagBehavior,
64    distributed: Option<Arc<dyn DistributedCache>>,
65    serializer: Option<Arc<dyn DistributedSerializer<V>>>,
66    backplane: Option<Arc<dyn Backplane>>,
67    distributed_locker: Option<Arc<dyn DistributedLocker>>,
68    circuit_l2: CircuitBreaker,
69    circuit_backplane: CircuitBreaker,
70    plugins: PluginHost,
71    recovery: Option<Arc<AutoRecoveryService>>,
72    default_options_provider: Option<Arc<dyn DefaultEntryOptionsProvider>>,
73    ignore_incoming_backplane: bool,
74    distributed_wire_version: Arc<str>,
75    distributed_key_modifier_mode: KeyModifierMode,
76    disable_tagging: bool,
77    wait_for_initial_backplane_subscribe: bool,
78}
79
80impl<V: Clone + Send + Sync + 'static> Clone for Cache<V> {
81    fn clone(&self) -> Self {
82        Self {
83            inner: Arc::clone(&self.inner),
84        }
85    }
86}
87
88/// The classification of an L1 read after tag markers are applied.
89enum L1Read<V> {
90    Fresh(Entry<V>),
91    Stale(Entry<V>),
92    Miss,
93}
94
95impl<V: Clone + Send + Sync + 'static> Cache<V> {
96    /// Starts building a cache.
97    pub fn builder() -> CacheBuilder<V> {
98        CacheBuilder::new()
99    }
100
101    /// Creates a cache with all default settings and the real system clock.
102    #[must_use]
103    pub fn new() -> Self {
104        CacheBuilder::new().build()
105    }
106
107    /// The cache's name.
108    #[must_use]
109    pub fn name(&self) -> &str {
110        &self.inner.name
111    }
112
113    /// The event hub; subscribe to observe hits, misses, fail-safe activations…
114    #[must_use]
115    pub fn events(&self) -> &Events {
116        &self.inner.events
117    }
118
119    /// A clone of the cache's default entry options, ready to tweak per call.
120    #[must_use]
121    pub fn entry_options(&self) -> EntryOptions {
122        self.inner.default_options.clone()
123    }
124
125    /// Whether the cache establishes its backplane subscription before `build`
126    /// returns (FusionCache `WaitForInitialBackplaneSubscribe`). In `amalgam` this
127    /// is inherently satisfied — see
128    /// [`CacheBuilder::wait_for_initial_backplane_subscribe`].
129    #[must_use]
130    pub fn wait_for_initial_backplane_subscribe(&self) -> bool {
131        self.inner.wait_for_initial_backplane_subscribe
132    }
133
134    // --------------------------------------------------------------- get_or_set
135
136    /// Returns the cached value for `key`, or runs `factory` to produce it.
137    ///
138    /// Uses the cache's default options, no tags and no fail-safe default. See
139    /// [`get_or_set_full`](Self::get_or_set_full) for the complete form.
140    pub async fn get_or_set<F, Fut>(&self, key: impl AsRef<str>, factory: F) -> Result<V>
141    where
142        F: FnOnce(FactoryContext<V>) -> Fut + Send + 'static,
143        Fut: Future<Output = std::result::Result<FactoryProduct<V>, FactoryError>> + Send + 'static,
144    {
145        self.get_or_set_full(key, factory, None, Box::from([]), MaybeValue::none())
146            .await
147    }
148
149    /// Like [`get_or_set`](Self::get_or_set) but with explicit per-call options.
150    pub async fn get_or_set_with<F, Fut>(
151        &self,
152        key: impl AsRef<str>,
153        factory: F,
154        options: EntryOptions,
155    ) -> Result<V>
156    where
157        F: FnOnce(FactoryContext<V>) -> Fut + Send + 'static,
158        Fut: Future<Output = std::result::Result<FactoryProduct<V>, FactoryError>> + Send + 'static,
159    {
160        self.get_or_set_full(
161            key,
162            factory,
163            Some(options),
164            Box::from([]),
165            MaybeValue::none(),
166        )
167        .await
168    }
169
170    /// Returns the cached value for `key`, or stores and returns the constant
171    /// `value` if it is absent.
172    ///
173    /// The constant-value form of [`get_or_set`](Self::get_or_set): it goes
174    /// through the same L1 → L2 → single-flight flow, but the "factory" simply
175    /// yields `value`.
176    pub async fn get_or_set_value(
177        &self,
178        key: impl AsRef<str>,
179        value: V,
180        options: Option<EntryOptions>,
181    ) -> Result<V> {
182        self.get_or_set_full(
183            key,
184            move |ctx| async move { Ok(ctx.value(value)) },
185            options,
186            Box::from([]),
187            MaybeValue::none(),
188        )
189        .await
190    }
191
192    /// The full `get_or_set`: per-call `options`, `tags` for the produced entry,
193    /// and a `fail_safe_default` served as a last resort when the factory fails
194    /// and no stale value exists.
195    #[tracing::instrument(
196        level = "debug",
197        name = "amalgam.get_or_set",
198        skip_all,
199        fields(cache = %self.inner.name, key = key.as_ref())
200    )]
201    pub async fn get_or_set_full<F, Fut>(
202        &self,
203        key: impl AsRef<str>,
204        factory: F,
205        options: Option<EntryOptions>,
206        tags: Box<[Tag]>,
207        fail_safe_default: MaybeValue<V>,
208    ) -> Result<V>
209    where
210        F: FnOnce(FactoryContext<V>) -> Fut + Send + 'static,
211        Fut: Future<Output = std::result::Result<FactoryProduct<V>, FactoryError>> + Send + 'static,
212    {
213        let opts = self.resolve_options(key.as_ref(), options);
214        let full_key = self.full_key(key.as_ref());
215        let now = self.inner.clock.now();
216
217        let mut stale_entry: Option<Entry<V>> = None;
218
219        // 1. Hot path: a fresh L1 entry short-circuits everything.
220        if !opts.skip_memory_read() {
221            match self.read_l1(&full_key, now).await {
222                L1Read::Fresh(entry) => {
223                    if entry.should_eager_refresh(now) {
224                        // Hand the factory to a non-blocking background refresh and
225                        // return the still-fresh value immediately.
226                        self.spawn_eager_refresh(
227                            full_key.clone(),
228                            opts.clone(),
229                            entry.clone(),
230                            factory,
231                        );
232                    }
233                    self.emit(CacheEvent::Hit {
234                        key: full_key,
235                        stale: false,
236                    });
237                    return Ok(entry.value_cloned());
238                }
239                L1Read::Stale(entry) => stale_entry = Some(entry),
240                L1Read::Miss => {}
241            }
242        }
243
244        // 2. Single-flight: acquire the per-key lock.
245        let guard = match self
246            .acquire_lock(&full_key, &opts, stale_entry.as_ref())
247            .await
248        {
249            LockOutcome::Acquired(guard) => guard,
250            LockOutcome::ServedStale(value) => return Ok(value),
251        };
252
253        // 3. Re-check L1 now we hold the lock: another flight may have populated it.
254        let now = self.inner.clock.now();
255        if !opts.skip_memory_read() {
256            match self.read_l1(&full_key, now).await {
257                L1Read::Fresh(entry) => {
258                    self.emit(CacheEvent::Hit {
259                        key: full_key,
260                        stale: false,
261                    });
262                    return Ok(entry.value_cloned());
263                }
264                L1Read::Stale(entry) => stale_entry = Some(entry),
265                L1Read::Miss => {}
266            }
267        }
268
269        // 3b. L2 read-through: a fresh L2 entry populates L1 and returns; a stale
270        // one becomes a (possibly newer) fail-safe fallback.
271        if self.inner.distributed.is_some()
272            && !opts.skip_distributed_read()
273            && !(stale_entry.is_some() && opts.skip_distributed_read_when_stale())
274        {
275            let now = self.inner.clock.now();
276            let l2_has_fallback = stale_entry.is_some() || fail_safe_default.has_value();
277            if let Some(entry) = self
278                .read_l2_guarded(&full_key, now, &opts, l2_has_fallback)
279                .await?
280            {
281                match self.evaluate_tags(entry.meta().created(), entry.meta().tags()) {
282                    TagVerdict::Remove => {
283                        self.remove_l2_guarded(&full_key).await;
284                    }
285                    verdict => {
286                        if !opts.skip_memory_write() {
287                            self.inner
288                                .memory
289                                .insert(Arc::clone(&full_key), entry.clone())
290                                .await;
291                        }
292                        let fresh =
293                            matches!(verdict, TagVerdict::Valid) && entry.freshness(now).is_fresh();
294                        if fresh {
295                            self.emit(CacheEvent::Hit {
296                                key: full_key,
297                                stale: false,
298                            });
299                            return Ok(entry.value_cloned());
300                        }
301                        stale_entry = Some(newer_of(stale_entry, entry));
302                    }
303                }
304            }
305        }
306
307        // 4. Run the factory under soft/hard timeout, with fail-safe behind it.
308        let stale_info = stale_entry.as_ref().map(stale_info_of);
309        let ctx = FactoryContext::new(full_key.clone(), opts.clone(), tags, stale_info);
310        let has_fallback = stale_entry.is_some() || fail_safe_default.has_value();
311        let factory_timeout = opts.appropriate_factory_timeout(has_fallback);
312        let allow_bg = opts.allow_timed_out_factory_background_completion();
313
314        match run_factory(factory, ctx, factory_timeout, allow_bg).await {
315            FactoryRun::Produced(Ok(product)) => {
316                let value = self.store_product(&full_key, product).await;
317                drop(guard);
318                Ok(value)
319            }
320            FactoryRun::Produced(Err(factory_err)) => {
321                tracing::warn!(
322                    cache = %self.inner.name,
323                    key = %full_key,
324                    error = factory_err.message(),
325                    "factory failed"
326                );
327                self.emit(CacheEvent::FactoryError {
328                    key: full_key.clone(),
329                    message: factory_err.message().to_owned(),
330                });
331                let now = self.inner.clock.now();
332                match self
333                    .try_serve_fallback(
334                        &full_key,
335                        &opts,
336                        now,
337                        stale_entry.as_ref(),
338                        &fail_safe_default,
339                    )
340                    .await
341                {
342                    Some(value) => Ok(value),
343                    None => Err(factory_err.into()),
344                }
345            }
346            FactoryRun::TimedOut(handle) => {
347                self.emit(CacheEvent::FactorySyntheticTimeout {
348                    key: full_key.clone(),
349                });
350                if let Some(handle) = handle {
351                    self.spawn_background_completion(full_key.clone(), handle, guard);
352                }
353                let now = self.inner.clock.now();
354                match self
355                    .try_serve_fallback(
356                        &full_key,
357                        &opts,
358                        now,
359                        stale_entry.as_ref(),
360                        &fail_safe_default,
361                    )
362                    .await
363                {
364                    Some(value) => Ok(value),
365                    None => Err(Error::FactoryTimeout {
366                        elapsed: factory_timeout.as_duration().unwrap_or(Duration::ZERO),
367                    }),
368                }
369            }
370        }
371    }
372
373    // ------------------------------------------------------------- simple ops
374
375    /// Writes `value` with the cache's default options and no tags.
376    pub async fn set(&self, key: impl AsRef<str>, value: V) {
377        self.set_full(key, value, None, Box::from([])).await;
378    }
379
380    /// Writes `value` with explicit options and tags.
381    pub async fn set_full(
382        &self,
383        key: impl AsRef<str>,
384        value: V,
385        options: Option<EntryOptions>,
386        tags: Box<[Tag]>,
387    ) {
388        let opts = self.resolve_options(key.as_ref(), options);
389        let full_key = self.full_key(key.as_ref());
390        let now = self.inner.clock.now();
391        let entry = Entry::fresh(value, &opts, now, tags, None, None);
392        self.write_entry(&full_key, &entry, &opts).await;
393        self.emit(CacheEvent::Set { key: full_key });
394    }
395
396    /// Reads a value without ever running a factory. A miss (or a stale entry,
397    /// unless `allow_stale_on_read_only` is set) yields [`MaybeValue::none`].
398    pub async fn try_get(
399        &self,
400        key: impl AsRef<str>,
401        options: Option<EntryOptions>,
402    ) -> MaybeValue<V> {
403        let opts = self.resolve_options(key.as_ref(), options);
404        let full_key = self.full_key(key.as_ref());
405        if opts.skip_memory_read() {
406            self.emit(CacheEvent::Miss { key: full_key });
407            return MaybeValue::none();
408        }
409        let now = self.inner.clock.now();
410        match self.read_l1(&full_key, now).await {
411            L1Read::Fresh(entry) => {
412                self.emit(CacheEvent::Hit {
413                    key: full_key,
414                    stale: false,
415                });
416                MaybeValue::from_value(entry.value_cloned())
417            }
418            L1Read::Stale(entry) if opts.allow_stale_on_read_only() => {
419                self.emit(CacheEvent::Hit {
420                    key: full_key,
421                    stale: true,
422                });
423                MaybeValue::from_value(entry.value_cloned())
424            }
425            _ => {
426                self.emit(CacheEvent::Miss { key: full_key });
427                MaybeValue::none()
428            }
429        }
430    }
431
432    /// Reads a value or returns `default` (never runs a factory).
433    pub async fn get_or_default(
434        &self,
435        key: impl AsRef<str>,
436        default: V,
437        options: Option<EntryOptions>,
438    ) -> V {
439        self.try_get(key, options).await.value_or(default)
440    }
441
442    /// Removes an entry from L1 and L2 and tells peers to evict it.
443    pub async fn remove(&self, key: impl AsRef<str>) {
444        let full_key = self.full_key(key.as_ref());
445        let opts = self.inner.default_options.clone();
446        self.inner.memory.remove(&full_key).await;
447        self.remove_l2_guarded(&full_key).await;
448        self.publish_guarded(BackplaneAction::Remove, &full_key, &opts)
449            .await;
450        self.emit(CacheEvent::Remove { key: full_key });
451    }
452
453    /// Logically expires an entry: it is no longer fresh, but fail-safe can still
454    /// serve it as a stale fallback. Propagated to L2 and peers.
455    pub async fn expire(&self, key: impl AsRef<str>) {
456        let full_key = self.full_key(key.as_ref());
457        let now = self.inner.clock.now();
458        let opts = self.inner.default_options.clone();
459        if let Some(entry) = self.inner.memory.get(&full_key).await {
460            let expired = entry.with_logical_expiration(now);
461            self.inner
462                .memory
463                .insert(Arc::clone(&full_key), expired.clone())
464                .await;
465            self.write_l2_guarded(&full_key, &expired, &opts).await;
466        }
467        self.publish_guarded(BackplaneAction::Expire, &full_key, &opts)
468            .await;
469        self.emit(CacheEvent::Expire { key: full_key });
470    }
471
472    /// Invalidates every entry carrying `tag` (lazily — entries are dropped on
473    /// their next read). No-op for a blank tag.
474    pub async fn remove_by_tag(&self, tag: impl AsRef<str>) {
475        if self.inner.disable_tagging {
476            tracing::warn!(
477                cache = %self.inner.name,
478                tag = tag.as_ref(),
479                "remove_by_tag ignored: tagging is disabled (DisableTagging)"
480            );
481            return;
482        }
483        if let Some(t) = Tag::new(tag.as_ref()) {
484            let now = self.inner.clock.now();
485            self.inner.tags.mark_tag(t, now);
486            let marker: Arc<str> = Arc::from(format!("{TAG_MARKER_PREFIX}{}", tag.as_ref()));
487            self.publish_marker(marker, now).await;
488            self.emit(CacheEvent::RemoveByTag {
489                tag: tag.as_ref().to_owned(),
490            });
491        }
492    }
493
494    /// Invalidates every entry carrying any of `tags`.
495    pub async fn remove_by_tags<I, S>(&self, tags: I)
496    where
497        I: IntoIterator<Item = S>,
498        S: AsRef<str>,
499    {
500        for tag in tags {
501            self.remove_by_tag(tag).await;
502        }
503    }
504
505    /// Clears the whole cache.
506    ///
507    /// With `allow_fail_safe = true` (the FusionCache default) every entry is
508    /// logically expired so fail-safe can still serve stale values; with `false`
509    /// they are hard-removed.
510    pub async fn clear(&self, allow_fail_safe: bool) {
511        if self.inner.disable_tagging {
512            tracing::warn!(
513                cache = %self.inner.name,
514                "clear ignored: tagging is disabled (DisableTagging); clear relies on tag markers"
515            );
516            return;
517        }
518        let now = self.inner.clock.now();
519        let marker: Arc<str> = if allow_fail_safe {
520            self.inner.tags.mark_clear_expire(now);
521            Arc::from(CLEAR_EXPIRE_KEY)
522        } else {
523            self.inner.tags.mark_clear_remove(now);
524            self.inner.memory.invalidate_all();
525            Arc::from(CLEAR_REMOVE_KEY)
526        };
527        self.publish_marker(marker, now).await;
528        self.emit(CacheEvent::Clear);
529    }
530
531    /// Runs the L1 backend's pending maintenance (eviction, expiry). Primarily
532    /// for deterministic tests.
533    pub async fn run_pending_tasks(&self) {
534        self.inner.memory.run_pending_tasks().await;
535    }
536
537    // ----------------------------------------------------------------- internals
538
539    fn full_key(&self, key: &str) -> Arc<str> {
540        match &self.inner.key_prefix {
541            Some(prefix) => Arc::from(format!("{prefix}{key}")),
542            None => Arc::from(key),
543        }
544    }
545
546    fn emit(&self, event: CacheEvent) {
547        if !self.inner.plugins.is_empty() {
548            self.inner.plugins.notify(&event);
549        }
550        self.inner.events.emit(event);
551    }
552
553    /// Resolves the options for an operation: explicit per-call options, else a
554    /// dynamic provider's options for this key, else the static defaults.
555    fn resolve_options(&self, key: &str, options: Option<EntryOptions>) -> EntryOptions {
556        if let Some(options) = options {
557            return options;
558        }
559        if let Some(provider) = &self.inner.default_options_provider
560            && let Some(options) = provider.options_for(key)
561        {
562            return options;
563        }
564        self.inner.default_options.clone()
565    }
566
567    /// Evaluates tag/clear markers for an entry, honouring the cache-wide
568    /// `disable_tagging` switch: when tagging is disabled every entry is `Valid`
569    /// (the marker registry is never consulted).
570    fn evaluate_tags(&self, created: Timestamp, tags: &[Tag]) -> TagVerdict {
571        if self.inner.disable_tagging {
572            return TagVerdict::Valid;
573        }
574        self.inner
575            .tags
576            .evaluate(created, tags, self.inner.remove_by_tag_behavior)
577    }
578
579    /// Reads L1, applies tag markers, and classifies the result.
580    async fn read_l1(&self, full_key: &str, now: Timestamp) -> L1Read<V> {
581        let Some(entry) = self.inner.memory.get(full_key).await else {
582            return L1Read::Miss;
583        };
584        match self.evaluate_tags(entry.meta().created(), entry.meta().tags()) {
585            TagVerdict::Remove => {
586                self.inner.memory.remove(full_key).await;
587                L1Read::Miss
588            }
589            TagVerdict::Expire => L1Read::Stale(entry),
590            TagVerdict::Valid => {
591                if entry.freshness(now).is_fresh() {
592                    L1Read::Fresh(entry)
593                } else {
594                    L1Read::Stale(entry)
595                }
596            }
597        }
598    }
599
600    async fn acquire_lock(
601        &self,
602        full_key: &str,
603        opts: &EntryOptions,
604        stale_entry: Option<&Entry<V>>,
605    ) -> LockOutcome<V> {
606        let local = match opts.memory_lock_timeout() {
607            Timeout::Infinite => self.inner.locks.lock(full_key).await,
608            Timeout::After(d) => {
609                match tokio::time::timeout(d, self.inner.locks.lock(full_key)).await {
610                    Ok(guard) => guard,
611                    Err(_) => {
612                        // Lock timed out: with fail-safe and a stale value, return it
613                        // rather than queue behind the lock holder.
614                        if opts.is_fail_safe_enabled()
615                            && let Some(stale) = stale_entry
616                        {
617                            self.emit(CacheEvent::Hit {
618                                key: Arc::from(full_key),
619                                stale: true,
620                            });
621                            return LockOutcome::ServedStale(stale.value_cloned());
622                        }
623                        // Otherwise the lock is best-effort: wait for it.
624                        self.inner.locks.lock(full_key).await
625                    }
626                }
627            }
628        };
629        let distributed = self.acquire_distributed_lock(full_key, opts).await;
630        LockOutcome::Acquired(LockGuard {
631            _local: local,
632            _distributed: distributed,
633        })
634    }
635
636    /// Acquires the cross-node distributed lock (if a locker is configured and not
637    /// skipped). Best-effort: a timeout or error proceeds without it.
638    async fn acquire_distributed_lock(
639        &self,
640        full_key: &str,
641        opts: &EntryOptions,
642    ) -> Option<DistributedReleaseGuard> {
643        if opts.skip_distributed_locker() {
644            return None;
645        }
646        let locker = self.inner.distributed_locker.as_ref()?;
647        let lock_key = format!("amalgam:lock:{full_key}");
648        match locker
649            .acquire(
650                &lock_key,
651                opts.physical_ttl(),
652                opts.distributed_lock_timeout(),
653            )
654            .await
655        {
656            Ok(Some(token)) => Some(DistributedReleaseGuard {
657                locker: Arc::clone(locker),
658                key: Arc::from(lock_key),
659                token,
660            }),
661            _ => None,
662        }
663    }
664
665    /// Stores a freshly-produced (or `NotModified`-reused) value as a fresh entry,
666    /// fires the appropriate events, and returns the value.
667    async fn store_product(&self, full_key: &Arc<str>, product: FactoryProduct<V>) -> V {
668        let now = self.inner.clock.now();
669        let reused = product.reused_stale;
670        let value = product.value.clone();
671        let opts = product.options;
672        let entry = Entry::fresh(
673            product.value,
674            &opts,
675            now,
676            product.tags,
677            product.etag,
678            product.last_modified,
679        );
680        self.write_entry(full_key, &entry, &opts).await;
681        if !reused {
682            self.emit(CacheEvent::FactorySuccess {
683                key: Arc::clone(full_key),
684            });
685        }
686        self.emit(CacheEvent::Set {
687            key: Arc::clone(full_key),
688        });
689        value
690    }
691
692    /// Write-through: stores `entry` into L1 and (if configured) L2, then
693    /// publishes a backplane `Set` so peers drop their local copy.
694    async fn write_entry(&self, full_key: &Arc<str>, entry: &Entry<V>, opts: &EntryOptions) {
695        if !opts.skip_memory_write() {
696            self.inner
697                .memory
698                .insert(Arc::clone(full_key), entry.clone())
699                .await;
700        }
701        if !opts.skip_distributed_write() {
702            self.write_l2_guarded(full_key, entry, opts).await;
703        }
704        if !opts.skip_backplane_notifications() {
705            self.publish_guarded(BackplaneAction::Set, full_key, opts)
706                .await;
707        }
708    }
709
710    /// L2 write, gated by the circuit breaker, with event emission and
711    /// auto-recovery enqueue on failure.
712    async fn write_l2_guarded(&self, full_key: &Arc<str>, entry: &Entry<V>, opts: &EntryOptions) {
713        if self.inner.distributed.is_none() {
714            return;
715        }
716        let now = self.inner.clock.now();
717        if !self.inner.circuit_l2.is_closed(now) {
718            self.enqueue_recovery(full_key, RecoveryAction::Set, now);
719            return;
720        }
721        let ttl = opts.distributed_physical_ttl();
722        if opts.allow_background_distributed_operations() {
723            // Fire-and-forget: don't make the caller wait on L2.
724            let this = self.clone();
725            let full_key = Arc::clone(full_key);
726            let entry = entry.clone();
727            tokio::spawn(async move {
728                this.do_l2_write(&full_key, &entry, ttl).await;
729            });
730        } else {
731            self.do_l2_write(full_key, entry, ttl).await;
732        }
733    }
734
735    async fn do_l2_write(&self, full_key: &Arc<str>, entry: &Entry<V>, ttl: Duration) {
736        match self.inner.l2_write(full_key, entry, ttl).await {
737            Ok(()) => self.close_circuit_l2(),
738            Err(err) => {
739                self.on_l2_error(full_key, &err);
740                self.enqueue_recovery(full_key, RecoveryAction::Set, self.inner.clock.now());
741            }
742        }
743    }
744
745    /// L2 read, gated by the circuit breaker and the appropriate distributed
746    /// timeout (soft when fail-safe + a fallback exists, otherwise hard).
747    /// Returns `Err` only when `rethrow_distributed_exceptions` is set.
748    async fn read_l2_guarded(
749        &self,
750        full_key: &Arc<str>,
751        now: Timestamp,
752        opts: &EntryOptions,
753        has_fallback: bool,
754    ) -> Result<Option<Entry<V>>> {
755        if self.inner.distributed.is_none() || !self.inner.circuit_l2.is_closed(now) {
756            return Ok(None);
757        }
758        let read = self.inner.l2_read(full_key, now);
759        let timeout = opts.appropriate_distributed_timeout(has_fallback);
760        let result = match timeout {
761            Timeout::After(d) => match tokio::time::timeout(d, read).await {
762                Ok(r) => r,
763                Err(_) => {
764                    // A *soft* distributed timeout (fail-safe on + a fallback
765                    // exists) is an expected, benign bail-out to the fallback: it
766                    // must not mark L2 unhealthy. Only a *hard* timeout is a
767                    // genuine L2 failure that trips the circuit breaker.
768                    if opts.is_fail_safe_enabled()
769                        && has_fallback
770                        && timeout == opts.distributed_soft_timeout()
771                    {
772                        return Ok(None);
773                    }
774                    Err(Error::Distributed("l2 read timed out".to_owned()))
775                }
776            },
777            Timeout::Infinite => read.await,
778        };
779        match result {
780            Ok(entry) => {
781                self.close_circuit_l2();
782                Ok(entry)
783            }
784            Err(err) => {
785                self.on_l2_error(full_key, &err);
786                // Serialization/deserialization failures honour
787                // `rethrow_serialization_exceptions`; transport failures honour
788                // `rethrow_distributed_exceptions`. Otherwise an L2 hiccup is a
789                // miss, not an error.
790                let rethrow = match err {
791                    Error::Serialization(_) | Error::Deserialization(_) => {
792                        opts.rethrow_serialization_exceptions()
793                    }
794                    _ => opts.rethrow_distributed_exceptions(),
795                };
796                if rethrow { Err(err) } else { Ok(None) }
797            }
798        }
799    }
800
801    /// L2 remove, gated by the circuit breaker, with auto-recovery on failure.
802    async fn remove_l2_guarded(&self, full_key: &Arc<str>) {
803        if self.inner.distributed.is_none() {
804            return;
805        }
806        let now = self.inner.clock.now();
807        if !self.inner.circuit_l2.is_closed(now) {
808            self.enqueue_recovery(full_key, RecoveryAction::Remove, now);
809            return;
810        }
811        match self.inner.l2_remove(full_key).await {
812            Ok(()) => self.close_circuit_l2(),
813            Err(err) => {
814                self.on_l2_error(full_key, &err);
815                self.enqueue_recovery(full_key, RecoveryAction::Remove, now);
816            }
817        }
818    }
819
820    fn on_l2_error(&self, full_key: &Arc<str>, err: &Error) {
821        match err {
822            Error::Serialization(message) => self.emit(CacheEvent::SerializationError {
823                key: Arc::clone(full_key),
824                message: message.clone(),
825            }),
826            Error::Deserialization(message) => self.emit(CacheEvent::DeserializationError {
827                key: Arc::clone(full_key),
828                message: message.clone(),
829            }),
830            _ => {}
831        }
832        if self.inner.circuit_l2.trip(self.inner.clock.now()) {
833            self.emit(CacheEvent::CircuitBreakerChange {
834                component: CircuitComponent::Distributed,
835                closed: false,
836            });
837        }
838    }
839
840    fn close_circuit_l2(&self) {
841        if self.inner.circuit_l2.close() {
842            self.emit(CacheEvent::CircuitBreakerChange {
843                component: CircuitComponent::Distributed,
844                closed: true,
845            });
846        }
847    }
848
849    fn close_circuit_backplane(&self) {
850        if self.inner.circuit_backplane.close() {
851            self.emit(CacheEvent::CircuitBreakerChange {
852                component: CircuitComponent::Backplane,
853                closed: true,
854            });
855        }
856    }
857
858    fn enqueue_recovery(&self, full_key: &Arc<str>, action: RecoveryAction, now: Timestamp) {
859        if let Some(recovery) = &self.inner.recovery {
860            recovery.enqueue(RecoveryItem {
861                key: Arc::clone(full_key),
862                action,
863                timestamp: now,
864                expires_at: now.saturating_add(RECOVERY_ITEM_TTL),
865                remaining_retries: None,
866            });
867        }
868    }
869
870    /// Publishes a backplane notification, gated by the circuit breaker. Honours
871    /// `allow_background_backplane_operations` (fire-and-forget vs awaited).
872    async fn publish_guarded(
873        &self,
874        action: BackplaneAction,
875        full_key: &Arc<str>,
876        opts: &EntryOptions,
877    ) {
878        if self.inner.backplane.is_none() {
879            return;
880        }
881        let now = self.inner.clock.now();
882        if !self.inner.circuit_backplane.is_closed(now) {
883            self.enqueue_recovery(full_key, recovery_action_of(action), now);
884            return;
885        }
886        if opts.allow_background_backplane_operations() {
887            let this = self.clone();
888            let key = Arc::clone(full_key);
889            tokio::spawn(async move {
890                this.do_publish(action, &key).await;
891            });
892        } else {
893            self.do_publish(action, full_key).await;
894        }
895    }
896
897    async fn do_publish(&self, action: BackplaneAction, full_key: &Arc<str>) {
898        let now = self.inner.clock.now();
899        match self.inner.backplane_send(action, full_key, now).await {
900            Ok(()) => {
901                self.close_circuit_backplane();
902                self.emit(CacheEvent::MessagePublished {
903                    key: Arc::clone(full_key),
904                });
905            }
906            Err(_) => {
907                if self.inner.circuit_backplane.trip(now) {
908                    self.emit(CacheEvent::CircuitBreakerChange {
909                        component: CircuitComponent::Backplane,
910                        closed: false,
911                    });
912                }
913                self.enqueue_recovery(full_key, recovery_action_of(action), now);
914            }
915        }
916    }
917
918    /// Publishes a reserved-key tag/clear marker so peers update their tag
919    /// registry (best-effort, awaited).
920    async fn publish_marker(&self, key: Arc<str>, now: Timestamp) {
921        if self.inner.backplane.is_none() {
922            return;
923        }
924        if let Ok(()) = self
925            .inner
926            .backplane_send(BackplaneAction::Set, &key, now)
927            .await
928        {
929            self.emit(CacheEvent::MessagePublished { key });
930        }
931    }
932
933    /// Applies an incoming backplane notification to the local L1 / tag registry.
934    async fn apply_backplane(&self, message: BackplaneMessage) {
935        if self.inner.ignore_incoming_backplane {
936            return;
937        }
938        self.emit(CacheEvent::MessageReceived {
939            key: Arc::clone(&message.key),
940        });
941        // A received message proves the backplane is healthy.
942        self.close_circuit_backplane();
943
944        // Tag / clear markers ride on `Set` messages with reserved keys.
945        if message.action == BackplaneAction::Set {
946            if let Some(tag) = message.key.strip_prefix(TAG_MARKER_PREFIX) {
947                if let Some(tag) = Tag::new(tag) {
948                    self.inner.tags.mark_tag(tag, message.timestamp);
949                }
950                return;
951            }
952            if &*message.key == CLEAR_EXPIRE_KEY {
953                self.inner.tags.mark_clear_expire(message.timestamp);
954                return;
955            }
956            if &*message.key == CLEAR_REMOVE_KEY {
957                self.inner.tags.mark_clear_remove(message.timestamp);
958                self.inner.memory.invalidate_all();
959                return;
960            }
961        }
962
963        match message.action {
964            BackplaneAction::Remove => {
965                self.inner.memory.remove(&message.key).await;
966            }
967            BackplaneAction::Expire => {
968                if let Some(entry) = self.inner.memory.get(&message.key).await {
969                    let expired = entry.with_logical_expiration(message.timestamp);
970                    self.inner
971                        .memory
972                        .insert(Arc::clone(&message.key), expired)
973                        .await;
974                }
975            }
976            BackplaneAction::Set => {
977                // FusionCache "passive update": a node that already holds the key
978                // in L1 eagerly refreshes it from L2; nodes that don't hold it do
979                // nothing. The refresh runs off the listener so a slow L2 never
980                // stalls processing of other backplane messages.
981                if self.inner.memory.get(&message.key).await.is_some() {
982                    let cache = self.clone();
983                    let key = Arc::clone(&message.key);
984                    tokio::spawn(async move {
985                        cache.refresh_l1_from_l2(&key).await;
986                    });
987                }
988            }
989        }
990    }
991
992    /// Eagerly refreshes a key's L1 entry from L2 — FusionCache's passive update
993    /// on a backplane `Set`. On an L2 miss or error, or when a tag/clear marker
994    /// invalidates the value, the stale L1 copy is dropped so the next read
995    /// re-pulls or re-runs the factory. Honours the same tag evaluation as the
996    /// `get_or_set` L2 read-through.
997    async fn refresh_l1_from_l2(&self, full_key: &Arc<str>) {
998        if self.inner.distributed.is_none() {
999            self.inner.memory.remove(full_key).await;
1000            return;
1001        }
1002        let now = self.inner.clock.now();
1003        let opts = self.inner.default_options.clone();
1004        match self.read_l2_guarded(full_key, now, &opts, false).await {
1005            Ok(Some(entry)) => {
1006                match self.evaluate_tags(entry.meta().created(), entry.meta().tags()) {
1007                    TagVerdict::Remove => {
1008                        self.inner.memory.remove(full_key).await;
1009                    }
1010                    _ => {
1011                        self.inner.memory.insert(Arc::clone(full_key), entry).await;
1012                    }
1013                }
1014            }
1015            Ok(None) | Err(_) => {
1016                self.inner.memory.remove(full_key).await;
1017            }
1018        }
1019    }
1020
1021    /// Spawns the background listener that applies remote backplane messages.
1022    /// Requires a tokio runtime; only spawned when a backplane is configured.
1023    fn spawn_backplane_listener(&self) {
1024        let Some(backplane) = &self.inner.backplane else {
1025            return;
1026        };
1027        let mut receiver = backplane.subscribe();
1028        let weak: Weak<CacheInner<V>> = Arc::downgrade(&self.inner);
1029        let instance_id = Arc::clone(&self.inner.instance_id);
1030        tokio::spawn(async move {
1031            loop {
1032                match receiver.recv().await {
1033                    Ok(message) => {
1034                        if message.source_id == instance_id {
1035                            continue; // ignore our own notifications
1036                        }
1037                        match weak.upgrade() {
1038                            Some(inner) => Cache { inner }.apply_backplane(message).await,
1039                            None => break, // the cache was dropped
1040                        }
1041                    }
1042                    Err(broadcast::error::RecvError::Closed) => break,
1043                    Err(broadcast::error::RecvError::Lagged(_)) => continue,
1044                }
1045            }
1046        });
1047    }
1048
1049    /// Attempts to serve a stale value (fail-safe). Returns the served value, or
1050    /// `None` if fail-safe is disabled or no fallback exists.
1051    async fn try_serve_fallback(
1052        &self,
1053        full_key: &Arc<str>,
1054        opts: &EntryOptions,
1055        now: Timestamp,
1056        stale_entry: Option<&Entry<V>>,
1057        fail_safe_default: &MaybeValue<V>,
1058    ) -> Option<V> {
1059        if !opts.is_fail_safe_enabled() {
1060            return None;
1061        }
1062        if let Some(stale) = stale_entry
1063            && let Some(throttled) = Entry::throttled(stale, opts, now)
1064        {
1065            let value = throttled.value_cloned();
1066            if !opts.skip_memory_write() {
1067                self.inner
1068                    .memory
1069                    .insert(Arc::clone(full_key), throttled)
1070                    .await;
1071            }
1072            self.emit_fail_safe(full_key);
1073            return Some(value);
1074        }
1075        if let Some(default) = fail_safe_default.value() {
1076            let value = default.clone();
1077            if !opts.skip_memory_write() {
1078                let entry = Entry::from_fail_safe_default(value.clone(), opts, now);
1079                self.inner.memory.insert(Arc::clone(full_key), entry).await;
1080            }
1081            self.emit_fail_safe(full_key);
1082            return Some(value);
1083        }
1084        None
1085    }
1086
1087    fn emit_fail_safe(&self, full_key: &Arc<str>) {
1088        tracing::debug!(
1089            cache = %self.inner.name,
1090            key = %full_key,
1091            "fail-safe activated; serving stale value"
1092        );
1093        self.emit(CacheEvent::FailSafeActivate {
1094            key: Arc::clone(full_key),
1095        });
1096        self.emit(CacheEvent::Hit {
1097            key: Arc::clone(full_key),
1098            stale: true,
1099        });
1100    }
1101
1102    fn spawn_background_completion(
1103        &self,
1104        full_key: Arc<str>,
1105        handle: JoinHandle<std::result::Result<FactoryProduct<V>, FactoryError>>,
1106        guard: LockGuard,
1107    ) {
1108        let this = self.clone();
1109        tokio::spawn(async move {
1110            // Hold the single-flight lock until the background factory resolves.
1111            let _guard = guard;
1112            match handle.await {
1113                Ok(Ok(product)) => {
1114                    this.store_product(&full_key, product).await;
1115                    this.emit(CacheEvent::BackgroundFactorySuccess { key: full_key });
1116                }
1117                Ok(Err(factory_err)) => {
1118                    // Per FusionCache, a background failure does NOT re-activate
1119                    // fail-safe; the throttled stale value already returned stands.
1120                    this.emit(CacheEvent::BackgroundFactoryError {
1121                        key: full_key,
1122                        message: factory_err.message().to_owned(),
1123                    });
1124                }
1125                Err(_) => {
1126                    this.emit(CacheEvent::BackgroundFactoryError {
1127                        key: full_key,
1128                        message: "factory task panicked".to_owned(),
1129                    });
1130                }
1131            }
1132        });
1133    }
1134
1135    fn spawn_eager_refresh<F, Fut>(
1136        &self,
1137        full_key: Arc<str>,
1138        opts: EntryOptions,
1139        current: Entry<V>,
1140        factory: F,
1141    ) where
1142        F: FnOnce(FactoryContext<V>) -> Fut + Send + 'static,
1143        Fut: Future<Output = std::result::Result<FactoryProduct<V>, FactoryError>> + Send + 'static,
1144    {
1145        // Non-blocking: if another caller already holds the key, skip silently.
1146        let Some(guard) = self.inner.locks.try_lock(&full_key) else {
1147            return;
1148        };
1149        self.emit(CacheEvent::EagerRefresh {
1150            key: Arc::clone(&full_key),
1151        });
1152        let this = self.clone();
1153        tokio::spawn(async move {
1154            let _guard = guard;
1155            let ctx = FactoryContext::new(
1156                Arc::clone(&full_key),
1157                opts,
1158                current.meta().tags().to_vec().into_boxed_slice(),
1159                Some(stale_info_of(&current)),
1160            );
1161            match factory(ctx).await {
1162                Ok(product) => {
1163                    this.store_product(&full_key, product).await;
1164                    this.emit(CacheEvent::BackgroundFactorySuccess { key: full_key });
1165                }
1166                Err(factory_err) => {
1167                    this.emit(CacheEvent::BackgroundFactoryError {
1168                        key: full_key,
1169                        message: factory_err.message().to_owned(),
1170                    });
1171                }
1172            }
1173        });
1174    }
1175}
1176
1177impl<V: Clone + Send + Sync + 'static> Default for Cache<V> {
1178    fn default() -> Self {
1179        Self::new()
1180    }
1181}
1182
1183impl<V: Clone + Send + Sync + 'static> CacheInner<V> {
1184    /// The L2 storage key (wire-version-prefixed so cache versions can share L2).
1185    fn l2_key(&self, full_key: &str) -> String {
1186        match self.distributed_key_modifier_mode {
1187            KeyModifierMode::Prefix => format!("{}:{}", self.distributed_wire_version, full_key),
1188            KeyModifierMode::Suffix => format!("{}:{}", full_key, self.distributed_wire_version),
1189            KeyModifierMode::None => full_key.to_owned(),
1190        }
1191    }
1192
1193    /// Serializes and writes an entry to L2 with the given physical TTL. `Err` on
1194    /// serialize/backend failure.
1195    async fn l2_write(&self, full_key: &str, entry: &Entry<V>, ttl: Duration) -> Result<()> {
1196        let (Some(l2), Some(serializer)) = (&self.distributed, &self.serializer) else {
1197            return Ok(());
1198        };
1199        let dist = DistributedEntry::from_entry(entry);
1200        let bytes = serializer.serialize(&dist)?;
1201        l2.set(&self.l2_key(full_key), bytes, Some(ttl)).await
1202    }
1203
1204    /// Reads and deserializes an entry from L2, rehydrated relative to `now`.
1205    async fn l2_read(&self, full_key: &str, now: Timestamp) -> Result<Option<Entry<V>>> {
1206        let (Some(l2), Some(serializer)) = (&self.distributed, &self.serializer) else {
1207            return Ok(None);
1208        };
1209        let Some(bytes) = l2.get(&self.l2_key(full_key)).await? else {
1210            return Ok(None);
1211        };
1212        let dist = serializer.deserialize(&bytes)?;
1213        Ok(Some(dist.into_entry(now)))
1214    }
1215
1216    /// Removes an entry from L2.
1217    async fn l2_remove(&self, full_key: &str) -> Result<()> {
1218        if let Some(l2) = &self.distributed {
1219            l2.remove(&self.l2_key(full_key)).await
1220        } else {
1221            Ok(())
1222        }
1223    }
1224
1225    /// Publishes a backplane message (awaited). No-op without a backplane.
1226    async fn backplane_send(
1227        &self,
1228        action: BackplaneAction,
1229        full_key: &str,
1230        ts: Timestamp,
1231    ) -> Result<()> {
1232        if let Some(backplane) = &self.backplane {
1233            backplane
1234                .publish(BackplaneMessage {
1235                    source_id: Arc::clone(&self.instance_id),
1236                    timestamp: ts,
1237                    action,
1238                    key: Arc::from(full_key),
1239                })
1240                .await
1241        } else {
1242            Ok(())
1243        }
1244    }
1245}
1246
1247#[async_trait]
1248impl<V: Clone + Send + Sync + 'static> RecoveryExecutor for CacheInner<V> {
1249    async fn replay(&self, item: &RecoveryItem) -> Result<()> {
1250        match item.action {
1251            RecoveryAction::Set => {
1252                if let Some(entry) = self.memory.get(&item.key).await {
1253                    let ttl = entry.backend_ttl();
1254                    self.l2_write(&item.key, &entry, ttl).await?;
1255                }
1256                self.backplane_send(BackplaneAction::Set, &item.key, item.timestamp)
1257                    .await?;
1258            }
1259            RecoveryAction::Remove => {
1260                self.l2_remove(&item.key).await?;
1261                self.backplane_send(BackplaneAction::Remove, &item.key, item.timestamp)
1262                    .await?;
1263            }
1264            RecoveryAction::Expire => {
1265                self.backplane_send(BackplaneAction::Expire, &item.key, item.timestamp)
1266                    .await?;
1267            }
1268        }
1269        Ok(())
1270    }
1271}
1272
1273/// Reserved key prefix for tag-invalidation markers propagated over the backplane.
1274const TAG_MARKER_PREFIX: &str = "__amalgam:t:";
1275/// Reserved key for a "clear (expire all)" marker.
1276const CLEAR_EXPIRE_KEY: &str = "__amalgam:clear:expire";
1277/// Reserved key for a "clear (remove all)" marker.
1278const CLEAR_REMOVE_KEY: &str = "__amalgam:clear:remove";
1279/// How long a queued auto-recovery item remains eligible for replay.
1280const RECOVERY_ITEM_TTL: Duration = Duration::from_secs(600);
1281
1282/// Maps a backplane action to the equivalent recovery action.
1283fn recovery_action_of(action: BackplaneAction) -> RecoveryAction {
1284    match action {
1285        BackplaneAction::Set => RecoveryAction::Set,
1286        BackplaneAction::Remove => RecoveryAction::Remove,
1287        BackplaneAction::Expire => RecoveryAction::Expire,
1288    }
1289}
1290
1291/// The combined single-flight guard: the local lock plus an optional cross-node
1292/// distributed lock. Dropping it releases both.
1293struct LockGuard {
1294    _local: KeyGuard,
1295    _distributed: Option<DistributedReleaseGuard>,
1296}
1297
1298/// Releases a held distributed lock when dropped (best-effort, in the background).
1299struct DistributedReleaseGuard {
1300    locker: Arc<dyn DistributedLocker>,
1301    key: Arc<str>,
1302    token: String,
1303}
1304
1305impl Drop for DistributedReleaseGuard {
1306    fn drop(&mut self) {
1307        let locker = Arc::clone(&self.locker);
1308        let key = Arc::clone(&self.key);
1309        let token = std::mem::take(&mut self.token);
1310        tokio::spawn(async move {
1311            let _ = locker.release(&key, &token).await;
1312        });
1313    }
1314}
1315
1316/// The result of trying to acquire the single-flight lock.
1317enum LockOutcome<V> {
1318    /// The lock was acquired; proceed to read L2 / run the factory.
1319    Acquired(LockGuard),
1320    /// The lock timed out but a stale value was served instead.
1321    ServedStale(V),
1322}
1323
1324/// The result of running a factory under a (possibly infinite) timeout.
1325///
1326/// This is a short-lived local return value (one per factory call), never stored
1327/// in bulk, so the size gap between variants is irrelevant — boxing the common
1328/// `Produced` path would only add a hot-path allocation.
1329#[allow(clippy::large_enum_variant)]
1330enum FactoryRun<V> {
1331    /// The factory completed (successfully or with a failure) within the timeout.
1332    Produced(std::result::Result<FactoryProduct<V>, FactoryError>),
1333    /// The factory exceeded its timeout. `Some(handle)` when it is still running
1334    /// in the background (to be completed later); `None` when it was aborted.
1335    TimedOut(Option<JoinHandle<std::result::Result<FactoryProduct<V>, FactoryError>>>),
1336}
1337
1338/// Runs `factory` against `ctx`, enforcing `timeout`.
1339///
1340/// For a finite timeout the factory is spawned so it can outlive the timeout and
1341/// finish in the background (when `allow_background` is set); for an infinite
1342/// timeout it runs inline.
1343async fn run_factory<V, F, Fut>(
1344    factory: F,
1345    ctx: FactoryContext<V>,
1346    timeout: Timeout,
1347    allow_background: bool,
1348) -> FactoryRun<V>
1349where
1350    V: Send + 'static,
1351    F: FnOnce(FactoryContext<V>) -> Fut + Send + 'static,
1352    Fut: Future<Output = std::result::Result<FactoryProduct<V>, FactoryError>> + Send + 'static,
1353{
1354    match timeout {
1355        Timeout::Infinite => FactoryRun::Produced(factory(ctx).await),
1356        Timeout::After(duration) => {
1357            let mut handle = tokio::spawn(factory(ctx));
1358            tokio::select! {
1359                joined = &mut handle => match joined {
1360                    Ok(result) => FactoryRun::Produced(result),
1361                    Err(_) => FactoryRun::Produced(Err(FactoryError::new("factory task panicked"))),
1362                },
1363                () = tokio::time::sleep(duration) => {
1364                    if allow_background {
1365                        FactoryRun::TimedOut(Some(handle))
1366                    } else {
1367                        handle.abort();
1368                        FactoryRun::TimedOut(None)
1369                    }
1370                }
1371            }
1372        }
1373    }
1374}
1375
1376/// Snapshots a stale entry into the conditional-refresh information handed to a
1377/// factory.
1378fn stale_info_of<V: Clone>(entry: &Entry<V>) -> StaleInfo<V> {
1379    StaleInfo {
1380        value: entry.value_cloned(),
1381        etag: entry.meta().etag().map(str::to_owned),
1382        last_modified: entry.meta().last_modified(),
1383        tags: entry.meta().tags().to_vec().into_boxed_slice(),
1384    }
1385}
1386
1387/// Builder for a [`Cache`].
1388///
1389/// ```
1390/// use amalgam::{Cache, EntryOptions};
1391/// use std::time::Duration;
1392///
1393/// let cache: Cache<String> = Cache::builder()
1394///     .name("users")
1395///     .key_prefix("u:")
1396///     .default_options(EntryOptions::new(Duration::from_secs(60)))
1397///     .build();
1398/// ```
1399#[must_use = "a builder does nothing until `.build()` is called"]
1400pub struct CacheBuilder<V> {
1401    name: Option<Arc<str>>,
1402    instance_id: Option<Arc<str>>,
1403    key_prefix: Option<Arc<str>>,
1404    default_options: EntryOptions,
1405    clock: Option<Arc<dyn Clock>>,
1406    max_capacity: Option<u64>,
1407    lock_shards: usize,
1408    remove_by_tag_behavior: RemoveByTagBehavior,
1409    events_capacity: usize,
1410    distributed: Option<Arc<dyn DistributedCache>>,
1411    serializer: Option<Arc<dyn DistributedSerializer<V>>>,
1412    backplane: Option<Arc<dyn Backplane>>,
1413    distributed_locker: Option<Arc<dyn DistributedLocker>>,
1414    plugins: Vec<Arc<dyn Plugin>>,
1415    distributed_circuit_breaker: Duration,
1416    backplane_circuit_breaker: Duration,
1417    recovery_config: RecoveryConfig,
1418    default_options_provider: Option<Arc<dyn DefaultEntryOptionsProvider>>,
1419    ignore_incoming_backplane: bool,
1420    distributed_wire_version: Arc<str>,
1421    distributed_key_modifier_mode: KeyModifierMode,
1422    disable_tagging: bool,
1423    wait_for_initial_backplane_subscribe: bool,
1424    _marker: std::marker::PhantomData<fn() -> V>,
1425}
1426
1427impl<V> CacheBuilder<V> {
1428    /// Creates a builder with default settings.
1429    pub fn new() -> Self {
1430        Self {
1431            name: None,
1432            instance_id: None,
1433            key_prefix: None,
1434            default_options: EntryOptions::default(),
1435            clock: None,
1436            max_capacity: None,
1437            lock_shards: 1024,
1438            remove_by_tag_behavior: RemoveByTagBehavior::default(),
1439            events_capacity: 256,
1440            distributed: None,
1441            serializer: None,
1442            backplane: None,
1443            distributed_locker: None,
1444            plugins: Vec::new(),
1445            distributed_circuit_breaker: Duration::ZERO,
1446            backplane_circuit_breaker: Duration::ZERO,
1447            recovery_config: RecoveryConfig::default(),
1448            default_options_provider: None,
1449            ignore_incoming_backplane: false,
1450            distributed_wire_version: Arc::from("v1"),
1451            distributed_key_modifier_mode: KeyModifierMode::default(),
1452            disable_tagging: false,
1453            wait_for_initial_backplane_subscribe: true,
1454            _marker: std::marker::PhantomData,
1455        }
1456    }
1457
1458    /// Sets this instance's id (used to ignore its own backplane messages). A
1459    /// random id is generated if not set.
1460    pub fn instance_id(mut self, id: impl AsRef<str>) -> Self {
1461        self.instance_id = Some(Arc::from(id.as_ref()));
1462        self
1463    }
1464
1465    /// Attaches an L2 distributed cache backend. Pair with
1466    /// [`serializer`](Self::serializer).
1467    pub fn distributed(mut self, distributed: Arc<dyn DistributedCache>) -> Self {
1468        self.distributed = Some(distributed);
1469        self
1470    }
1471
1472    /// Sets the serializer used for the L2 wire format (e.g.
1473    /// [`JsonSerializer`](crate::JsonSerializer)).
1474    pub fn serializer(mut self, serializer: Arc<dyn DistributedSerializer<V>>) -> Self {
1475        self.serializer = Some(serializer);
1476        self
1477    }
1478
1479    /// Attaches a backplane for multi-node L1 invalidation.
1480    ///
1481    /// Note: when a backplane is configured, [`build`](Self::build) spawns a
1482    /// listener task and so must be called from within a tokio runtime.
1483    pub fn backplane(mut self, backplane: Arc<dyn Backplane>) -> Self {
1484        self.backplane = Some(backplane);
1485        self
1486    }
1487
1488    /// Sets the cache's name (used in events/diagnostics).
1489    pub fn name(mut self, name: impl AsRef<str>) -> Self {
1490        self.name = Some(Arc::from(name.as_ref()));
1491        self
1492    }
1493
1494    /// Sets a prefix prepended to every key.
1495    pub fn key_prefix(mut self, prefix: impl AsRef<str>) -> Self {
1496        self.key_prefix = Some(Arc::from(prefix.as_ref()));
1497        self
1498    }
1499
1500    /// Sets the default entry options merged into every operation that does not
1501    /// supply its own.
1502    pub fn default_options(mut self, options: EntryOptions) -> Self {
1503        self.default_options = options;
1504        self
1505    }
1506
1507    /// Injects a custom [`Clock`] (e.g. [`ManualClock`](crate::ManualClock) in
1508    /// tests).
1509    pub fn clock(mut self, clock: Arc<dyn Clock>) -> Self {
1510        self.clock = Some(clock);
1511        self
1512    }
1513
1514    /// Caps the number of entries held in L1 (`None` = unbounded).
1515    pub fn max_capacity(mut self, capacity: u64) -> Self {
1516        self.max_capacity = Some(capacity);
1517        self
1518    }
1519
1520    /// Sets the number of single-flight lock shards (rounded up to a power of two).
1521    pub fn lock_shards(mut self, shards: usize) -> Self {
1522        self.lock_shards = shards;
1523        self
1524    }
1525
1526    /// Sets what `remove_by_tag` does to matched entries.
1527    pub fn remove_by_tag_behavior(mut self, behavior: RemoveByTagBehavior) -> Self {
1528        self.remove_by_tag_behavior = behavior;
1529        self
1530    }
1531
1532    /// Sets the event channel's buffer capacity.
1533    pub fn events_capacity(mut self, capacity: usize) -> Self {
1534        self.events_capacity = capacity;
1535        self
1536    }
1537
1538    /// Attaches a cross-node distributed locker (e.g. Redis-backed) for
1539    /// cluster-wide single-flight.
1540    pub fn distributed_locker(mut self, locker: Arc<dyn DistributedLocker>) -> Self {
1541        self.distributed_locker = Some(locker);
1542        self
1543    }
1544
1545    /// Registers a [`Plugin`] to observe events and lifecycle.
1546    pub fn plugin(mut self, plugin: Arc<dyn Plugin>) -> Self {
1547        self.plugins.push(plugin);
1548        self
1549    }
1550
1551    /// Opens the L2 circuit breaker for `duration` after an L2 failure
1552    /// (`Duration::ZERO` disables it — the default).
1553    pub fn distributed_circuit_breaker(mut self, duration: Duration) -> Self {
1554        self.distributed_circuit_breaker = duration;
1555        self
1556    }
1557
1558    /// Opens the backplane circuit breaker for `duration` after a backplane
1559    /// failure (`Duration::ZERO` disables it — the default).
1560    pub fn backplane_circuit_breaker(mut self, duration: Duration) -> Self {
1561        self.backplane_circuit_breaker = duration;
1562        self
1563    }
1564
1565    /// Configures (or disables) auto-recovery of failed L2 / backplane operations.
1566    pub fn auto_recovery(mut self, config: RecoveryConfig) -> Self {
1567        self.recovery_config = config;
1568        self
1569    }
1570
1571    /// Sets a provider of dynamic, per-key default options (consulted when a call
1572    /// supplies no explicit options).
1573    pub fn default_options_provider(
1574        mut self,
1575        provider: Arc<dyn DefaultEntryOptionsProvider>,
1576    ) -> Self {
1577        self.default_options_provider = Some(provider);
1578        self
1579    }
1580
1581    /// Drops all incoming backplane notifications (dangerous; for testing).
1582    pub fn ignore_incoming_backplane(mut self, ignore: bool) -> Self {
1583        self.ignore_incoming_backplane = ignore;
1584        self
1585    }
1586
1587    /// Sets the L2 wire-format version combined with distributed keys.
1588    pub fn distributed_wire_version(mut self, version: impl AsRef<str>) -> Self {
1589        self.distributed_wire_version = Arc::from(version.as_ref());
1590        self
1591    }
1592
1593    /// Sets how the wire-format version is combined with the L2 key (prefix,
1594    /// suffix, or none).
1595    pub fn distributed_key_modifier_mode(mut self, mode: KeyModifierMode) -> Self {
1596        self.distributed_key_modifier_mode = mode;
1597        self
1598    }
1599
1600    /// Disables tagging cache-wide (FusionCache `DisableTagging`).
1601    ///
1602    /// When enabled, the per-read tag/clear-marker check is skipped entirely (a
1603    /// small read-path saving for caches that never tag), and
1604    /// [`remove_by_tag`](Cache::remove_by_tag) /
1605    /// [`remove_by_tags`](Cache::remove_by_tags) / [`clear`](Cache::clear) are
1606    /// ignored and logged at `warn` (never silently) — the Rust-idiomatic
1607    /// counterpart of FusionCache throwing on tag use when tagging is disabled.
1608    /// Off by default.
1609    pub fn disable_tagging(mut self, disable: bool) -> Self {
1610        self.disable_tagging = disable;
1611        self
1612    }
1613
1614    /// Whether to establish the backplane subscription before [`build`](Self::build)
1615    /// returns (FusionCache `WaitForInitialBackplaneSubscribe`, default `true`).
1616    ///
1617    /// In `amalgam` this is inherently satisfied: the listener calls
1618    /// `backplane.subscribe()` synchronously inside `build`, and the Redis backplane
1619    /// completes its `SUBSCRIBE` during `connect()` before being handed to the
1620    /// builder — so there is no startup window where the node serves while
1621    /// unsubscribed. The flag is retained for configuration parity and is queryable
1622    /// via [`Cache::wait_for_initial_backplane_subscribe`].
1623    pub fn wait_for_initial_backplane_subscribe(mut self, wait: bool) -> Self {
1624        self.wait_for_initial_backplane_subscribe = wait;
1625        self
1626    }
1627}
1628
1629impl<V> Default for CacheBuilder<V> {
1630    fn default() -> Self {
1631        Self::new()
1632    }
1633}
1634
1635impl<V: Clone + Send + Sync + 'static> CacheBuilder<V> {
1636    /// Builds the cache.
1637    #[must_use]
1638    pub fn build(self) -> Cache<V> {
1639        let clock: Arc<dyn Clock> = self.clock.unwrap_or_else(|| Arc::new(SystemClock));
1640        // Auto-recovery only matters when there is an L2 or backplane to recover.
1641        let recovery = if self.recovery_config.enabled
1642            && (self.distributed.is_some() || self.backplane.is_some())
1643        {
1644            Some(AutoRecoveryService::new(
1645                self.recovery_config,
1646                Arc::clone(&clock),
1647            ))
1648        } else {
1649            None
1650        };
1651        let events = Events::with_capacity(self.events_capacity);
1652        let inner = Arc::new(CacheInner {
1653            name: self.name.unwrap_or_else(|| Arc::from("amalgam")),
1654            instance_id: self.instance_id.unwrap_or_else(generate_instance_id),
1655            memory: MemoryStore::new(self.max_capacity, events.clone()),
1656            locks: Arc::new(KeyedLock::new(self.lock_shards)),
1657            tags: Arc::new(TagRegistry::new()),
1658            events,
1659            clock,
1660            default_options: self.default_options,
1661            key_prefix: self.key_prefix,
1662            remove_by_tag_behavior: self.remove_by_tag_behavior,
1663            distributed: self.distributed,
1664            serializer: self.serializer,
1665            backplane: self.backplane,
1666            distributed_locker: self.distributed_locker,
1667            circuit_l2: CircuitBreaker::new(self.distributed_circuit_breaker),
1668            circuit_backplane: CircuitBreaker::new(self.backplane_circuit_breaker),
1669            plugins: PluginHost::new(self.plugins),
1670            recovery: recovery.clone(),
1671            default_options_provider: self.default_options_provider,
1672            ignore_incoming_backplane: self.ignore_incoming_backplane,
1673            distributed_wire_version: self.distributed_wire_version,
1674            distributed_key_modifier_mode: self.distributed_key_modifier_mode,
1675            disable_tagging: self.disable_tagging,
1676            wait_for_initial_backplane_subscribe: self.wait_for_initial_backplane_subscribe,
1677        });
1678        // Wire the recovery executor (the cache) as a Weak so it never keeps the
1679        // cache alive, then start the background drain loop.
1680        if let Some(recovery) = &recovery {
1681            let executor: Arc<dyn RecoveryExecutor> =
1682                Arc::clone(&inner) as Arc<dyn RecoveryExecutor>;
1683            recovery.set_executor(Arc::downgrade(&executor));
1684            recovery.spawn();
1685        }
1686        let cache = Cache { inner };
1687        cache.spawn_backplane_listener();
1688        cache
1689    }
1690}
1691
1692/// Generates a random per-instance id for backplane self-filtering.
1693fn generate_instance_id() -> Arc<str> {
1694    Arc::from(format!("amalgam-{:016x}", fastrand::u64(..)))
1695}
1696
1697/// Picks the entry with the later creation timestamp (the "newer" fallback).
1698fn newer_of<V: Clone>(existing: Option<Entry<V>>, candidate: Entry<V>) -> Entry<V> {
1699    match existing {
1700        Some(existing) if existing.meta().created() >= candidate.meta().created() => existing,
1701        _ => candidate,
1702    }
1703}