use std::collections::VecDeque;
use std::sync::{Arc, Mutex};
use web_time::Instant;
use crate::bus::TraceSubscriber;
use crate::event::RosaceTrace;
type EventFilter = Arc<dyn Fn(&RosaceTrace) -> bool + Send + Sync>;
pub struct RingBufferSubscriber {
buffer: Arc<Mutex<VecDeque<(Instant, RosaceTrace)>>>,
capacity: usize,
filter: Option<EventFilter>,
}
impl RingBufferSubscriber {
pub fn new(capacity: usize) -> Self {
Self {
buffer: Arc::new(Mutex::new(VecDeque::with_capacity(capacity))),
capacity,
filter: None,
}
}
pub fn filtered(
capacity: usize,
filter: impl Fn(&RosaceTrace) -> bool + Send + Sync + 'static,
) -> Self {
Self {
buffer: Arc::new(Mutex::new(VecDeque::with_capacity(capacity))),
capacity,
filter: Some(Arc::new(filter)),
}
}
pub fn len(&self) -> usize {
self.buffer
.lock()
.expect("RingBufferSubscriber lock poisoned")
.len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn drain(&self) -> Vec<RosaceTrace> {
self.buffer
.lock()
.expect("RingBufferSubscriber lock poisoned")
.drain(..)
.map(|(_, e)| e)
.collect()
}
pub fn snapshot(&self) -> Vec<RosaceTrace> {
self.buffer
.lock()
.expect("RingBufferSubscriber lock poisoned")
.iter()
.map(|(_, e)| e.clone())
.collect()
}
pub fn snapshot_timestamped(&self) -> Vec<(Instant, RosaceTrace)> {
self.buffer
.lock()
.expect("RingBufferSubscriber lock poisoned")
.iter()
.cloned()
.collect()
}
pub fn export_perfetto_json(&self) -> String {
super::perfetto::to_chrome_trace_json(&self.snapshot_timestamped())
}
}
impl TraceSubscriber for RingBufferSubscriber {
fn on_trace(&self, event: &RosaceTrace) {
if let Some(f) = &self.filter {
if !f(event) {
return;
}
}
let mut buf = self
.buffer
.lock()
.expect("RingBufferSubscriber lock poisoned");
if buf.len() >= self.capacity {
buf.pop_front();
}
buf.push_back((Instant::now(), event.clone()));
}
}