//! 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 { 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> for GapTimeout { type Out = GapTransport; fn connect( &self, _details: &ConnectionDetails, chained: Option>, ) -> std::result::Result, ureq::Error> { Ok(chained.map(|inner| GapTransport { inner, gap: self.0 })) } } #[derive(Debug)] struct GapTransport { inner: Box, 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 { 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 }, } }