rig-http 0.44.0

Rig's transport contracts: the HTTP client and websocket traits a transport implements, and the wasm-compatible bounds they are written with.
Documentation
//! Target-dependent thread bounds, boxed futures, and executor-independent timers.
//!
//! ```
//! use rig_http::wasm_compat::WasmBoxedFuture;
//! let future: WasmBoxedFuture<'_, u32> = Box::pin(async { 42 });
//! ```

use bytes::Bytes;
use std::pin::Pin;

use futures::Stream;

// Browser markers assume single-threaded execution; atomics would invalidate
// that assumption. Relaxed bounds do not make non-Send handlers thread-safe.
#[cfg(all(
    target_arch = "wasm32",
    target_os = "unknown",
    target_feature = "atomics"
))]
compile_error!(
    "rig-http does not support threaded wasm (`+atomics`): its wasm-compat markers assume a \
     single-threaded target"
);

#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
/// `Send` on native targets, a no-op marker on browser wasm.
///
/// ```compile_fail
/// use std::rc::Rc;
/// use rig_http::wasm_compat::WasmCompatSend;
///
/// fn shared<T: WasmCompatSend>(_: T) {}
/// shared(Rc::new(0u8));
/// ```
#[diagnostic::on_unimplemented(
    message = "`{Self}` is not `Send`, and Rig needs `Send` natively",
    label = "not `Send`",
    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"
)]
pub trait WasmCompatSend: Send {}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
/// `Send` on native targets, a no-op marker on browser wasm.
pub trait WasmCompatSend {}

#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
impl<T> WasmCompatSend for T where T: Send {}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
impl<T> WasmCompatSend for T {}

#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
/// Streaming response bound that includes `Send` on native targets.
pub trait WasmCompatSendStream:
    Stream<Item = Result<Bytes, crate::http_client::Error>> + Send
{
    type InnerItem: Send;
}

#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
/// Streaming response bound without `Send` on browser wasm.
pub trait WasmCompatSendStream: Stream<Item = Result<Bytes, crate::http_client::Error>> {
    type InnerItem;
}

#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
impl<T> WasmCompatSendStream for T
where
    T: Stream<Item = Result<Bytes, crate::http_client::Error>> + Send,
{
    type InnerItem = Result<Bytes, crate::http_client::Error>;
}

#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
impl<T> WasmCompatSendStream for T
where
    T: Stream<Item = Result<Bytes, crate::http_client::Error>>,
{
    type InnerItem = Result<Bytes, crate::http_client::Error>;
}

#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
/// `Sync` on native targets, a no-op marker on browser wasm.
///
/// ```compile_fail
/// use std::cell::Cell;
/// use rig_http::wasm_compat::WasmCompatSync;
///
/// fn shared<T: WasmCompatSync>(_: T) {}
/// shared(Cell::new(0u8));
/// ```
#[diagnostic::on_unimplemented(
    message = "`{Self}` is not `Sync`, and Rig needs `Sync` natively",
    label = "not `Sync`",
    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"
)]
pub trait WasmCompatSync: Sync {}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
/// `Sync` on native targets, a no-op marker on browser wasm.
pub trait WasmCompatSync {}

#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
impl<T> WasmCompatSync for T where T: Sync {}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
impl<T> WasmCompatSync for T {}

#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
/// Boxed future with `Send` on the same targets as [`WasmCompatSend`].
pub type WasmBoxedFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;

#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
/// Boxed future type without `Send`, on browser wasm.
pub type WasmBoxedFuture<'a, T> = Pin<Box<dyn Future<Output = T> + 'a>>;

#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
/// Boxed stream with `Send` on the same targets as [`WasmCompatSend`].
pub type WasmBoxedStream<'a, T> = Pin<Box<dyn Stream<Item = T> + Send + 'a>>;

#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
/// Boxed stream type without `Send`, on browser wasm.
pub type WasmBoxedStream<'a, T> = Pin<Box<dyn Stream<Item = T> + 'a>>;

/// Error returned by [`timeout`] when the future does not complete in time.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Elapsed;

impl std::fmt::Display for Elapsed {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        f.write_str("future timed out")
    }
}

impl std::error::Error for Elapsed {}

/// Await `future` until `duration` elapses, then drop it and return [`Elapsed`].
/// The future is polled first, even for zero duration; cancellation runs only
/// its drop cleanup. Uses a native timer thread or browser `setTimeout`, without
/// requiring an executor timer or a feature flag.
///
/// # Panics
/// May panic if the timer cannot represent the deadline for `duration`.
pub async fn timeout<F>(duration: std::time::Duration, future: F) -> Result<F::Output, Elapsed>
where
    F: Future,
{
    use futures::future::{Either, select};

    let delay = futures_timer::Delay::new(duration);
    futures::pin_mut!(future);
    futures::pin_mut!(delay);
    match select(future, delay).await {
        Either::Left((output, _)) => Ok(output),
        Either::Right(((), _)) => Err(Elapsed),
    }
}

/// Sleep for `duration` using the native or browser timer backend of [`timeout`].
///
/// # Panics
/// May panic if the timer cannot represent the deadline for `duration`.
pub async fn sleep(duration: std::time::Duration) {
    futures_timer::Delay::new(duration).await;
}

#[cfg(test)]
mod tests;