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}