amalgam/backplane.rs
1//! The backplane: multi-node invalidation messaging.
2//!
3//! In a multi-node deployment each node has its own L1. When one node changes a
4//! key, it publishes a small [`BackplaneMessage`] (never the value itself) so the
5//! other nodes can invalidate or expire their local L1 copy and re-pull the
6//! authoritative value from L2 on the next read.
7//!
8//! The [`InProcessBackplane`] here is a faithful reference backed by a broadcast
9//! channel — perfect for tests and single-process multi-instance setups. A Redis
10//! pub/sub adapter is a feature-gated follow-up.
11
12use std::sync::Arc;
13
14use async_trait::async_trait;
15use tokio::sync::broadcast;
16
17use crate::error::Result;
18use crate::time::Timestamp;
19
20/// The kind of change a backplane message announces.
21#[derive(Debug, Clone, Copy, PartialEq, Eq)]
22pub enum BackplaneAction {
23 /// A key was written; receivers should drop their local copy and re-pull L2.
24 Set,
25 /// A key was removed; receivers should hard-evict it.
26 Remove,
27 /// A key was logically expired; receivers should expire (keep for fail-safe).
28 Expire,
29}
30
31/// A backplane notification. Carries no value — only what changed and when.
32#[derive(Debug, Clone)]
33pub struct BackplaneMessage {
34 /// The id of the publishing cache instance (so receivers ignore their own).
35 pub source_id: Arc<str>,
36 /// When the change happened (for "newer wins" expiration).
37 pub timestamp: Timestamp,
38 /// What kind of change occurred.
39 pub action: BackplaneAction,
40 /// The (prefixed) key affected.
41 pub key: Arc<str>,
42}
43
44/// A multi-node notification channel.
45///
46/// Implement over Redis pub/sub, NATS, etc. The cache publishes on local changes
47/// and subscribes to apply remote ones.
48#[async_trait]
49pub trait Backplane: Send + Sync {
50 /// Publishes a message to all other subscribers.
51 ///
52 /// # Errors
53 /// Returns [`Error::Backplane`](crate::Error::Backplane) on transport failure.
54 async fn publish(&self, message: BackplaneMessage) -> Result<()>;
55
56 /// Subscribes to the message stream.
57 fn subscribe(&self) -> broadcast::Receiver<BackplaneMessage>;
58}
59
60/// An in-process reference [`Backplane`] backed by a broadcast channel.
61///
62/// Share one instance (via `Arc`) between several [`Cache`](crate::Cache)
63/// instances to simulate a multi-node cluster within one process.
64#[derive(Clone)]
65pub struct InProcessBackplane {
66 sender: broadcast::Sender<BackplaneMessage>,
67}
68
69impl InProcessBackplane {
70 /// Creates a backplane with the given message buffer capacity.
71 #[must_use]
72 pub fn with_capacity(capacity: usize) -> Self {
73 let (sender, _) = broadcast::channel(capacity.max(1));
74 Self { sender }
75 }
76}
77
78impl Default for InProcessBackplane {
79 fn default() -> Self {
80 Self::with_capacity(256)
81 }
82}
83
84#[async_trait]
85impl Backplane for InProcessBackplane {
86 async fn publish(&self, message: BackplaneMessage) -> Result<()> {
87 // A send error only means "no live subscribers"; that is fine.
88 let _ = self.sender.send(message);
89 Ok(())
90 }
91
92 fn subscribe(&self) -> broadcast::Receiver<BackplaneMessage> {
93 self.sender.subscribe()
94 }
95}