use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::Arc;
use concurrent_queue::ConcurrentQueue;
use foundation_core::valtron::{EventReadiness, QueueReadiness};
use crate::types::Messages;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u32)]
pub enum CancelCode {
None = 0,
PauseForPriority = 1,
Abort = 2,
}
impl CancelCode {
pub fn load(signal: &Arc<AtomicU32>) -> Self {
match signal.load(Ordering::SeqCst) {
1 => CancelCode::PauseForPriority,
2 => CancelCode::Abort,
_ => CancelCode::None,
}
}
pub fn store(self, signal: &Arc<AtomicU32>) {
signal.store(self as u32, Ordering::SeqCst);
}
pub fn reset(signal: &Arc<AtomicU32>) {
CancelCode::None.store(signal);
}
#[must_use]
pub fn is_active(self) -> bool {
self != CancelCode::None
}
}
pub struct SteeringQueues {
pub priority: Arc<ConcurrentQueue<Messages>>,
pub follow_up: Arc<ConcurrentQueue<Messages>>,
pub cancel_signal: Arc<AtomicU32>,
}
impl SteeringQueues {
#[must_use]
pub fn new() -> Self {
Self {
priority: Arc::new(ConcurrentQueue::unbounded()),
follow_up: Arc::new(ConcurrentQueue::unbounded()),
cancel_signal: Arc::new(AtomicU32::new(CancelCode::None as u32)),
}
}
#[must_use]
pub fn from_shared(
priority: Arc<ConcurrentQueue<Messages>>,
follow_up: Arc<ConcurrentQueue<Messages>>,
cancel_signal: Arc<AtomicU32>,
) -> Self {
Self {
priority,
follow_up,
cancel_signal,
}
}
pub fn push_priority(&self, msg: Messages) {
let _ = self.priority.push(msg);
CancelCode::PauseForPriority.store(&self.cancel_signal);
}
pub fn push_follow_up(&self, msg: Messages) {
let _ = self.follow_up.push(msg);
}
#[must_use]
pub fn pop_priority(&self) -> Option<Messages> {
self.priority.pop().ok()
}
#[must_use]
pub fn pop_follow_up(&self) -> Option<Messages> {
self.follow_up.pop().ok()
}
#[must_use]
pub fn has_priority(&self) -> bool {
!self.priority.is_empty()
}
#[must_use]
pub fn has_follow_up(&self) -> bool {
!self.follow_up.is_empty()
}
#[must_use]
pub fn cancel_code(&self) -> CancelCode {
CancelCode::load(&self.cancel_signal)
}
pub fn abort(&self) {
CancelCode::Abort.store(&self.cancel_signal);
}
#[must_use]
pub fn is_aborted(&self) -> bool {
self.cancel_code() == CancelCode::Abort
}
pub fn reset_cancel(&self) {
CancelCode::reset(&self.cancel_signal);
}
#[must_use]
pub fn is_active(&self) -> bool {
self.has_priority() || self.cancel_code().is_active()
}
#[must_use]
pub fn priority_readiness(&self) -> Arc<dyn EventReadiness> {
Arc::new(QueueReadiness::new(self.priority.clone()))
}
#[must_use]
pub fn followup_readiness(&self) -> Arc<dyn EventReadiness> {
Arc::new(QueueReadiness::new(self.follow_up.clone()))
}
#[must_use]
pub fn drain_priority(&self) -> Vec<Messages> {
let mut out = Vec::new();
while let Ok(msg) = self.priority.pop() {
out.push(msg);
}
out
}
#[must_use]
pub fn drain_follow_up(&self) -> Vec<Messages> {
let mut out = Vec::new();
while let Ok(msg) = self.follow_up.pop() {
out.push(msg);
}
out
}
}
impl Default for SteeringQueues {
fn default() -> Self {
Self::new()
}
}