use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Weak};
use beamr::atom::Atom;
use beamr::scheduler::Scheduler;
#[derive(Clone)]
pub struct ReadyWaker {
scheduler: Weak<Scheduler>,
pid: u64,
ready_atom: Atom,
ready_pending: Arc<AtomicBool>,
#[cfg(test)]
fire_probe: Option<std::sync::Arc<std::sync::atomic::AtomicU64>>,
}
impl std::fmt::Debug for ReadyWaker {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("ReadyWaker")
.field("pid", &self.pid)
.field("scheduler_live", &(self.scheduler.strong_count() > 0))
.finish_non_exhaustive()
}
}
impl ReadyWaker {
pub(crate) fn new(
scheduler: &Arc<Scheduler>,
pid: u64,
ready_atom: Atom,
ready_pending: Arc<AtomicBool>,
) -> Self {
Self {
scheduler: std::sync::Arc::downgrade(scheduler),
pid,
ready_atom,
ready_pending,
#[cfg(test)]
fire_probe: None,
}
}
#[cfg(test)]
pub(crate) fn for_test(probe: Arc<std::sync::atomic::AtomicU64>) -> Self {
Self {
scheduler: Weak::new(),
pid: 0,
ready_atom: Atom::OK,
ready_pending: Arc::new(AtomicBool::new(false)),
fire_probe: Some(probe),
}
}
pub(crate) fn fire(&self) -> bool {
#[cfg(test)]
if let Some(probe) = &self.fire_probe {
probe.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
return true;
}
let Some(scheduler) = self.scheduler.upgrade() else {
return false;
};
self.ready_pending.store(true, Ordering::Release);
scheduler.enqueue_atom_message(self.pid, self.ready_atom)
}
}