Luigit
repositories / smith

smith

There are many coding harnesses - but this one is fast

owned by admin

smith-harness/src/http.rs

Raw
//! Host-owned HTTP engine over the `smith::http` vocabulary.
//!
//! One executor serves provider traffic today and plugin net-scope
//! effects later; every call is bounded, cancellable, and recordable.

use smith::error::{ProviderFault, Result, SmithError};
use smith::http::{
    HttpBodyEnd, HttpBodyItem, HttpChunk, HttpExchange, HttpExecutor, HttpExecutorRef, HttpHead,
    HttpResponse,
};
use smith::tool::CancelHandle;
use std::io::Read;
use std::sync::Arc;
use std::sync::mpsc;
use std::time::Duration;
use ureq::unversioned::resolver::DefaultResolver;
use ureq::unversioned::transport::{
    Buffers, ConnectionDetails, Connector, DefaultConnector, NextTimeout, Transport,
};

/// Largest accepted response body, in bytes.
pub const MAX_BODY_BYTES: u64 = 64 * 1024 * 1024;

/// Largest accepted single request body, in bytes.
pub const MAX_REQUEST_BYTES: usize = 16 * 1024 * 1024;

/// Timeout for the response head.
pub const HEAD_TIMEOUT: Duration = Duration::from_secs(120);

/// Deadline for one chunk gap.
pub const CHUNK_TIMEOUT: Duration = Duration::from_secs(300);

/// The HTTP engine: `ureq`, blocking, pure Rust, no async runtime
/// (`SMH-SPEC-SPEC0001`, network transport).
#[derive(Clone)]
pub struct UreqHttpExecutor {
    agent: ureq::Agent,
}

impl UreqHttpExecutor {
    /// A default executor with [`HEAD_TIMEOUT`] and [`CHUNK_TIMEOUT`].
    #[must_use]
    pub fn new() -> Self {
        Self::with_timeouts(HEAD_TIMEOUT, CHUNK_TIMEOUT)
    }

    /// An executor bounding the connect and the response head wait each by
    /// `head`, and every single socket wait by `gap`.
    #[must_use]
    pub fn with_timeouts(head: Duration, gap: Duration) -> Self {
        let config = ureq::Agent::config_builder()
            .http_status_as_error(false)
            .timeout_connect(Some(head))
            .timeout_recv_response(Some(head))
            .build();
        let connector = DefaultConnector::default().chain(GapTimeout(gap));
        Self {
            agent: ureq::Agent::with_parts(config, connector, DefaultResolver::default()),
        }
    }
}

impl Default for UreqHttpExecutor {
    fn default() -> Self {
        Self::new()
    }
}

impl HttpExecutor for UreqHttpExecutor {
    fn execute(&self, exchange: HttpExchange, cancel: CancelHandle) -> Result<HttpResponse> {
        admit(&exchange, &cancel)?;
        let mut prepared = match exchange.method {
            smith::http::HttpMethod::Post => self.agent.post(&exchange.url),
        };
        for header in &exchange.headers {
            prepared = prepared.header(&header.name, &header.value);
        }
        let response = prepared
            .send(&exchange.body[..])
            .map_err(|e| transient(e.to_string()))?;
        let head = HttpHead {
            status: response.status().as_u16(),
        };
        Ok(stream_body(
            head,
            response.into_body().into_reader(),
            cancel,
        ))
    }
}

/// Connector stage capping every socket wait at a gap duration, so a
/// stalled body fails per gap instead of per whole-body budget.
#[derive(Debug)]
struct GapTimeout(Duration);

impl Connector<Box<dyn Transport>> for GapTimeout {
    type Out = GapTransport;

    fn connect(
        &self,
        _details: &ConnectionDetails,
        chained: Option<Box<dyn Transport>>,
    ) -> std::result::Result<Option<Self::Out>, ureq::Error> {
        Ok(chained.map(|inner| GapTransport { inner, gap: self.0 }))
    }
}

#[derive(Debug)]
struct GapTransport {
    inner: Box<dyn Transport>,
    gap: Duration,
}

impl GapTransport {
    fn cap(&self, timeout: NextTimeout) -> NextTimeout {
        NextTimeout {
            after: timeout.after.min(self.gap.into()),
            reason: timeout.reason,
        }
    }
}

impl Transport for GapTransport {
    fn buffers(&mut self) -> &mut dyn Buffers {
        self.inner.buffers()
    }

    fn transmit_output(
        &mut self,
        amount: usize,
        timeout: NextTimeout,
    ) -> std::result::Result<(), ureq::Error> {
        let timeout = self.cap(timeout);
        self.inner.transmit_output(amount, timeout)
    }

    fn await_input(&mut self, timeout: NextTimeout) -> std::result::Result<bool, ureq::Error> {
        let timeout = self.cap(timeout);
        self.inner.await_input(timeout)
    }

    fn is_open(&mut self) -> bool {
        self.inner.is_open()
    }

    fn is_tls(&self) -> bool {
        self.inner.is_tls()
    }
}

/// The shared default executor.
#[must_use]
pub fn default_executor() -> HttpExecutorRef {
    Arc::new(UreqHttpExecutor::new())
}

fn admit(exchange: &HttpExchange, cancel: &CancelHandle) -> Result<()> {
    if exchange.body.len() > MAX_REQUEST_BYTES {
        return Err(transient(
            "http effect: request body exceeds the size bound".to_string(),
        ));
    }
    if cancel.is_cancelled() {
        return Err(SmithError::Cancelled);
    }
    Ok(())
}

/// Body chunks in flight before the reader thread blocks on the consumer.
const BODY_CHANNEL_CHUNKS: usize = 16;

/// Pump `body` into ordered items on a reader thread.
///
/// The body ends with [`HttpBodyEnd::Complete`] at end of body,
/// [`HttpBodyEnd::Failed`] on a read error or gap timeout,
/// [`HttpBodyEnd::Truncated`] past [`MAX_BODY_BYTES`], or
/// [`HttpBodyEnd::Cancelled`] once `cancel` fires; no chunk read after
/// cancellation is delivered.
fn stream_body(
    head: HttpHead,
    mut body: impl Read + Send + 'static,
    cancel: CancelHandle,
) -> HttpResponse {
    let (sender, receiver) = mpsc::sync_channel(BODY_CHANNEL_CHUNKS);
    std::thread::spawn(move || {
        let mut buffer = [0u8; 16 * 1024];
        let mut total: u64 = 0;
        let end = loop {
            let read = match body.read(&mut buffer) {
                Ok(0) => break HttpBodyEnd::Complete,
                Err(error) => {
                    break HttpBodyEnd::Failed {
                        message: error.to_string(),
                    };
                }
                Ok(read) => read,
            };
            if cancel.is_cancelled() {
                break HttpBodyEnd::Cancelled;
            }
            total = total.saturating_add(u64::try_from(read).unwrap_or(u64::MAX));
            if total > MAX_BODY_BYTES {
                break HttpBodyEnd::Truncated {
                    limit: MAX_BODY_BYTES,
                };
            }
            let chunk = HttpChunk {
                bytes: buffer[..read].to_vec(),
            };
            if sender.send(HttpBodyItem::Chunk(chunk)).is_err() {
                return;
            }
        };
        let _ = sender.send(HttpBodyItem::End(end));
    });
    HttpResponse {
        head,
        body: receiver,
    }
}

const fn transient(message: String) -> SmithError {
    SmithError::Provider {
        fault: ProviderFault::Transient { message },
    }
}