ntex 4.0.0-beta.17

Framework for composable network services
#[cfg(feature = "compress")]
use crate::http::{Payload, encoding::Decoder, header};
use crate::{Ctx, Service, SharedCfg, error::Error, http::body::MessageBody};

use super::{ClientConfig, ClientRawRequest, Connect, ServiceRequest, ServiceResponse};
use super::{connector::Connector, error::ClientError};

#[derive(Debug)]
pub struct Sender {
    connector: Connector,
}

impl Sender {
    pub(super) fn new(connector: Connector) -> Self {
        Self { connector }
    }
}

#[allow(unused_variables)]
impl Service<SharedCfg, ServiceRequest> for Sender {
    type Res = ServiceResponse;
    type Error = Error<ClientError>;

    crate::forward_ready!(SharedCfg, connector);
    crate::forward_shutdown!(SharedCfg, connector);

    async fn call(
        &self,
        req: ServiceRequest,
        ctx: Ctx<'_, Self, SharedCfg>,
    ) -> Result<Self::Res, Self::Error> {
        let ServiceRequest {
            head,
            addr,
            body,
            headers,
            timeout,
            response_decompress,
        } = req;

        let uri = head.uri.clone();
        let con = ctx.call(&self.connector, Connect { uri, addr }).await?;
        let config = ctx.st().get::<ClientConfig>();

        let timeout = timeout.unwrap_or_else(|| config.response_timeout());

        let req = ClientRawRequest {
            head,
            headers,
            size: body.size(),
        };

        #[allow(unused_mut)]
        let (mut head, payload) = con.send_request(req, body, timeout).await?;

        #[cfg(feature = "compress")]
        if response_decompress {
            let decoder = Decoder::from_headers(payload, &head.headers);
            if decoder.is_decoding() {
                // the headers describe the encoded payload
                head.headers.remove(&header::CONTENT_ENCODING);
                head.headers.remove(&header::CONTENT_LENGTH);
            }
            let payload = Payload::from_stream(decoder);
            return Ok(ServiceResponse {
                head,
                payload,
                config: config.clone(),
            });
        }

        Ok(ServiceResponse {
            head,
            payload,
            config,
        })
    }
}