use std::sync::Arc;
use async_trait::async_trait;
use tokio::sync::broadcast;
use crate::error::Result;
use crate::time::Timestamp;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum BackplaneAction {
Set,
Remove,
Expire,
}
#[derive(Debug, Clone)]
pub struct BackplaneMessage {
pub source_id: Arc<str>,
pub timestamp: Timestamp,
pub action: BackplaneAction,
pub key: Arc<str>,
}
#[async_trait]
pub trait Backplane: Send + Sync {
async fn publish(&self, message: BackplaneMessage) -> Result<()>;
fn subscribe(&self) -> broadcast::Receiver<BackplaneMessage>;
}
#[derive(Clone)]
pub struct InProcessBackplane {
sender: broadcast::Sender<BackplaneMessage>,
}
impl InProcessBackplane {
#[must_use]
pub fn with_capacity(capacity: usize) -> Self {
let (sender, _) = broadcast::channel(capacity.max(1));
Self { sender }
}
}
impl Default for InProcessBackplane {
fn default() -> Self {
Self::with_capacity(256)
}
}
#[async_trait]
impl Backplane for InProcessBackplane {
async fn publish(&self, message: BackplaneMessage) -> Result<()> {
let _ = self.sender.send(message);
Ok(())
}
fn subscribe(&self) -> broadcast::Receiver<BackplaneMessage> {
self.sender.subscribe()
}
}