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 },
}
}