use super::{EventMask, Signals};
use core::convert::Infallible;
use heapless::mpmc::MpMcQueue;
use nb::Error::WouldBlock;
#[cfg_attr(docsrs, doc(cfg(feature = "heapless")))]
pub struct Mpmc<T, const N: usize>(MpMcQueue<T, N>);
impl<T, const N: usize> Default for Mpmc<T, N> {
fn default() -> Self {
Mpmc(Default::default())
}
}
impl<T, const N: usize> Mpmc<T, N> {
pub const fn new() -> Self {
Mpmc(MpMcQueue::new())
}
pub fn inner(&self) -> &MpMcQueue<T, N> {
&self.0
}
pub async fn enqueue<S: EventMask>(&self, item: T, signals: &Signals<'_, S>, ev: S) {
let mut item = Some(item);
let queued = signals.drive_infallible(ev, move || {
self.try_enqueue(item.take().ok_or(WouldBlock)?, signals, ev)
.map_err(|returned| {
item = Some(returned);
WouldBlock
})
});
queued.await
}
pub fn try_enqueue<S: EventMask>(
&self,
item: T,
signals: &Signals<'_, S>,
ev: S,
) -> Result<(), T> {
let result = self.0.enqueue(item);
if let Ok(()) = result {
signals.raise(ev);
}
result
}
pub async fn dequeue<S: EventMask>(&self, signals: &Signals<'_, S>, ev: S) -> T {
signals
.drive_infallible(ev, || self.try_dequeue(signals, ev))
.await
}
pub fn try_dequeue<S: EventMask>(
&self,
signals: &Signals<'_, S>,
ev: S,
) -> nb::Result<T, Infallible> {
let item = self.0.dequeue().ok_or(WouldBlock)?;
signals.raise(ev);
Ok(item)
}
}
impl<T, const N: usize> From<MpMcQueue<T, N>> for Mpmc<T, N> {
fn from(queue: MpMcQueue<T, N>) -> Self {
Mpmc(queue)
}
}
impl<T, const N: usize> From<Mpmc<T, N>> for MpMcQueue<T, N> {
fn from(queue: Mpmc<T, N>) -> Self {
queue.0
}
}