use core::hint;
use core::mem;
use core::ops::Deref;
use alloc::sync::Arc;
use alloc::sync::Weak;
use once_cell::unsync::Lazy;
use parking_lot::Mutex;
use parking_lot::RwLock;
use parking_lot::RwLockReadGuard;
pub struct Listener<T> {
trackers: Mutex<Vec<Weak<RwLock<Vec<T>>>>>,
}
impl<T> Listener<T> {
pub fn emit(&self, event: &T)
where
T: Clone,
{
self.emit_with(|| event.clone());
}
pub fn emit_with(&self, mut event_source: impl FnMut() -> T) {
let should_prune = self
.trackers
.lock()
.iter()
.fold(false, |should_prune, tracker| {
let Some(tracker) = tracker.upgrade() else {
return true;
};
tracker.write().push(event_source());
should_prune
});
if should_prune {
self.prune();
}
}
pub fn emit_with_cached(&self, cached_event_source: impl FnOnce() -> T)
where
T: Clone,
{
self.emit_with({
let initializer = Lazy::new(cached_event_source);
move || initializer.clone()
});
}
pub fn track(&self) -> Tracker<T> {
let events = Arc::default();
self.trackers.lock().push(Arc::downgrade(&events));
Tracker { events }
}
}
impl<T> Listener<T> {
fn prune(&self) {
self.trackers.lock().retain(|arc| 0 < arc.strong_count());
}
}
impl<T> Default for Listener<T> {
fn default() -> Self {
let trackers = Mutex::default();
Self { trackers }
}
}
#[derive(Clone, Debug)]
pub struct Tracker<T> {
events: Arc<RwLock<Vec<T>>>,
}
impl<T> Tracker<T> {
#[must_use]
pub fn consume(&self) -> Vec<T> {
let mut events = self.events.write();
mem::take(&mut events)
}
#[must_use]
pub fn data(&self) -> TrackerData<T> {
TrackerData(self.events.read())
}
#[must_use]
pub fn stop(mut self) -> Vec<T> {
loop {
match Arc::try_unwrap(self.events) {
Ok(events) => return events.into_inner(),
Err(ev) => {
self.events = ev;
hint::spin_loop();
}
};
}
}
}
#[derive(Debug)]
pub struct TrackerData<'l, T>(RwLockReadGuard<'l, Vec<T>>);
impl<T> Deref for TrackerData<'_, T> {
type Target = [T];
fn deref(&self) -> &Self::Target {
&self.0
}
}