| 1 | //! The sync worker's and background's loops: threads natively, and in the browser, which |
| 2 | //! gives a module one thread, tasks on its event loop. Native waits block their thread, so |
| 3 | //! `complete` runs the block to the end in one poll. |
| 4 | |
| 5 | use std::{io, time::Duration}; |
| 6 | |
| 7 | #[cfg(not(target_arch = "wasm32"))] |
| 8 | pub(crate) use std::{ |
| 9 | sync::mpsc::{Receiver, SyncSender as Sender}, |
| 10 | thread::JoinHandle, |
| 11 | }; |
| 12 | |
| 13 | #[cfg(all(feature = "live", not(target_arch = "wasm32")))] |
| 14 | pub(crate) type Answer<T> = std::sync::mpsc::Sender<T>; |
| 15 | |
| 16 | #[cfg(all(feature = "live", not(target_arch = "wasm32")))] |
| 17 | pub(crate) fn response<T>() -> (Answer<T>, std::sync::mpsc::Receiver<T>) { |
| 18 | std::sync::mpsc::channel() |
| 19 | } |
| 20 | |
| 21 | #[cfg(all(feature = "live", not(target_arch = "wasm32")))] |
| 22 | pub(crate) async fn answered<T>( |
| 23 | receiver: &std::sync::mpsc::Receiver<T>, |
| 24 | timeout: Duration, |
| 25 | ) -> Result<T, std::sync::mpsc::RecvTimeoutError> { |
| 26 | receiver.recv_timeout(timeout) |
| 27 | } |
| 28 | |
| 29 | /// A wake that waits, at most one, as `mpsc::sync_channel(1)` keeps one. |
| 30 | #[cfg(not(target_arch = "wasm32"))] |
| 31 | pub(crate) fn channel() -> (Sender<()>, Receiver<()>) { |
| 32 | std::sync::mpsc::sync_channel(1) |
| 33 | } |
| 34 | |
| 35 | /// Waits for a wake or `timeout`, without one for ever. |
| 36 | #[cfg(not(target_arch = "wasm32"))] |
| 37 | pub(crate) async fn wait(receiver: &Receiver<()>, timeout: Option<Duration>) { |
| 38 | let _ = match timeout { |
| 39 | Some(timeout) => receiver.recv_timeout(timeout).ok(), |
| 40 | None => receiver.recv().ok(), |
| 41 | }; |
| 42 | } |
| 43 | |
| 44 | /// Runs the loop `start` makes, whose waits block, on a thread of its own named `name`. |
| 45 | #[cfg(not(target_arch = "wasm32"))] |
| 46 | pub(crate) fn spawn<T: Send + 'static, F: Future<Output = T>>( |
| 47 | name: &str, |
| 48 | start: impl FnOnce() -> F + Send + 'static, |
| 49 | ) -> io::Result<JoinHandle<T>> { |
| 50 | std::thread::Builder::new() |
| 51 | .name(name.into()) |
| 52 | .spawn(move || complete(start())) |
| 53 | } |
| 54 | |
| 55 | /// `work`'s output, where nothing it awaits is pending. |
| 56 | #[cfg(not(target_arch = "wasm32"))] |
| 57 | pub(crate) fn complete<T>(work: impl Future<Output = T>) -> T { |
| 58 | let mut work = std::pin::pin!(work); |
| 59 | match work |
| 60 | .as_mut() |
| 61 | .poll(&mut std::task::Context::from_waker(std::task::Waker::noop())) |
| 62 | { |
| 63 | std::task::Poll::Ready(output) => output, |
| 64 | std::task::Poll::Pending => unreachable!("Native waits block rather than pend"), |
| 65 | } |
| 66 | } |
| 67 | |
| 68 | /// Polls a synchronous entry point once; browser I/O must use the async entry point. |
| 69 | pub(crate) fn ready<T>(work: impl Future<Output = T>) -> io::Result<T> { |
| 70 | let mut work = std::pin::pin!(work); |
| 71 | match work |
| 72 | .as_mut() |
| 73 | .poll(&mut std::task::Context::from_waker(std::task::Waker::noop())) |
| 74 | { |
| 75 | std::task::Poll::Ready(output) => Ok(output), |
| 76 | std::task::Poll::Pending => Err(io::ErrorKind::WouldBlock.into()), |
| 77 | } |
| 78 | } |
| 79 | |
| 80 | #[cfg(target_arch = "wasm32")] |
| 81 | pub(crate) use web::*; |
| 82 | |
| 83 | #[cfg(target_arch = "wasm32")] |
| 84 | mod web { |
| 85 | use super::*; |
| 86 | use std::{ |
| 87 | sync::{Arc, Mutex}, |
| 88 | task::{Poll, Waker}, |
| 89 | }; |
| 90 | use wasm_bindgen::{JsCast, prelude::*}; |
| 91 | |
| 92 | struct Response<T> { |
| 93 | value: Option<T>, |
| 94 | closed: bool, |
| 95 | waiting: Option<Waker>, |
| 96 | } |
| 97 | |
| 98 | pub(crate) struct Answer<T>(Arc<Mutex<Response<T>>>); |
| 99 | pub(crate) struct Answered<T>(Arc<Mutex<Response<T>>>); |
| 100 | |
| 101 | pub(crate) fn response<T>() -> (Answer<T>, Answered<T>) { |
| 102 | let response = Arc::new(Mutex::new(Response { |
| 103 | value: None, |
| 104 | closed: false, |
| 105 | waiting: None, |
| 106 | })); |
| 107 | (Answer(response.clone()), Answered(response)) |
| 108 | } |
| 109 | |
| 110 | impl<T> Answer<T> { |
| 111 | pub(crate) fn send(self, value: T) -> Result<(), ()> { |
| 112 | let wake = { |
| 113 | let mut response = self.0.lock().map_err(|_| ())?; |
| 114 | response.value = Some(value); |
| 115 | response.waiting.take() |
| 116 | }; |
| 117 | if let Some(wake) = wake { |
| 118 | wake.wake(); |
| 119 | } |
| 120 | Ok(()) |
| 121 | } |
| 122 | } |
| 123 | |
| 124 | impl<T> Drop for Answer<T> { |
| 125 | fn drop(&mut self) { |
| 126 | let wake = { |
| 127 | let mut response = self.0.lock().unwrap(); |
| 128 | response.closed = true; |
| 129 | response.waiting.take() |
| 130 | }; |
| 131 | if let Some(wake) = wake { |
| 132 | wake.wake(); |
| 133 | } |
| 134 | } |
| 135 | } |
| 136 | |
| 137 | pub(crate) async fn answered<T>( |
| 138 | receiver: &Answered<T>, |
| 139 | timeout: Duration, |
| 140 | ) -> Result<T, std::sync::mpsc::RecvTimeoutError> { |
| 141 | let deadline = web_time::Instant::now() + timeout; |
| 142 | let mut armed = false; |
| 143 | std::future::poll_fn(|context| { |
| 144 | let mut response = receiver.0.lock().unwrap(); |
| 145 | if let Some(value) = response.value.take() { |
| 146 | return Poll::Ready(Ok(value)); |
| 147 | } |
| 148 | if response.closed { |
| 149 | return Poll::Ready(Err(std::sync::mpsc::RecvTimeoutError::Disconnected)); |
| 150 | } |
| 151 | if web_time::Instant::now() >= deadline { |
| 152 | return Poll::Ready(Err(std::sync::mpsc::RecvTimeoutError::Timeout)); |
| 153 | } |
| 154 | response.waiting = Some(context.waker().clone()); |
| 155 | if !std::mem::replace(&mut armed, true) { |
| 156 | after(timeout, context.waker().clone()); |
| 157 | } |
| 158 | Poll::Pending |
| 159 | }) |
| 160 | .await |
| 161 | } |
| 162 | |
| 163 | #[derive(Default)] |
| 164 | struct Bell { |
| 165 | rung: bool, |
| 166 | waiting: Option<Waker>, |
| 167 | } |
| 168 | |
| 169 | pub(crate) struct Sender<T>(Arc<Mutex<Bell>>, std::marker::PhantomData<T>); |
| 170 | pub(crate) struct Receiver<T>(Arc<Mutex<Bell>>, std::marker::PhantomData<T>); |
| 171 | |
| 172 | pub(crate) fn channel() -> (Sender<()>, Receiver<()>) { |
| 173 | let bell = Arc::new(Mutex::new(Bell::default())); |
| 174 | ( |
| 175 | Sender(Arc::clone(&bell), Default::default()), |
| 176 | Receiver(bell, Default::default()), |
| 177 | ) |
| 178 | } |
| 179 | |
| 180 | impl Sender<()> { |
| 181 | pub(crate) fn try_send(&self, (): ()) -> Result<(), ()> { |
| 182 | let waiting = { |
| 183 | let mut bell = self.0.lock().map_err(|_| ())?; |
| 184 | bell.rung = true; |
| 185 | bell.waiting.take() |
| 186 | }; |
| 187 | if let Some(waiting) = waiting { |
| 188 | waiting.wake(); |
| 189 | } |
| 190 | Ok(()) |
| 191 | } |
| 192 | } |
| 193 | |
| 194 | /// Wakes `waker` after `delay` through the global scope's `setTimeout`. |
| 195 | fn after(delay: Duration, waker: Waker) { |
| 196 | let global = js_sys::global(); |
| 197 | let Ok(set_timeout) = js_sys::Reflect::get(&global, &"setTimeout".into()) |
| 198 | .and_then(|function| function.dyn_into::<js_sys::Function>()) |
| 199 | else { |
| 200 | return waker.wake(); |
| 201 | }; |
| 202 | let wake = Closure::once_into_js(move || waker.wake()); |
| 203 | let _ = set_timeout.call2(&global, &wake, &(delay.as_secs_f64() * 1e3).into()); |
| 204 | } |
| 205 | |
| 206 | pub(crate) async fn wait(receiver: &Receiver<()>, timeout: Option<Duration>) { |
| 207 | let deadline = timeout.map(|timeout| web_time::Instant::now() + timeout); |
| 208 | let mut armed = false; |
| 209 | std::future::poll_fn(|context| { |
| 210 | let Ok(mut bell) = receiver.0.lock() else { |
| 211 | return Poll::Ready(()); |
| 212 | }; |
| 213 | if std::mem::take(&mut bell.rung) |
| 214 | || deadline.is_some_and(|deadline| deadline <= web_time::Instant::now()) |
| 215 | { |
| 216 | return Poll::Ready(()); |
| 217 | } |
| 218 | bell.waiting = Some(context.waker().clone()); |
| 219 | if let Some(deadline) = deadline |
| 220 | && !std::mem::replace(&mut armed, true) |
| 221 | { |
| 222 | after( |
| 223 | deadline.saturating_duration_since(web_time::Instant::now()), |
| 224 | context.waker().clone(), |
| 225 | ); |
| 226 | } |
| 227 | Poll::Pending |
| 228 | }) |
| 229 | .await |
| 230 | } |
| 231 | |
| 232 | /// A task the browser's event loop runs. It cannot be waited for, and ends on its own. |
| 233 | pub(crate) struct JoinHandle<T>(Arc<Mutex<Option<T>>>); |
| 234 | |
| 235 | impl<T> JoinHandle<T> { |
| 236 | /// What the task ended with, if it has ended. |
| 237 | pub(crate) fn finished(self) -> Option<T> { |
| 238 | self.0.lock().ok()?.take() |
| 239 | } |
| 240 | } |
| 241 | |
| 242 | pub(crate) fn spawn<T: 'static, F: Future<Output = T> + 'static>( |
| 243 | _: &str, |
| 244 | start: impl FnOnce() -> F, |
| 245 | ) -> io::Result<JoinHandle<T>> { |
| 246 | let output = Arc::new(Mutex::new(None)); |
| 247 | let ended = Arc::clone(&output); |
| 248 | let work = start(); |
| 249 | wasm_bindgen_futures::spawn_local(async move { |
| 250 | let result = work.await; |
| 251 | if let Ok(mut ended) = ended.lock() { |
| 252 | *ended = Some(result); |
| 253 | } |
| 254 | }); |
| 255 | Ok(JoinHandle(output)) |
| 256 | } |
| 257 | } |