Skip to main content

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}