use std::collections::VecDeque;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
pub const DEFAULT_CAPACITY: usize = 256;
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct SinkStats {
pub accepted: u64,
pub dropped: u64,
}
impl SinkStats {
pub fn offered(&self) -> u64 {
self.accepted + self.dropped
}
pub fn is_complete(&self) -> bool {
self.dropped == 0
}
}
#[derive(Clone, Debug)]
pub struct SampleSink<T> {
inner: Arc<Inner<T>>,
}
#[derive(Debug)]
struct Inner<T> {
queue: Mutex<VecDeque<T>>,
capacity: usize,
accepted: AtomicU64,
dropped: AtomicU64,
}
impl<T> SampleSink<T> {
pub fn with_capacity(capacity: usize) -> Self {
Self {
inner: Arc::new(Inner {
queue: Mutex::new(VecDeque::with_capacity(capacity.max(1))),
capacity: capacity.max(1),
accepted: AtomicU64::new(0),
dropped: AtomicU64::new(0),
}),
}
}
pub fn offer(&self, sample: T) -> bool {
let mut queue = match self.inner.queue.lock() {
Ok(q) => q,
Err(poisoned) => poisoned.into_inner(),
};
if queue.len() >= self.inner.capacity {
self.inner.dropped.fetch_add(1, Ordering::Relaxed);
return false;
}
queue.push_back(sample);
self.inner.accepted.fetch_add(1, Ordering::Relaxed);
true
}
pub fn drain(&self) -> Vec<T> {
let mut queue = match self.inner.queue.lock() {
Ok(q) => q,
Err(poisoned) => poisoned.into_inner(),
};
queue.drain(..).collect()
}
pub fn len(&self) -> usize {
match self.inner.queue.lock() {
Ok(q) => q.len(),
Err(poisoned) => poisoned.into_inner().len(),
}
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn stats(&self) -> SinkStats {
SinkStats {
accepted: self.inner.accepted.load(Ordering::Relaxed),
dropped: self.inner.dropped.load(Ordering::Relaxed),
}
}
}
impl<T> Default for SampleSink<T> {
fn default() -> Self {
Self::with_capacity(DEFAULT_CAPACITY)
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::AtomicBool;
use std::time::{Duration, Instant};
#[test]
fn accepts_up_to_capacity_then_drops() {
let sink = SampleSink::with_capacity(4);
for i in 0..4 {
assert!(sink.offer(i), "sample {i} should fit");
}
assert!(!sink.offer(99), "the 5th sample must be dropped");
let stats = sink.stats();
assert_eq!(stats.accepted, 4);
assert_eq!(stats.dropped, 1);
assert_eq!(stats.offered(), 5);
assert!(!stats.is_complete());
}
#[test]
fn draining_makes_room_again() {
let sink = SampleSink::with_capacity(2);
sink.offer(1);
sink.offer(2);
assert!(!sink.offer(3));
assert_eq!(sink.drain(), vec![1, 2]);
assert!(
sink.offer(4),
"a drained sink must accept again — drops are transient, not terminal"
);
}
#[test]
fn drain_preserves_order() {
let sink = SampleSink::with_capacity(8);
for i in 0..5 {
sink.offer(i);
}
assert_eq!(sink.drain(), vec![0, 1, 2, 3, 4]);
}
#[test]
fn zero_capacity_is_treated_as_one() {
let sink = SampleSink::with_capacity(0);
assert!(
sink.offer(1),
"a sink that can never accept would report 100% drops and hide \
whether the consumer ever worked"
);
}
#[test]
fn producer_never_blocks_on_a_stalled_consumer() {
let sink: SampleSink<u64> = SampleSink::with_capacity(8);
let start = Instant::now();
for i in 0..10_000 {
sink.offer(i);
}
let elapsed = start.elapsed();
assert!(
elapsed < Duration::from_secs(2),
"producing 10k samples into a full sink took {elapsed:?}; offer() must not wait"
);
let stats = sink.stats();
assert_eq!(stats.accepted, 8, "only capacity should be retained");
assert_eq!(stats.dropped, 9_992);
assert_eq!(stats.offered(), 10_000, "every offer must be accounted for");
}
#[test]
fn slow_consumer_causes_counted_drops_not_blocking() {
let sink: SampleSink<u64> = SampleSink::with_capacity(16);
let stop = Arc::new(AtomicBool::new(false));
let consumer = {
let sink = sink.clone();
let stop = Arc::clone(&stop);
std::thread::spawn(move || {
let mut seen = 0u64;
while !stop.load(Ordering::Relaxed) {
seen += sink.drain().len() as u64;
std::thread::sleep(Duration::from_millis(5));
}
seen + sink.drain().len() as u64
})
};
for i in 0..5_000 {
sink.offer(i);
}
stop.store(true, Ordering::Relaxed);
let consumed = consumer.join().unwrap();
let stats = sink.stats();
assert_eq!(
stats.offered(),
5_000,
"accepted + dropped must equal what was offered"
);
assert!(
stats.dropped > 0,
"a consumer sleeping 5ms per batch cannot keep up with 5k offers"
);
assert!(
consumed <= stats.accepted,
"consumed {consumed} exceeds accepted {}",
stats.accepted
);
}
#[test]
fn clones_share_the_same_queue_and_counters() {
let a = SampleSink::with_capacity(4);
let b = a.clone();
a.offer(1);
b.offer(2);
assert_eq!(b.len(), 2);
assert_eq!(a.stats().accepted, 2);
assert_eq!(a.drain(), vec![1, 2]);
assert!(b.is_empty());
}
#[test]
fn a_poisoned_lock_does_not_panic_the_producer() {
let sink: SampleSink<u64> = SampleSink::with_capacity(4);
let poisoner = {
let sink = sink.clone();
std::thread::spawn(move || {
let _guard = sink.inner.queue.lock().unwrap();
panic!("consumer died holding the lock");
})
};
assert!(poisoner.join().is_err(), "the helper thread should panic");
assert!(sink.offer(1));
assert_eq!(sink.stats().accepted, 1);
}
}