2023-10-16 15:50:22 +02:00
|
|
|
use std::future;
|
|
|
|
|
use std::rc::Rc;
|
|
|
|
|
use std::sync::atomic::{AtomicBool, Ordering};
|
2024-01-12 11:51:19 +01:00
|
|
|
use std::sync::mpsc::{self, RecvError, SendError, TryRecvError};
|
2023-10-16 15:50:22 +02:00
|
|
|
use std::sync::{Arc, Mutex};
|
|
|
|
|
use std::task::Poll;
|
|
|
|
|
|
2024-01-12 11:51:19 +01:00
|
|
|
use super::AtomicWaker;
|
|
|
|
|
|
|
|
|
|
pub fn channel<T>() -> (Sender<T>, Receiver<T>) {
|
2023-10-16 15:50:22 +02:00
|
|
|
let (sender, receiver) = mpsc::channel();
|
2023-12-22 22:20:41 +01:00
|
|
|
let shared = Arc::new(Shared {
|
2023-10-16 15:50:22 +02:00
|
|
|
closed: AtomicBool::new(false),
|
|
|
|
|
waker: AtomicWaker::new(),
|
|
|
|
|
});
|
|
|
|
|
|
2024-01-12 11:51:19 +01:00
|
|
|
let sender = Sender(Arc::new(SenderInner {
|
2023-12-22 22:20:41 +01:00
|
|
|
sender: Mutex::new(sender),
|
|
|
|
|
shared: Arc::clone(&shared),
|
|
|
|
|
}));
|
2024-01-12 11:51:19 +01:00
|
|
|
let receiver = Receiver {
|
2023-10-16 15:50:22 +02:00
|
|
|
receiver: Rc::new(receiver),
|
2023-12-22 22:20:41 +01:00
|
|
|
shared,
|
2023-10-16 15:50:22 +02:00
|
|
|
};
|
|
|
|
|
|
|
|
|
|
(sender, receiver)
|
|
|
|
|
}
|
|
|
|
|
|
2024-01-12 11:51:19 +01:00
|
|
|
pub struct Sender<T>(Arc<SenderInner<T>>);
|
2023-12-22 22:20:41 +01:00
|
|
|
|
|
|
|
|
struct SenderInner<T> {
|
2023-10-16 15:50:22 +02:00
|
|
|
// We need to wrap it into a `Mutex` to make it `Sync`. So the sender can't
|
|
|
|
|
// be accessed on the main thread, as it could block. Additionally we need
|
2023-12-22 22:20:41 +01:00
|
|
|
// to wrap `Sender` in an `Arc` to make it clonable on the main thread without
|
2023-10-16 15:50:22 +02:00
|
|
|
// having to block.
|
2024-01-12 11:51:19 +01:00
|
|
|
sender: Mutex<mpsc::Sender<T>>,
|
2023-12-22 22:20:41 +01:00
|
|
|
shared: Arc<Shared>,
|
2023-10-16 15:50:22 +02:00
|
|
|
}
|
|
|
|
|
|
2024-01-12 11:51:19 +01:00
|
|
|
impl<T> Sender<T> {
|
2023-10-16 15:50:22 +02:00
|
|
|
pub fn send(&self, event: T) -> Result<(), SendError<T>> {
|
2023-12-22 22:20:41 +01:00
|
|
|
self.0.sender.lock().unwrap().send(event)?;
|
|
|
|
|
self.0.shared.waker.wake();
|
2023-10-16 15:50:22 +02:00
|
|
|
|
|
|
|
|
Ok(())
|
|
|
|
|
}
|
2023-12-22 22:20:41 +01:00
|
|
|
}
|
2023-10-16 15:50:22 +02:00
|
|
|
|
2023-12-22 22:20:41 +01:00
|
|
|
impl<T> SenderInner<T> {
|
|
|
|
|
fn close(&self) {
|
|
|
|
|
self.shared.closed.store(true, Ordering::Relaxed);
|
|
|
|
|
self.shared.waker.wake();
|
2023-10-16 15:50:22 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2024-01-12 11:51:19 +01:00
|
|
|
impl<T> Clone for Sender<T> {
|
2023-10-16 15:50:22 +02:00
|
|
|
fn clone(&self) -> Self {
|
2023-12-22 22:20:41 +01:00
|
|
|
Self(Arc::clone(&self.0))
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl<T> Drop for SenderInner<T> {
|
|
|
|
|
fn drop(&mut self) {
|
|
|
|
|
self.close();
|
2023-10-16 15:50:22 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2024-01-12 11:51:19 +01:00
|
|
|
pub struct Receiver<T> {
|
|
|
|
|
receiver: Rc<mpsc::Receiver<T>>,
|
2023-12-22 22:20:41 +01:00
|
|
|
shared: Arc<Shared>,
|
2023-10-16 15:50:22 +02:00
|
|
|
}
|
|
|
|
|
|
2024-01-12 11:51:19 +01:00
|
|
|
impl<T> Receiver<T> {
|
2023-10-16 15:50:22 +02:00
|
|
|
pub async fn next(&self) -> Result<T, RecvError> {
|
|
|
|
|
future::poll_fn(|cx| match self.receiver.try_recv() {
|
|
|
|
|
Ok(event) => Poll::Ready(Ok(event)),
|
|
|
|
|
Err(TryRecvError::Empty) => {
|
2023-12-22 22:20:41 +01:00
|
|
|
self.shared.waker.register(cx.waker());
|
2023-10-16 15:50:22 +02:00
|
|
|
|
|
|
|
|
match self.receiver.try_recv() {
|
|
|
|
|
Ok(event) => Poll::Ready(Ok(event)),
|
|
|
|
|
Err(TryRecvError::Empty) => {
|
2023-12-22 22:20:41 +01:00
|
|
|
if self.shared.closed.load(Ordering::Relaxed) {
|
2023-10-16 15:50:22 +02:00
|
|
|
Poll::Ready(Err(RecvError))
|
|
|
|
|
} else {
|
|
|
|
|
Poll::Pending
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
Err(TryRecvError::Disconnected) => Poll::Ready(Err(RecvError)),
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
Err(TryRecvError::Disconnected) => Poll::Ready(Err(RecvError)),
|
|
|
|
|
})
|
|
|
|
|
.await
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn try_recv(&self) -> Result<Option<T>, RecvError> {
|
|
|
|
|
match self.receiver.try_recv() {
|
|
|
|
|
Ok(value) => Ok(Some(value)),
|
|
|
|
|
Err(TryRecvError::Empty) => Ok(None),
|
|
|
|
|
Err(TryRecvError::Disconnected) => Err(RecvError),
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2024-01-12 11:51:19 +01:00
|
|
|
impl<T> Clone for Receiver<T> {
|
2023-10-16 15:50:22 +02:00
|
|
|
fn clone(&self) -> Self {
|
|
|
|
|
Self {
|
|
|
|
|
receiver: Rc::clone(&self.receiver),
|
2023-12-22 22:20:41 +01:00
|
|
|
shared: Arc::clone(&self.shared),
|
2023-10-16 15:50:22 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2024-01-12 11:51:19 +01:00
|
|
|
impl<T> Drop for Receiver<T> {
|
2023-12-22 22:20:41 +01:00
|
|
|
fn drop(&mut self) {
|
|
|
|
|
self.shared.closed.store(true, Ordering::Relaxed);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
struct Shared {
|
2023-10-16 15:50:22 +02:00
|
|
|
closed: AtomicBool,
|
|
|
|
|
waker: AtomicWaker,
|
|
|
|
|
}
|