use std::collections::VecDeque;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use super::scheduler::{wake_handle, WakeHandle};
use super::scope::Scope;
use super::signal::Signal;
pub type CoalesceFn<T> = Arc<dyn Fn(&mut T, T) + Send + Sync>;
pub enum OverflowPolicy<T> {
DropOldest,
DropNewest,
Coalesce(CoalesceFn<T>),
}
impl<T> OverflowPolicy<T> {
pub fn coalesce(fold: impl Fn(&mut T, T) + Send + Sync + 'static) -> OverflowPolicy<T> {
OverflowPolicy::Coalesce(Arc::new(fold))
}
}
impl<T> Clone for OverflowPolicy<T> {
fn clone(&self) -> Self {
match self {
OverflowPolicy::DropOldest => OverflowPolicy::DropOldest,
OverflowPolicy::DropNewest => OverflowPolicy::DropNewest,
OverflowPolicy::Coalesce(f) => OverflowPolicy::Coalesce(f.clone()),
}
}
}
impl<T> std::fmt::Debug for OverflowPolicy<T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(match self {
OverflowPolicy::DropOldest => "DropOldest",
OverflowPolicy::DropNewest => "DropNewest",
OverflowPolicy::Coalesce(_) => "Coalesce(..)",
})
}
}
#[derive(Copy, Clone, Debug, Default, PartialEq, Eq)]
pub struct IngestStats {
pub delivered: u64,
pub dropped: u64,
pub coalesced: u64,
pub fold_panics: u64,
}
struct BoundedShared<T> {
wake: WakeHandle,
queue: Mutex<VecDeque<T>>,
capacity: usize,
policy: OverflowPolicy<T>,
drain_scheduled: AtomicBool,
dropped: AtomicU64,
coalesced: AtomicU64,
fold_panics: AtomicU64,
dead_sends: AtomicU64,
events: Signal<Vec<T>>,
stats: Signal<IngestStats>,
}
pub struct BoundedSender<T> {
shared: Arc<BoundedShared<T>>,
}
impl<T> Clone for BoundedSender<T> {
fn clone(&self) -> Self {
BoundedSender {
shared: self.shared.clone(),
}
}
}
fn run_fold<T>(fold: &CoalesceFn<T>, target: &mut T, value: T) -> bool {
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| fold(target, value))).is_ok()
}
fn admit_transit<T>(
buf: &mut VecDeque<T>,
value: T,
capacity: usize,
policy: &OverflowPolicy<T>,
dropped: &mut u64,
coalesced: &mut u64,
fold_panics: &mut u64,
) {
if buf.len() < capacity {
buf.push_back(value);
return;
}
match policy {
OverflowPolicy::DropOldest => {
buf.pop_front();
buf.push_back(value);
*dropped += 1;
}
OverflowPolicy::DropNewest => {
*dropped += 1;
}
OverflowPolicy::Coalesce(fold) => {
let newest = buf.back_mut().expect("capacity >= 1: full is non-empty");
if run_fold(fold, newest, value) {
*coalesced += 1;
} else {
*dropped += 1; *fold_panics += 1;
}
}
}
}
impl<T: Send + 'static> BoundedSender<T> {
pub fn send(&self, value: T) {
let shared = &self.shared;
{
let mut queue = shared.queue.lock().expect("bounded-source queue");
let (mut d, mut c, mut p) = (0u64, 0u64, 0u64);
admit_transit(
&mut queue,
value,
shared.capacity,
&shared.policy,
&mut d,
&mut c,
&mut p,
);
if d > 0 {
shared.dropped.fetch_add(d, Ordering::Relaxed);
}
if c > 0 {
shared.coalesced.fetch_add(c, Ordering::Relaxed);
}
if p > 0 {
shared.fold_panics.fetch_add(p, Ordering::Relaxed);
}
} if !shared.drain_scheduled.swap(true, Ordering::AcqRel) {
let shared = shared.clone();
self.shared.wake.post(move || drain(&shared));
}
}
pub fn dead_sends(&self) -> u64 {
self.shared.dead_sends.load(Ordering::Relaxed)
}
}
fn drain<T: 'static>(shared: &Arc<BoundedShared<T>>) {
shared.drain_scheduled.store(false, Ordering::Release);
let batch: VecDeque<T> = {
let mut queue = shared.queue.lock().expect("bounded-source queue");
std::mem::take(&mut *queue)
}; if batch.is_empty() {
return; }
if !shared.events.is_alive() {
shared
.dead_sends
.fetch_add(batch.len() as u64, Ordering::Relaxed);
return;
}
let (mut delivered, mut dropped, mut coalesced, mut fold_panics) = (0u64, 0u64, 0u64, 0u64);
super::runtime::batch(|| {
shared.events.update(|vec| {
let cap = shared.capacity;
let incoming = batch.len();
match &shared.policy {
OverflowPolicy::DropOldest => {
let unfittable = incoming.saturating_sub(cap);
dropped += unfittable as u64;
delivered += (incoming - unfittable) as u64;
vec.extend(batch.into_iter().skip(unfittable));
let aged = vec.len().saturating_sub(cap);
vec.drain(..aged);
}
OverflowPolicy::DropNewest => {
let room = cap.saturating_sub(vec.len());
let admitted = room.min(incoming);
delivered += admitted as u64;
dropped += (incoming - admitted) as u64;
vec.extend(batch.into_iter().take(admitted));
}
OverflowPolicy::Coalesce(fold) => {
for value in batch {
if vec.len() < cap {
vec.push(value);
delivered += 1;
} else if run_fold(fold, vec.last_mut().expect("cap >= 1"), value) {
coalesced += 1;
} else {
dropped += 1; fold_panics += 1;
}
}
}
}
});
dropped += shared.dropped.swap(0, Ordering::Relaxed);
coalesced += shared.coalesced.swap(0, Ordering::Relaxed);
fold_panics += shared.fold_panics.swap(0, Ordering::Relaxed);
if shared.stats.is_alive() {
shared.stats.update(|s| {
s.delivered += delivered;
s.dropped += dropped;
s.coalesced += coalesced;
s.fold_panics += fold_panics;
});
}
});
}
pub fn bounded_source<T: Send + 'static>(
cx: Scope,
capacity: usize,
policy: OverflowPolicy<T>,
) -> (BoundedSender<T>, Signal<Vec<T>>, Signal<IngestStats>) {
assert!(
capacity >= 1,
"abstracttui reactive: bounded_source capacity must be >= 1 — a zero-capacity \
window admits nothing. FIX: size capacity for the burst you accept (it is also \
the retained-window length) and choose the OverflowPolicy that names what \
overflow should mean"
);
let events = cx.signal(Vec::new());
let stats = cx.signal(IngestStats::default());
let sender = BoundedSender {
shared: Arc::new(BoundedShared {
wake: wake_handle(),
queue: Mutex::new(VecDeque::with_capacity(capacity.min(4096))),
capacity,
policy,
drain_scheduled: AtomicBool::new(false),
dropped: AtomicU64::new(0),
coalesced: AtomicU64::new(0),
fold_panics: AtomicU64::new(0),
dead_sends: AtomicU64::new(0),
events,
stats,
}),
};
(sender, events, stats)
}
#[cfg(test)]
#[path = "ingest_tests.rs"]
mod tests;