use choreo_proto::DaemonMessage;
use std::sync::mpsc;
use tracing::warn;
pub(crate) fn try_send_keep_on_full(
tx: &mpsc::SyncSender<DaemonMessage>,
client_id: u64,
path: &str,
message: &DaemonMessage,
) -> bool {
match tx.try_send(message.clone()) {
Ok(()) => true,
Err(mpsc::TrySendError::Full(_)) => {
crate::metrics::record_broadcast_dropped(path);
true
}
Err(mpsc::TrySendError::Disconnected(_)) => {
warn!("removing disconnected subscriber {client_id}");
false
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use choreo_proto::SessionStatus;
use std::sync::mpsc;
fn status_msg(session_id: u64) -> DaemonMessage {
DaemonMessage::SessionStatusChanged {
session_id,
status: SessionStatus::Inactive,
last_modified: 0,
}
}
#[test]
fn try_send_keep_on_full_delivers_and_keeps_subscriber() {
let (tx, rx) = mpsc::sync_channel::<DaemonMessage>(1);
let msg = status_msg(1);
assert!(try_send_keep_on_full(&tx, 7, "summary", &msg));
assert_eq!(rx.recv().unwrap(), msg);
}
#[test]
fn try_send_keep_on_full_drops_on_full_buffer_and_keeps_subscriber() {
let (tx, rx) = mpsc::sync_channel::<DaemonMessage>(1);
let filler = status_msg(1);
tx.send(filler.clone()).unwrap();
let broadcast = status_msg(2);
assert!(try_send_keep_on_full(&tx, 7, "session", &broadcast));
assert_eq!(rx.recv().unwrap(), filler);
assert!(rx.try_recv().is_err(), "full-buffer message was dropped");
assert!(try_send_keep_on_full(&tx, 7, "session", &broadcast));
assert_eq!(rx.recv().unwrap(), broadcast);
}
#[test]
fn try_send_keep_on_full_evicts_on_disconnected() {
let (tx, rx) = mpsc::sync_channel::<DaemonMessage>(1);
drop(rx); assert!(!try_send_keep_on_full(&tx, 7, "summary", &status_msg(1)));
}
}