use std::sync::{
Arc,
atomic::{AtomicU64, Ordering},
};
use std::time::{Duration, Instant};
use tokio::sync::mpsc;
const DEPTH: usize = 8;
const REPORT_INTERVAL: Duration = Duration::from_secs(1);
pub(super) fn bounded<T>() -> (Sender<T>, Receiver<T>) {
let (tx, rx) = mpsc::channel(DEPTH);
let dropped = Arc::new(AtomicU64::new(0));
let sender = Sender {
tx,
dropped: dropped.clone(),
};
let receiver = Receiver {
rx,
dropped,
gap: false,
unreported: 0,
last_report: None,
};
(sender, receiver)
}
pub(super) struct Sender<T> {
tx: mpsc::Sender<T>,
dropped: Arc<AtomicU64>,
}
impl<T> Sender<T> {
pub(super) fn push(&self, item: T) {
match self.tx.try_send(item) {
Ok(()) | Err(mpsc::error::TrySendError::Closed(_)) => {}
Err(mpsc::error::TrySendError::Full(_)) => {
self.dropped.fetch_add(1, Ordering::Relaxed);
}
}
}
}
pub(super) struct Receiver<T> {
rx: mpsc::Receiver<T>,
dropped: Arc<AtomicU64>,
gap: bool,
unreported: u64,
last_report: Option<Instant>,
}
impl<T> Receiver<T> {
pub(super) async fn recv(&mut self) -> Option<T> {
let item = self.rx.recv().await;
self.observe();
item
}
pub(super) fn gap(&mut self) -> bool {
std::mem::take(&mut self.gap)
}
fn observe(&mut self) {
let dropped = self.dropped.swap(0, Ordering::Relaxed);
if dropped > 0 {
self.gap = true;
self.unreported += dropped;
}
if self.unreported == 0 {
return;
}
let now = Instant::now();
if self
.last_report
.is_some_and(|last| now.duration_since(last) < REPORT_INTERVAL)
{
return;
}
tracing::warn!(
dropped = self.unreported,
capacity = DEPTH,
"dropped audio capture buffers"
);
self.last_report = Some(now);
self.unreported = 0;
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn drops_newest_when_full() {
let (tx, mut rx) = bounded();
for id in 0..DEPTH + 2 {
tx.push(id);
}
drop(tx);
let mut received = Vec::new();
while let Some(id) = rx.recv().await {
received.push(id);
}
assert_eq!(received, (0..DEPTH).collect::<Vec<_>>());
}
#[tokio::test]
async fn drop_marks_a_gap() {
let (tx, mut rx) = bounded();
for id in 0..DEPTH + 1 {
tx.push(id);
}
assert_eq!(rx.recv().await, Some(0));
assert!(rx.gap());
assert_eq!(rx.recv().await, Some(1));
assert!(!rx.gap());
}
#[tokio::test]
async fn no_gap_without_drops() {
let (tx, mut rx) = bounded();
tx.push(1);
assert_eq!(rx.recv().await, Some(1));
assert!(!rx.gap());
}
#[tokio::test]
async fn close_returns_none_after_draining() {
let (tx, mut rx) = bounded();
tx.push(1);
drop(tx);
assert_eq!(rx.recv().await, Some(1));
assert_eq!(rx.recv().await, None);
}
}