Skip to main content

fizzy_sdk/http/
reqwest.rs

1use std::time::Duration;
2
3use async_trait::async_trait;
4use bytes::Bytes;
5use futures_util::stream;
6use reqwest::redirect::Policy;
7
8use crate::client::DEFAULT_TIMEOUT;
9use crate::error::Error;
10use crate::http::{Body, HttpClient, Request, Response};
11
12/// The [`HttpClient`] the SDK ships, over [`reqwest`] with rustls and HTTP/2. This is what a
13/// client gets when nothing else is supplied.
14///
15/// It never follows a redirect, whatever it was built from: the SDK does that itself.
16#[derive(Debug, Clone)]
17pub struct ReqwestClient {
18    http: reqwest::Client,
19}
20
21impl ReqwestClient {
22    /// A client that gives an answer `timeout` to arrive.
23    pub fn with_timeout(timeout: Duration) -> Result<ReqwestClient, Error> {
24        ReqwestClient::from_builder(reqwest::Client::builder().timeout(timeout))
25    }
26
27    /// A client built from settings of the caller's own — a proxy, a root certificate, a set
28    /// of default headers. Whatever redirect policy the builder carries is replaced with
29    /// none, since following one here would hide it from the SDK.
30    pub fn from_builder(builder: reqwest::ClientBuilder) -> Result<ReqwestClient, Error> {
31        let http = builder
32            .redirect(Policy::none())
33            .build()
34            .map_err(|error| Error::usage(format!("HTTP client: {error}")))?;
35        Ok(ReqwestClient { http })
36    }
37}
38
39impl Default for ReqwestClient {
40    #[allow(clippy::expect_used)] // reqwest builds a client from its own defaults
41    fn default() -> ReqwestClient {
42        ReqwestClient::with_timeout(DEFAULT_TIMEOUT)
43            .expect("reqwest builds a client from its defaults")
44    }
45}
46
47#[async_trait]
48impl HttpClient for ReqwestClient {
49    async fn send(&self, request: Request<Bytes>) -> Result<Response<Body>, Error> {
50        let request = reqwest::Request::try_from(request).map_err(Error::network)?;
51        let answered = self.http.execute(request).await.map_err(Error::network)?;
52
53        let status = answered.status();
54        let version = answered.version();
55        let headers = answered.headers().clone();
56        let content_length = answered.content_length();
57
58        let mut response = Response::new(Body::from_stream(chunks(answered), content_length));
59        *response.status_mut() = status;
60        *response.version_mut() = version;
61        *response.headers_mut() = headers;
62        Ok(response)
63    }
64}
65
66/// The body a chunk at a time: reqwest hands it out with `chunk()`, so the stream is that
67/// call repeated until it answers `None`.
68fn chunks(response: reqwest::Response) -> impl stream::Stream<Item = Result<Bytes, Error>> + Send {
69    stream::try_unfold(response, |mut response| async move {
70        match response.chunk().await {
71            Ok(Some(chunk)) => Ok(Some((chunk, response))),
72            Ok(None) => Ok(None),
73            Err(error) => Err(Error::network(error)),
74        }
75    })
76}