use std::sync::{Arc, Mutex};
use tokio::sync::broadcast;
const CHANNEL_CAPACITY: usize = 100;
pub(crate) struct CachingSender<T> {
tx: broadcast::Sender<T>,
last: Arc<Mutex<Option<T>>>,
}
impl<T> Clone for CachingSender<T> {
fn clone(&self) -> Self {
Self {
tx: self.tx.clone(),
last: Arc::clone(&self.last),
}
}
}
impl<T: Clone> CachingSender<T> {
pub(crate) fn new() -> Self {
let (tx, _rx) = broadcast::channel(CHANNEL_CAPACITY);
Self {
tx,
last: Arc::new(Mutex::new(None)),
}
}
pub(crate) fn publish(&self, value: T) -> bool {
let mut last = self.last.lock().unwrap_or_else(|e| e.into_inner());
*last = Some(value.clone());
self.tx.send(value).is_ok()
}
pub(crate) fn publish_if_changed(&self, value: &T) -> bool
where
T: PartialEq,
{
let mut last = self.last.lock().unwrap_or_else(|e| e.into_inner());
if last.as_ref() == Some(value) {
return self.tx.receiver_count() > 0;
}
*last = Some(value.clone());
self.tx.send(value.clone()).is_ok()
}
pub(crate) fn receiver_count(&self) -> usize {
self.tx.receiver_count()
}
pub(crate) fn subscribe_with_current(&self) -> broadcast::Receiver<T> {
let last = self.last.lock().unwrap_or_else(|e| e.into_inner());
let rx = self.tx.subscribe();
if let Some(value) = last.clone() {
let _ = self.tx.send(value);
}
rx
}
pub(crate) fn downgrade(&self) -> WeakCachingSender<T> {
WeakCachingSender {
tx: self.tx.downgrade(),
last: Arc::clone(&self.last),
}
}
}
pub(crate) struct WeakCachingSender<T> {
tx: broadcast::WeakSender<T>,
last: Arc<Mutex<Option<T>>>,
}
impl<T> WeakCachingSender<T> {
pub(crate) fn upgrade(&self) -> Option<CachingSender<T>> {
self.tx.upgrade().map(|tx| CachingSender {
tx,
last: Arc::clone(&self.last),
})
}
}
#[cfg(test)]
mod tests {
use super::CachingSender;
#[test]
fn late_subscriber_is_re_sent_the_cached_value() {
let tx = CachingSender::new();
let _first = tx.subscribe_with_current();
assert!(tx.publish(7));
let mut late = tx.subscribe_with_current();
assert_eq!(
late.try_recv().expect("cached value re-sent on subscribe"),
7
);
}
#[test]
fn publish_reports_whether_anyone_is_listening() {
let tx = CachingSender::new();
assert!(!tx.publish(1), "no receivers yet");
let _rx = tx.subscribe_with_current();
assert!(tx.publish(2), "a receiver is listening");
}
#[test]
fn publish_if_changed_does_not_resend_an_unchanged_value() {
let tx = CachingSender::new();
let mut rx = tx.subscribe_with_current();
assert!(tx.publish_if_changed(&5));
assert!(tx.publish_if_changed(&5), "still has a receiver");
assert!(tx.publish_if_changed(&6));
assert_eq!(rx.try_recv().unwrap(), 5);
assert_eq!(rx.try_recv().unwrap(), 6);
assert!(rx.try_recv().is_err(), "the unchanged 5 was not re-sent");
}
#[test]
fn weak_handle_upgrades_and_shares_the_cache_while_a_producer_lives() {
let tx = CachingSender::new();
let _rx = tx.subscribe_with_current();
assert!(tx.publish(11));
let weak = tx.downgrade();
let _producer = tx.clone();
drop(tx);
let upgraded = weak
.upgrade()
.expect("upgradable while a strong sender lives");
let mut late = upgraded.subscribe_with_current();
assert_eq!(late.try_recv().expect("cached value re-sent"), 11);
}
#[test]
fn weak_handle_does_not_upgrade_after_the_producer_exits() {
let tx = CachingSender::<i32>::new();
let weak = tx.downgrade();
let _rx = tx.subscribe_with_current();
drop(tx);
assert!(weak.upgrade().is_none());
}
}