use core::{
future::Future,
pin::Pin,
task::{Context, Poll},
};
use heapless::mpmc;
#[derive(Clone)]
pub struct Consumer<'a, T, const N: usize> {
queue: &'a mpmc::MpMcQueue<T, N>,
}
impl<'a, T, const N: usize> Consumer<'a, T, N> {
pub async fn recv(&'a self) -> T
where
T: Copy,
{
ConsumeOnce { queue: self.queue }.await
}
pub fn try_recv(&'a self) -> Option<T>
where
T: Copy,
{
self.queue.dequeue()
}
}
struct ConsumeOnce<'a, T, const N: usize> {
queue: &'a mpmc::MpMcQueue<T, N>,
}
impl<'a, T, const N: usize> Future for ConsumeOnce<'a, T, N>
where
T: Copy,
{
type Output = T;
fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Self::Output> {
match self.queue.dequeue() {
Some(value) => Poll::Ready(value),
None => Poll::Pending,
}
}
}
#[derive(Clone)]
pub struct Producer<'a, T, const N: usize> {
queue: &'a mpmc::MpMcQueue<T, N>,
}
impl<'a, T, const N: usize> Producer<'a, T, N> {
pub async fn send(&'a self, item: T)
where
T: Copy,
{
ProduceOnce {
queue: self.queue,
item,
}
.await
}
pub fn try_send(&'a self, item: T) -> Result<(), T>
where
T: Copy,
{
self.queue.enqueue(item)
}
}
struct ProduceOnce<'a, T, const N: usize> {
queue: &'a mpmc::MpMcQueue<T, N>,
item: T,
}
impl<'a, T, const N: usize> Future for ProduceOnce<'a, T, N>
where
T: Copy,
{
type Output = ();
fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Self::Output> {
match self.queue.enqueue(self.item) {
Ok(()) => Poll::Ready(()),
Err(_) => Poll::Pending,
}
}
}
pub fn new_channel<'a, T, const N: usize>(
queue: &'a mpmc::MpMcQueue<T, N>,
) -> (Producer<'a, T, N>, Consumer<'a, T, N>) {
(Producer { queue }, Consumer { queue })
}
#[macro_export]
macro_rules! channel {
($name:ident, $type:ty, $size:expr) => {
let q = heapless::mpmc::MpMcQueue::<$type, $size>::new();
let $name = $crate::channels::new_channel(&q);
};
}