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
5use std::{io, time::Duration};
6
7#[cfg(not(target_arch = "wasm32"))]
8pub(crate) use std::{
9 sync::mpsc::{Receiver, SyncSender as Sender},
10 thread::JoinHandle,
11};
12
13#[cfg(all(feature = "live", not(target_arch = "wasm32")))]
14pub(crate) type Answer<T> = std::sync::mpsc::Sender<T>;
15
16#[cfg(all(feature = "live", not(target_arch = "wasm32")))]
17pub(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")))]
22pub(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"))]
31pub(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"))]
37pub(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"))]
46pub(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"))]
57pub(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.
69pub(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")]
81pub(crate) use web::*;
82
83#[cfg(target_arch = "wasm32")]
84mod 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}