2023-10-16 15:50:22 +02:00
|
|
|
use std::future;
|
2024-06-12 00:02:26 +02:00
|
|
|
use std::num::NonZeroUsize;
|
2023-10-16 15:50:22 +02:00
|
|
|
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
|
|
|
|
use std::sync::Arc;
|
|
|
|
|
use std::task::Poll;
|
|
|
|
|
|
2024-01-12 11:51:19 +01:00
|
|
|
use super::super::main_thread::MainThreadMarker;
|
|
|
|
|
use super::{AtomicWaker, Wrapper};
|
|
|
|
|
|
2024-06-12 00:02:26 +02:00
|
|
|
pub struct WakerSpawner<T: 'static>(Wrapper<Handler<T>, Sender, NonZeroUsize>);
|
2023-10-16 15:50:22 +02:00
|
|
|
|
2024-06-12 00:02:26 +02:00
|
|
|
pub struct Waker<T: 'static>(Wrapper<Handler<T>, Sender, NonZeroUsize>);
|
2023-10-16 15:50:22 +02:00
|
|
|
|
|
|
|
|
struct Handler<T> {
|
|
|
|
|
value: T,
|
2024-06-12 00:02:26 +02:00
|
|
|
handler: fn(&T, NonZeroUsize, bool),
|
2023-10-16 15:50:22 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[derive(Clone)]
|
|
|
|
|
struct Sender(Arc<Inner>);
|
|
|
|
|
|
|
|
|
|
impl<T> WakerSpawner<T> {
|
|
|
|
|
#[track_caller]
|
2024-06-12 00:02:26 +02:00
|
|
|
pub fn new(
|
|
|
|
|
main_thread: MainThreadMarker,
|
|
|
|
|
value: T,
|
|
|
|
|
handler: fn(&T, NonZeroUsize, bool),
|
|
|
|
|
) -> Option<Self> {
|
2023-10-16 15:50:22 +02:00
|
|
|
let inner = Arc::new(Inner {
|
|
|
|
|
counter: AtomicUsize::new(0),
|
|
|
|
|
waker: AtomicWaker::new(),
|
|
|
|
|
closed: AtomicBool::new(false),
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
let handler = Handler { value, handler };
|
|
|
|
|
|
|
|
|
|
let sender = Sender(Arc::clone(&inner));
|
|
|
|
|
|
|
|
|
|
let wrapper = Wrapper::new(
|
2023-12-26 01:22:10 +01:00
|
|
|
main_thread,
|
2023-10-16 15:50:22 +02:00
|
|
|
handler,
|
|
|
|
|
|handler, count| {
|
|
|
|
|
let handler = handler.borrow();
|
|
|
|
|
let handler = handler.as_ref().unwrap();
|
2024-06-12 00:02:26 +02:00
|
|
|
(handler.handler)(&handler.value, count, true);
|
2023-10-16 15:50:22 +02:00
|
|
|
},
|
|
|
|
|
{
|
|
|
|
|
let inner = Arc::clone(&inner);
|
|
|
|
|
|
|
|
|
|
move |handler| async move {
|
|
|
|
|
while let Some(count) = future::poll_fn(|cx| {
|
|
|
|
|
let count = inner.counter.swap(0, Ordering::Relaxed);
|
|
|
|
|
|
2024-06-12 00:02:26 +02:00
|
|
|
match NonZeroUsize::new(count) {
|
|
|
|
|
Some(count) => Poll::Ready(Some(count)),
|
|
|
|
|
None => {
|
|
|
|
|
inner.waker.register(cx.waker());
|
2023-10-16 15:50:22 +02:00
|
|
|
|
2024-06-12 00:02:26 +02:00
|
|
|
let count = inner.counter.swap(0, Ordering::Relaxed);
|
2023-10-16 15:50:22 +02:00
|
|
|
|
2024-06-12 00:02:26 +02:00
|
|
|
match NonZeroUsize::new(count) {
|
|
|
|
|
Some(count) => Poll::Ready(Some(count)),
|
|
|
|
|
None => {
|
|
|
|
|
if inner.closed.load(Ordering::Relaxed) {
|
|
|
|
|
return Poll::Ready(None);
|
|
|
|
|
}
|
2023-10-16 15:50:22 +02:00
|
|
|
|
2024-06-12 00:02:26 +02:00
|
|
|
Poll::Pending
|
|
|
|
|
},
|
|
|
|
|
}
|
|
|
|
|
},
|
2023-10-16 15:50:22 +02:00
|
|
|
}
|
|
|
|
|
})
|
|
|
|
|
.await
|
|
|
|
|
{
|
|
|
|
|
let handler = handler.borrow();
|
|
|
|
|
let handler = handler.as_ref().unwrap();
|
2024-06-12 00:02:26 +02:00
|
|
|
(handler.handler)(&handler.value, count, false);
|
2023-10-16 15:50:22 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
},
|
|
|
|
|
sender,
|
|
|
|
|
|inner, _| {
|
|
|
|
|
inner.0.counter.fetch_add(1, Ordering::Relaxed);
|
|
|
|
|
inner.0.waker.wake();
|
|
|
|
|
},
|
|
|
|
|
)?;
|
|
|
|
|
|
|
|
|
|
Some(Self(wrapper))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn waker(&self) -> Waker<T> {
|
|
|
|
|
Waker(self.0.clone())
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn fetch(&self) -> usize {
|
|
|
|
|
debug_assert!(
|
2023-12-26 01:22:10 +01:00
|
|
|
MainThreadMarker::new().is_some(),
|
2023-10-16 15:50:22 +02:00
|
|
|
"this should only be called from the main thread"
|
|
|
|
|
);
|
|
|
|
|
|
|
|
|
|
self.0.with_sender_data(|inner| inner.0.counter.swap(0, Ordering::Relaxed))
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl<T> Drop for WakerSpawner<T> {
|
|
|
|
|
fn drop(&mut self) {
|
|
|
|
|
self.0.with_sender_data(|inner| {
|
|
|
|
|
inner.0.closed.store(true, Ordering::Relaxed);
|
|
|
|
|
inner.0.waker.wake();
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl<T> Waker<T> {
|
|
|
|
|
pub fn wake(&self) {
|
2024-06-12 00:02:26 +02:00
|
|
|
self.0.send(NonZeroUsize::MIN)
|
2023-10-16 15:50:22 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl<T> Clone for Waker<T> {
|
|
|
|
|
fn clone(&self) -> Self {
|
|
|
|
|
Self(self.0.clone())
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
struct Inner {
|
|
|
|
|
counter: AtomicUsize,
|
|
|
|
|
waker: AtomicWaker,
|
|
|
|
|
closed: AtomicBool,
|
|
|
|
|
}
|