nyquest-backend-winrt 0.4.0

Windows.Web.Http.HttpClient backend for nyquest
Documentation
use std::io;

use nyquest_interface::blocking::{BlockingBackend, BlockingClient, BlockingResponse, Request};
use nyquest_interface::client::ClientOptions;
use nyquest_interface::Result as NyquestResult;
use timer_ext::BlockingTimeoutExt;
use windows::Web::Http::HttpCompletionOption;

#[cfg(feature = "blocking-stream")]
mod stream_content;
mod timer_ext;

use crate::client::WinrtClient;
use crate::error::IntoNyquestResult;
use crate::ibuffer::IBufferExt;
use crate::request::create_body;
use crate::response::WinrtResponse;
use crate::response_size_limiter::ResponseSizeLimiter;
use crate::timer::Timer;

pub struct WinrtBlockingResponse {
    inner: WinrtResponse,
}

impl crate::WinrtBackend {
    pub fn create_blocking_client(&self, options: ClientOptions) -> io::Result<WinrtClient> {
        WinrtClient::create(options)
    }
}

impl WinrtClient {
    fn send_request(&self, req: Request) -> NyquestResult<WinrtBlockingResponse> {
        let req_msg = self.create_request(&req)?;
        if let Some(body) = req.body {
            #[cfg(feature = "blocking-stream")]
            let body = create_body(body, &mut stream_content::transform_stream)?;
            #[cfg(not(feature = "blocking-stream"))]
            let body = create_body(body, &mut |_| {
                unreachable!("blocking-stream feature is disabled")
            })?;
            self.append_content_headers(&body, &req.additional_headers)?;
            req_msg.SetContent(&body).into_nyquest_result()?;
        }
        let mut timer = Timer::new(self.request_timeout);
        let res = self
            .client
            .SendRequestWithOptionAsync(&req_msg, HttpCompletionOption::ResponseHeadersRead)
            .into_nyquest_result()?
            .timeout_by(&mut timer)?;
        let inner =
            WinrtResponse::new(res, self.max_response_buffer_size, timer).into_nyquest_result()?;
        Ok(WinrtBlockingResponse { inner })
    }
}

impl BlockingClient for WinrtClient {
    type Response = WinrtBlockingResponse;
    fn request(&self, req: Request) -> NyquestResult<Self::Response> {
        self.send_request(req)
    }
}

impl BlockingBackend for crate::WinrtBackend {
    type BlockingClient = WinrtClient;

    fn create_blocking_client(
        &self,
        options: ClientOptions,
    ) -> NyquestResult<Self::BlockingClient> {
        self.create_blocking_client(options).into_nyquest_result()
    }
}

impl BlockingResponse for WinrtBlockingResponse {
    fn status(&self) -> u16 {
        self.inner.status
    }

    fn get_header(&self, header: &str) -> nyquest_interface::Result<Vec<String>> {
        self.inner.get_header(header).into_nyquest_result()
    }

    fn content_length(&self) -> Option<u64> {
        self.inner.content_length
    }

    fn text(&mut self) -> NyquestResult<String> {
        let task = self
            .inner
            .content()
            .into_nyquest_result()?
            .ReadAsStringAsync()
            .into_nyquest_result()?;
        let size_limiter =
            ResponseSizeLimiter::hook_progress(self.inner.max_response_buffer_size, &task)?;
        let res = task
            .timeout_by(&mut self.inner.request_timer)
            .map(|r| r.to_string_lossy());
        let content = size_limiter.assert_size(res)?;
        Ok(content)
    }

    fn bytes(&mut self) -> NyquestResult<Vec<u8>> {
        let task = self
            .inner
            .content()
            .into_nyquest_result()?
            .ReadAsBufferAsync()
            .into_nyquest_result()?;
        let size_limiter =
            ResponseSizeLimiter::hook_progress(self.inner.max_response_buffer_size, &task)?;
        let res = task
            .timeout_by(&mut self.inner.request_timer)
            .and_then(|b| b.to_vec());
        let arr = size_limiter.assert_size(res)?;
        Ok(arr)
    }
}

#[cfg(feature = "blocking-stream")]
impl io::Read for WinrtBlockingResponse {
    fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
        let reader = self.inner.reader_mut()?;

        let mut size = reader.UnconsumedBufferLength()?;
        if size == 0 {
            let loaded = reader.LoadAsync(buf.len() as u32)?.get()?;
            if loaded == 0 {
                return Ok(0);
            }
            size = reader.UnconsumedBufferLength()?;
        }
        let size = buf.len().min(size as usize);
        let buf = &mut buf[..size];
        reader.ReadBytes(buf)?;
        Ok(size)
    }
}