Skip to main content

amalgam/
events.rs

1//! The event hub.
2//!
3//! FusionCache exposes a rich set of events (hits, misses, fail-safe
4//! activations, factory errors, …) fired on background threads. The idiomatic
5//! Rust equivalent is a broadcast channel: subscribers receive a stream of
6//! [`CacheEvent`]s without blocking the cache's hot path, and handler execution
7//! is naturally decoupled from the operation that produced the event.
8
9use std::sync::Arc;
10
11use tokio::sync::broadcast;
12
13/// An observable cache event.
14#[derive(Debug, Clone, PartialEq, Eq)]
15#[non_exhaustive]
16pub enum CacheEvent {
17    /// A value was served. `stale` is `true` when it came from a fail-safe /
18    /// stale fallback rather than a fresh entry.
19    Hit {
20        /// The cache key.
21        key: Arc<str>,
22        /// Whether the served value was stale.
23        stale: bool,
24    },
25    /// Nothing servable was found for the key.
26    Miss {
27        /// The cache key.
28        key: Arc<str>,
29    },
30    /// A value was written to the cache.
31    Set {
32        /// The cache key.
33        key: Arc<str>,
34    },
35    /// A key was removed.
36    Remove {
37        /// The cache key.
38        key: Arc<str>,
39    },
40    /// A key was logically expired.
41    Expire {
42        /// The cache key.
43        key: Arc<str>,
44    },
45    /// The factory completed successfully on the foreground path.
46    FactorySuccess {
47        /// The cache key.
48        key: Arc<str>,
49    },
50    /// The factory returned an error.
51    FactoryError {
52        /// The cache key.
53        key: Arc<str>,
54        /// The failure message.
55        message: String,
56    },
57    /// The factory exceeded a (soft or hard) timeout.
58    FactorySyntheticTimeout {
59        /// The cache key.
60        key: Arc<str>,
61    },
62    /// A stale value was served because the factory failed or timed out.
63    FailSafeActivate {
64        /// The cache key.
65        key: Arc<str>,
66    },
67    /// A proactive background refresh was started.
68    EagerRefresh {
69        /// The cache key.
70        key: Arc<str>,
71    },
72    /// A timed-out factory later completed successfully in the background.
73    BackgroundFactorySuccess {
74        /// The cache key.
75        key: Arc<str>,
76    },
77    /// A timed-out factory later failed in the background.
78    BackgroundFactoryError {
79        /// The cache key.
80        key: Arc<str>,
81        /// The failure message.
82        message: String,
83    },
84    /// All entries carrying a tag were invalidated.
85    RemoveByTag {
86        /// The tag.
87        tag: String,
88    },
89    /// The whole cache was cleared.
90    Clear,
91    /// An entry was evicted from L1 by the backend's size/expiry policy.
92    Eviction {
93        /// The cache key.
94        key: Arc<str>,
95    },
96    /// A distributed-cache or backplane circuit breaker opened or closed.
97    CircuitBreakerChange {
98        /// Which component the breaker guards.
99        component: CircuitComponent,
100        /// `true` if the breaker is now closed (healthy), `false` if open.
101        closed: bool,
102    },
103    /// A value failed to serialize for L2.
104    SerializationError {
105        /// The cache key.
106        key: Arc<str>,
107        /// The error message.
108        message: String,
109    },
110    /// A value failed to deserialize from L2.
111    DeserializationError {
112        /// The cache key.
113        key: Arc<str>,
114        /// The error message.
115        message: String,
116    },
117    /// A backplane notification was published to peers.
118    MessagePublished {
119        /// The cache key.
120        key: Arc<str>,
121    },
122    /// A backplane notification was received from a peer.
123    MessageReceived {
124        /// The cache key.
125        key: Arc<str>,
126    },
127}
128
129/// Which subsystem a [`CacheEvent::CircuitBreakerChange`] refers to.
130#[derive(Debug, Clone, Copy, PartialEq, Eq)]
131pub enum CircuitComponent {
132    /// The L2 distributed cache.
133    Distributed,
134    /// The multi-node backplane.
135    Backplane,
136}
137
138/// A broadcaster of [`CacheEvent`]s.
139///
140/// Cloning an `Events` shares the same underlying channel.
141#[derive(Debug, Clone)]
142pub struct Events {
143    sender: broadcast::Sender<CacheEvent>,
144}
145
146impl Events {
147    /// Creates a hub with the given subscriber buffer capacity.
148    #[must_use]
149    pub fn with_capacity(capacity: usize) -> Self {
150        let (sender, _) = broadcast::channel(capacity.max(1));
151        Self { sender }
152    }
153
154    /// Subscribes to the event stream.
155    ///
156    /// A subscriber that falls behind by more than the buffer capacity will
157    /// observe a `Lagged` error from the receiver — events are best-effort
158    /// observability, never a correctness mechanism.
159    #[must_use]
160    pub fn subscribe(&self) -> broadcast::Receiver<CacheEvent> {
161        self.sender.subscribe()
162    }
163
164    /// Emits an event. Does nothing if there are no subscribers.
165    pub fn emit(&self, event: CacheEvent) {
166        // A send error only means "no live subscribers"; that is expected.
167        let _ = self.sender.send(event);
168    }
169}
170
171impl Default for Events {
172    fn default() -> Self {
173        Self::with_capacity(256)
174    }
175}