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
//! Type-erased HTTP transport with optional middleware.
//! Without middleware, requests pass through unchanged after body conversion to
//! bytes; lazy response bodies convert through `U::from`.
//!
//! ```
//! use rig_http::http_client::{DynHttpClient, HttpClientExt};
//!
//! fn erase(client: impl HttpClientExt + 'static) -> DynHttpClient {
//!     DynHttpClient::new(client)
//! }
//! ```

use std::any::Any;
use std::fmt;
use std::sync::Arc;

use bytes::Bytes;
use http::{Request, Response};

use super::middleware::HttpMiddleware;
use super::{HttpClientExt, LazyBody, MultipartForm, Result, StreamingResponse};
use crate::wasm_compat::{WasmBoxedFuture, WasmCompatSend, WasmCompatSync};

/// Object-safe mirror of [`HttpClientExt`] with the generics fixed to
/// [`Bytes`]. Private: the only way to reach it is through
/// [`DynHttpClient`], which re-exposes the generic surface.
pub(crate) trait ErasedHttpClient: WasmCompatSend + WasmCompatSync {
    fn send_bytes(
        &self,
        req: Request<Bytes>,
    ) -> WasmBoxedFuture<'static, Result<Response<LazyBody<Bytes>>>>;

    fn send_multipart_bytes(
        &self,
        req: Request<MultipartForm>,
    ) -> WasmBoxedFuture<'static, Result<Response<LazyBody<Bytes>>>>;

    /// Borrows `self`: [`HttpClientExt::send_streaming`] does not promise a
    /// `'static` future, so neither can its erasure.
    fn send_streaming_bytes(
        &self,
        req: Request<Bytes>,
    ) -> WasmBoxedFuture<'_, Result<StreamingResponse>>;
}

impl<H> ErasedHttpClient for H
where
    H: HttpClientExt + 'static,
{
    fn send_bytes(
        &self,
        req: Request<Bytes>,
    ) -> WasmBoxedFuture<'static, Result<Response<LazyBody<Bytes>>>> {
        Box::pin(self.send::<Bytes, Bytes>(req))
    }

    fn send_multipart_bytes(
        &self,
        req: Request<MultipartForm>,
    ) -> WasmBoxedFuture<'static, Result<Response<LazyBody<Bytes>>>> {
        Box::pin(self.send_multipart::<Bytes>(req))
    }

    fn send_streaming_bytes(
        &self,
        req: Request<Bytes>,
    ) -> WasmBoxedFuture<'_, Result<StreamingResponse>> {
        Box::pin(self.send_streaming::<Bytes>(req))
    }
}

/// A type-erased, cheaply cloneable HTTP transport.
///
/// Clones share the transport and middleware instances but copy the middleware
/// list. Erasing an already erased transport clones it without another layer.
/// `Debug` prints only the type name to avoid exposing transport credentials.
/// The transport is not serializable.
///
/// ```compile_fail
/// fn assert_serialize<T: serde::Serialize>() {}
/// assert_serialize::<rig_http::http_client::DynHttpClient>();
/// ```
#[derive(Clone)]
pub struct DynHttpClient {
    inner: Arc<dyn ErasedHttpClient>,
    /// Transport-boundary middleware, applied in attachment order around
    /// every request this handle sends. Cloned handles share the same stack.
    middleware: Vec<Arc<dyn HttpMiddleware>>,
}

impl DynHttpClient {
    /// Erase `http`. If `http` is already a `DynHttpClient`, this is a clone
    /// (its attached middleware included).
    pub fn new<H>(http: H) -> Self
    where
        H: HttpClientExt + 'static,
    {
        if let Some(already) = (&http as &dyn Any).downcast_ref::<DynHttpClient>() {
            return already.clone();
        }
        Self {
            inner: Arc::new(http),
            middleware: Vec::new(),
        }
    }

    /// Attach a transport-boundary [`HttpMiddleware`] to this handle.
    ///
    /// Hooks run in attachment order with the phases defined by [`HttpMiddleware`].
    /// Existing clones keep their previous stack. The underlying transport stays
    /// shared, so [`Self::ptr_eq`] is unaffected.
    pub fn with_middleware<M>(mut self, middleware: M) -> Self
    where
        M: HttpMiddleware + 'static,
    {
        self.middleware.push(Arc::new(middleware));
        self
    }

    /// Whether two handles share the same underlying transport.
    ///
    /// Compares the transport only: handles that differ in attached
    /// middleware but wrap the same transport still compare equal.
    pub fn ptr_eq(&self, other: &Self) -> bool {
        Arc::ptr_eq(&self.inner, &other.inner)
    }

    /// Run the request-side middleware phases: all header hooks in order,
    /// then all body hooks in order (each seeing the final headers).
    async fn apply_request_middleware(
        &self,
        parts: &mut http::request::Parts,
        body: Bytes,
    ) -> Result<Bytes> {
        for mw in &self.middleware {
            mw.before_request_headers(&parts.method, &parts.uri, &mut parts.headers)
                .await?;
        }
        let mut body = body;
        for mw in &self.middleware {
            body = mw
                .before_request_body(&parts.method, &parts.uri, &parts.headers, body)
                .await?;
        }
        Ok(body)
    }

    /// Run only the header hooks (the multipart path, which has no single
    /// serialized body to hand to the body hooks).
    async fn apply_header_middleware(&self, parts: &mut http::request::Parts) -> Result<()> {
        for mw in &self.middleware {
            mw.before_request_headers(&parts.method, &parts.uri, &mut parts.headers)
                .await?;
        }
        Ok(())
    }

    /// Run the response hooks in attachment order.
    async fn apply_response_middleware(
        &self,
        method: &http::Method,
        uri: &http::Uri,
        status: http::StatusCode,
        headers: &http::HeaderMap,
    ) -> Result<()> {
        for mw in &self.middleware {
            mw.after_response(method, uri, status, headers).await?;
        }
        Ok(())
    }
}

impl fmt::Debug for DynHttpClient {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        f.write_str("DynHttpClient")
    }
}

impl HttpClientExt for DynHttpClient {
    fn send<T, U>(
        &self,
        req: Request<T>,
    ) -> impl Future<Output = Result<Response<LazyBody<U>>>> + WasmCompatSend + 'static
    where
        T: Into<Bytes>,
        T: WasmCompatSend,
        U: From<Bytes>,
        U: WasmCompatSend + 'static,
    {
        let this = self.clone();
        let (mut parts, body) = req.map(Into::into).into_parts();
        async move {
            let body = this.apply_request_middleware(&mut parts, body).await?;
            let (method, uri) = (parts.method.clone(), parts.uri.clone());
            let response = this
                .inner
                .send_bytes(Request::from_parts(parts, body))
                .await?;
            this.apply_response_middleware(&method, &uri, response.status(), response.headers())
                .await?;
            Ok(convert_body::<U>(response))
        }
    }

    fn send_multipart<U>(
        &self,
        req: Request<MultipartForm>,
    ) -> impl Future<Output = Result<Response<LazyBody<U>>>> + WasmCompatSend + 'static
    where
        U: From<Bytes>,
        U: WasmCompatSend + 'static,
    {
        let this = self.clone();
        let (mut parts, body) = req.into_parts();
        async move {
            this.apply_header_middleware(&mut parts).await?;
            let (method, uri) = (parts.method.clone(), parts.uri.clone());
            let response = this
                .inner
                .send_multipart_bytes(Request::from_parts(parts, body))
                .await?;
            this.apply_response_middleware(&method, &uri, response.status(), response.headers())
                .await?;
            Ok(convert_body::<U>(response))
        }
    }

    fn send_streaming<T>(
        &self,
        req: Request<T>,
    ) -> impl Future<Output = Result<StreamingResponse>> + WasmCompatSend
    where
        T: Into<Bytes> + WasmCompatSend,
    {
        let (mut parts, body) = req.map(Into::into).into_parts();
        async move {
            let body = self.apply_request_middleware(&mut parts, body).await?;
            let (method, uri) = (parts.method.clone(), parts.uri.clone());
            let response = self
                .inner
                .send_streaming_bytes(Request::from_parts(parts, body))
                .await?;
            self.apply_response_middleware(&method, &uri, response.status(), response.headers())
                .await?;
            Ok(response)
        }
    }
}

fn convert_body<U>(response: Response<LazyBody<Bytes>>) -> Response<LazyBody<U>>
where
    U: From<Bytes> + WasmCompatSend + 'static,
{
    response.map(|body| -> LazyBody<U> { Box::pin(async move { body.await.map(U::from) }) })
}

#[cfg(test)]
mod tests;