Skip to main content

rig_http/
wasm_compat.rs

1//! Target-dependent thread bounds, boxed futures, and executor-independent timers.
2//!
3//! ```
4//! use rig_http::wasm_compat::WasmBoxedFuture;
5//! let future: WasmBoxedFuture<'_, u32> = Box::pin(async { 42 });
6//! ```
7
8use bytes::Bytes;
9use std::pin::Pin;
10
11use futures::Stream;
12
13// Browser markers assume single-threaded execution; atomics would invalidate
14// that assumption. Relaxed bounds do not make non-Send handlers thread-safe.
15#[cfg(all(
16    target_arch = "wasm32",
17    target_os = "unknown",
18    target_feature = "atomics"
19))]
20compile_error!(
21    "rig-http does not support threaded wasm (`+atomics`): its wasm-compat markers assume a \
22     single-threaded target"
23);
24
25#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
26/// `Send` on native targets, a no-op marker on browser wasm.
27///
28/// ```compile_fail
29/// use std::rc::Rc;
30/// use rig_http::wasm_compat::WasmCompatSend;
31///
32/// fn shared<T: WasmCompatSend>(_: T) {}
33/// shared(Rc::new(0u8));
34/// ```
35#[diagnostic::on_unimplemented(
36    message = "`{Self}` is not `Send`, and Rig needs `Send` natively",
37    label = "not `Send`",
38    note = "handlers and transports run inside a driver's task: hold shared state behind an `Arc` (never an `Rc`), or use a `!Send` value only on browser wasm, where this marker is a no-op"
39)]
40pub trait WasmCompatSend: Send {}
41#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
42/// `Send` on native targets, a no-op marker on browser wasm.
43pub trait WasmCompatSend {}
44
45#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
46impl<T> WasmCompatSend for T where T: Send {}
47#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
48impl<T> WasmCompatSend for T {}
49
50#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
51/// Streaming response bound that includes `Send` on native targets.
52pub trait WasmCompatSendStream:
53    Stream<Item = Result<Bytes, crate::http_client::Error>> + Send
54{
55    type InnerItem: Send;
56}
57
58#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
59/// Streaming response bound without `Send` on browser wasm.
60pub trait WasmCompatSendStream: Stream<Item = Result<Bytes, crate::http_client::Error>> {
61    type InnerItem;
62}
63
64#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
65impl<T> WasmCompatSendStream for T
66where
67    T: Stream<Item = Result<Bytes, crate::http_client::Error>> + Send,
68{
69    type InnerItem = Result<Bytes, crate::http_client::Error>;
70}
71
72#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
73impl<T> WasmCompatSendStream for T
74where
75    T: Stream<Item = Result<Bytes, crate::http_client::Error>>,
76{
77    type InnerItem = Result<Bytes, crate::http_client::Error>;
78}
79
80#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
81/// `Sync` on native targets, a no-op marker on browser wasm.
82///
83/// ```compile_fail
84/// use std::cell::Cell;
85/// use rig_http::wasm_compat::WasmCompatSync;
86///
87/// fn shared<T: WasmCompatSync>(_: T) {}
88/// shared(Cell::new(0u8));
89/// ```
90#[diagnostic::on_unimplemented(
91    message = "`{Self}` is not `Sync`, and Rig needs `Sync` natively",
92    label = "not `Sync`",
93    note = "handlers and transports are shared between a driver and its in-flight tasks: use `Mutex`/atomics instead of `Cell`/`RefCell`, or use a `!Sync` value only on browser wasm, where this marker is a no-op"
94)]
95pub trait WasmCompatSync: Sync {}
96#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
97/// `Sync` on native targets, a no-op marker on browser wasm.
98pub trait WasmCompatSync {}
99
100#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
101impl<T> WasmCompatSync for T where T: Sync {}
102#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
103impl<T> WasmCompatSync for T {}
104
105#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
106/// Boxed future with `Send` on the same targets as [`WasmCompatSend`].
107pub type WasmBoxedFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
108
109#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
110/// Boxed future type without `Send`, on browser wasm.
111pub type WasmBoxedFuture<'a, T> = Pin<Box<dyn Future<Output = T> + 'a>>;
112
113#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
114/// Boxed stream with `Send` on the same targets as [`WasmCompatSend`].
115pub type WasmBoxedStream<'a, T> = Pin<Box<dyn Stream<Item = T> + Send + 'a>>;
116
117#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
118/// Boxed stream type without `Send`, on browser wasm.
119pub type WasmBoxedStream<'a, T> = Pin<Box<dyn Stream<Item = T> + 'a>>;
120
121/// Error returned by [`timeout`] when the future does not complete in time.
122#[derive(Debug, Clone, Copy, PartialEq, Eq)]
123pub struct Elapsed;
124
125impl std::fmt::Display for Elapsed {
126    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
127        f.write_str("future timed out")
128    }
129}
130
131impl std::error::Error for Elapsed {}
132
133/// Await `future` until `duration` elapses, then drop it and return [`Elapsed`].
134/// The future is polled first, even for zero duration; cancellation runs only
135/// its drop cleanup. Uses a native timer thread or browser `setTimeout`, without
136/// requiring an executor timer or a feature flag.
137///
138/// # Panics
139/// May panic if the timer cannot represent the deadline for `duration`.
140pub async fn timeout<F>(duration: std::time::Duration, future: F) -> Result<F::Output, Elapsed>
141where
142    F: Future,
143{
144    use futures::future::{Either, select};
145
146    let delay = futures_timer::Delay::new(duration);
147    futures::pin_mut!(future);
148    futures::pin_mut!(delay);
149    match select(future, delay).await {
150        Either::Left((output, _)) => Ok(output),
151        Either::Right(((), _)) => Err(Elapsed),
152    }
153}
154
155/// Sleep for `duration` using the native or browser timer backend of [`timeout`].
156///
157/// # Panics
158/// May panic if the timer cannot represent the deadline for `duration`.
159pub async fn sleep(duration: std::time::Duration) {
160    futures_timer::Delay::new(duration).await;
161}
162
163#[cfg(test)]
164mod tests;