use std::{
collections::VecDeque,
sync::atomic::{AtomicU64, Ordering::Relaxed},
};
use parking_lot::Mutex;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PubSubMessageKind {
Channel,
Pattern,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PubSubMessage {
pub kind: PubSubMessageKind,
pub pattern: Option<Box<[u8]>>,
pub channel: Box<[u8]>,
pub value: Box<[u8]>,
}
pub trait PubSubSink: Send + Sync {
fn publish(&self, channel: &[u8], value: &[u8]);
fn pattern_publish(&self, pattern: &[u8], channel: &[u8], value: &[u8]);
}
pub struct PubSubMailbox {
capacity: usize,
queue: Mutex<VecDeque<PubSubMessage>>,
dropped: AtomicU64,
}
impl PubSubMailbox {
pub fn new(capacity: usize) -> Self {
Self {
capacity: capacity.max(1),
queue: Mutex::new(VecDeque::new()),
dropped: AtomicU64::new(0),
}
}
pub fn drain(&self) -> Vec<PubSubMessage> {
self.queue.lock().drain(..).collect()
}
pub fn len(&self) -> usize {
self.queue.lock().len()
}
pub fn is_empty(&self) -> bool {
self.queue.lock().is_empty()
}
pub fn dropped_count(&self) -> u64 {
self.dropped.load(Relaxed)
}
fn push(&self, message: PubSubMessage) {
let mut queue = self.queue.lock();
if queue.len() >= self.capacity {
queue.pop_front();
self.dropped.fetch_add(1, Relaxed);
}
queue.push_back(message);
}
}
impl PubSubSink for PubSubMailbox {
fn publish(&self, channel: &[u8], value: &[u8]) {
self.push(PubSubMessage {
kind: PubSubMessageKind::Channel,
pattern: None,
channel: channel.into(),
value: value.into(),
});
}
fn pattern_publish(&self, pattern: &[u8], channel: &[u8], value: &[u8]) {
self.push(PubSubMessage {
kind: PubSubMessageKind::Pattern,
pattern: Some(pattern.into()),
channel: channel.into(),
value: value.into(),
});
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn mailbox_bounded_and_drops_tail() {
let mailbox = PubSubMailbox::new(2);
mailbox.publish(b"a", b"1");
mailbox.publish(b"b", b"2");
mailbox.publish(b"c", b"3");
assert_eq!(mailbox.len(), 2);
assert_eq!(mailbox.dropped_count(), 1);
let messages = mailbox.drain();
assert_eq!(messages.len(), 2);
assert_eq!(messages[0].channel.as_ref(), b"b");
assert!(mailbox.is_empty());
}
#[test]
fn mailbox_pattern_message_keeps_pattern() {
let mailbox = PubSubMailbox::new(4);
mailbox.pattern_publish(b"a*", b"ab", b"v");
let messages = mailbox.drain();
assert_eq!(messages[0].kind, PubSubMessageKind::Pattern);
assert_eq!(messages[0].pattern.as_deref(), Some(b"a*".as_slice()));
}
}