1use 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
38pub 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
88enum L1Read<V> {
90 Fresh(Entry<V>),
91 Stale(Entry<V>),
92 Miss,
93}
94
95impl<V: Clone + Send + Sync + 'static> Cache<V> {
96 pub fn builder() -> CacheBuilder<V> {
98 CacheBuilder::new()
99 }
100
101 #[must_use]
103 pub fn new() -> Self {
104 CacheBuilder::new().build()
105 }
106
107 #[must_use]
109 pub fn name(&self) -> &str {
110 &self.inner.name
111 }
112
113 #[must_use]
115 pub fn events(&self) -> &Events {
116 &self.inner.events
117 }
118
119 #[must_use]
121 pub fn entry_options(&self) -> EntryOptions {
122 self.inner.default_options.clone()
123 }
124
125 #[must_use]
130 pub fn wait_for_initial_backplane_subscribe(&self) -> bool {
131 self.inner.wait_for_initial_backplane_subscribe
132 }
133
134 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 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 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 #[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 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 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 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 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 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 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 pub async fn set(&self, key: impl AsRef<str>, value: V) {
377 self.set_full(key, value, None, Box::from([])).await;
378 }
379
380 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 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 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 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 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 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 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 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 pub async fn run_pending_tasks(&self) {
534 self.inner.memory.run_pending_tasks().await;
535 }
536
537 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 self.close_circuit_backplane();
943
944 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 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 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 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; }
1037 match weak.upgrade() {
1038 Some(inner) => Cache { inner }.apply_backplane(message).await,
1039 None => break, }
1041 }
1042 Err(broadcast::error::RecvError::Closed) => break,
1043 Err(broadcast::error::RecvError::Lagged(_)) => continue,
1044 }
1045 }
1046 });
1047 }
1048
1049 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 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 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 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(¤t)),
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 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 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 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 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 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
1273const TAG_MARKER_PREFIX: &str = "__amalgam:t:";
1275const CLEAR_EXPIRE_KEY: &str = "__amalgam:clear:expire";
1277const CLEAR_REMOVE_KEY: &str = "__amalgam:clear:remove";
1279const RECOVERY_ITEM_TTL: Duration = Duration::from_secs(600);
1281
1282fn 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
1291struct LockGuard {
1294 _local: KeyGuard,
1295 _distributed: Option<DistributedReleaseGuard>,
1296}
1297
1298struct 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
1316enum LockOutcome<V> {
1318 Acquired(LockGuard),
1320 ServedStale(V),
1322}
1323
1324#[allow(clippy::large_enum_variant)]
1330enum FactoryRun<V> {
1331 Produced(std::result::Result<FactoryProduct<V>, FactoryError>),
1333 TimedOut(Option<JoinHandle<std::result::Result<FactoryProduct<V>, FactoryError>>>),
1336}
1337
1338async 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
1376fn 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#[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 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 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 pub fn distributed(mut self, distributed: Arc<dyn DistributedCache>) -> Self {
1468 self.distributed = Some(distributed);
1469 self
1470 }
1471
1472 pub fn serializer(mut self, serializer: Arc<dyn DistributedSerializer<V>>) -> Self {
1475 self.serializer = Some(serializer);
1476 self
1477 }
1478
1479 pub fn backplane(mut self, backplane: Arc<dyn Backplane>) -> Self {
1484 self.backplane = Some(backplane);
1485 self
1486 }
1487
1488 pub fn name(mut self, name: impl AsRef<str>) -> Self {
1490 self.name = Some(Arc::from(name.as_ref()));
1491 self
1492 }
1493
1494 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 pub fn default_options(mut self, options: EntryOptions) -> Self {
1503 self.default_options = options;
1504 self
1505 }
1506
1507 pub fn clock(mut self, clock: Arc<dyn Clock>) -> Self {
1510 self.clock = Some(clock);
1511 self
1512 }
1513
1514 pub fn max_capacity(mut self, capacity: u64) -> Self {
1516 self.max_capacity = Some(capacity);
1517 self
1518 }
1519
1520 pub fn lock_shards(mut self, shards: usize) -> Self {
1522 self.lock_shards = shards;
1523 self
1524 }
1525
1526 pub fn remove_by_tag_behavior(mut self, behavior: RemoveByTagBehavior) -> Self {
1528 self.remove_by_tag_behavior = behavior;
1529 self
1530 }
1531
1532 pub fn events_capacity(mut self, capacity: usize) -> Self {
1534 self.events_capacity = capacity;
1535 self
1536 }
1537
1538 pub fn distributed_locker(mut self, locker: Arc<dyn DistributedLocker>) -> Self {
1541 self.distributed_locker = Some(locker);
1542 self
1543 }
1544
1545 pub fn plugin(mut self, plugin: Arc<dyn Plugin>) -> Self {
1547 self.plugins.push(plugin);
1548 self
1549 }
1550
1551 pub fn distributed_circuit_breaker(mut self, duration: Duration) -> Self {
1554 self.distributed_circuit_breaker = duration;
1555 self
1556 }
1557
1558 pub fn backplane_circuit_breaker(mut self, duration: Duration) -> Self {
1561 self.backplane_circuit_breaker = duration;
1562 self
1563 }
1564
1565 pub fn auto_recovery(mut self, config: RecoveryConfig) -> Self {
1567 self.recovery_config = config;
1568 self
1569 }
1570
1571 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 pub fn ignore_incoming_backplane(mut self, ignore: bool) -> Self {
1583 self.ignore_incoming_backplane = ignore;
1584 self
1585 }
1586
1587 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 pub fn distributed_key_modifier_mode(mut self, mode: KeyModifierMode) -> Self {
1596 self.distributed_key_modifier_mode = mode;
1597 self
1598 }
1599
1600 pub fn disable_tagging(mut self, disable: bool) -> Self {
1610 self.disable_tagging = disable;
1611 self
1612 }
1613
1614 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 #[must_use]
1638 pub fn build(self) -> Cache<V> {
1639 let clock: Arc<dyn Clock> = self.clock.unwrap_or_else(|| Arc::new(SystemClock));
1640 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 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
1692fn generate_instance_id() -> Arc<str> {
1694 Arc::from(format!("amalgam-{:016x}", fastrand::u64(..)))
1695}
1696
1697fn 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}