#![allow(clippy::needless_doctest_main)]
use async_channel::Receiver;
use async_channel::TryRecvError;
use async_channel::unbounded;
use futures_lite::Stream;
use std::pin::Pin;
use std::rc::Rc;
use std::task::Context;
use std::task::Poll;
use wasm_bindgen::JsCast;
use wasm_bindgen::closure::Closure;
use web_sys::Window;
use web_sys::window;
struct IntervalInner {
target: Window,
handle: i32,
_closure: Closure<dyn FnMut()>,
}
impl Drop for IntervalInner {
fn drop(&mut self) {
self.target.clear_interval_with_handle(self.handle);
}
}
pub struct IntervalStream {
receiver: Receiver<f64>,
_inner: Rc<IntervalInner>,
}
impl Unpin for IntervalStream {}
impl IntervalStream {
pub fn new(interval_ms: i32) -> Self {
let target: Window = window().unwrap();
Self::new_with_target(interval_ms, target)
}
pub fn new_with_target(interval_ms: i32, target: Window) -> Self {
let (sender, receiver) = unbounded::<f64>();
let closure = Closure::wrap(Box::new(move || {
let ts = window()
.and_then(|w| w.performance())
.map(|p| p.now())
.unwrap_or(0.0);
let _ = sender.try_send(ts);
}) as Box<dyn FnMut()>);
let handle = target
.set_interval_with_callback_and_timeout_and_arguments_0(
closure.as_ref().unchecked_ref(),
interval_ms,
)
.unwrap();
let inner = Rc::new(IntervalInner {
target,
handle,
_closure: closure,
});
Self {
receiver,
_inner: inner,
}
}
pub fn forget(self) { std::mem::forget(self); }
pub async fn next_tick(&mut self) -> Option<f64> {
self.receiver.recv().await.ok()
}
}
impl Stream for IntervalStream {
type Item = f64;
fn poll_next(
self: Pin<&mut Self>,
cx: &mut Context<'_>,
) -> Poll<Option<Self::Item>> {
let this = self.get_mut();
match this.receiver.try_recv() {
Ok(item) => return Poll::Ready(Some(item)),
Err(TryRecvError::Closed) => return Poll::Ready(None),
Err(TryRecvError::Empty) => {}
}
let recv = this.receiver.clone();
let fut = recv.recv();
futures_lite::pin!(fut);
match fut.poll(cx) {
Poll::Ready(Ok(item)) => Poll::Ready(Some(item)),
Poll::Ready(Err(_closed)) => Poll::Ready(None),
Poll::Pending => Poll::Pending,
}
}
}
#[cfg(test)]
#[cfg(target_arch = "wasm32")]
mod tests {
use super::IntervalStream;
use crate::prelude::*;
use crate::web_utils::document_ext as doc;
#[ignore = "requires dom"]
#[test]
fn works() {
let _ = doc::document();
let _ = doc::head();
let _ = doc::body();
}
#[ignore = "requires dom"]
#[crate::test]
async fn yields_timestamps() {
doc::clear_body();
let mut interval = IntervalStream::new(10);
let a = interval.next_tick().await.unwrap();
let b = interval.next_tick().await.unwrap();
(b >= a).xpect_true();
}
}