use nmbrs_runtime::adapter::{
AdapterError, DriverAdapter, ExecutionError, JsonBody, OpDispenser, OpResult, ResultBody,
TextBody,
};
use nmbrs_workload::model::ParsedOp;
pub struct HttpConfig {
pub base_url: Option<String>,
pub timeout_ms: u64,
pub connect_timeout_ms: Option<u64>,
pub follow_redirects: bool,
}
impl Default for HttpConfig {
fn default() -> Self {
Self {
base_url: None,
timeout_ms: 30_000,
connect_timeout_ms: None,
follow_redirects: true,
}
}
}
impl HttpConfig {
pub fn from_params(params: &std::collections::HashMap<String, String>) -> Self {
Self {
base_url: params.get("base_url").or(params.get("host")).cloned(),
timeout_ms: params
.get("timeout")
.and_then(|s| s.parse().ok())
.unwrap_or(30_000),
connect_timeout_ms: None,
follow_redirects: true,
}
}
}
pub struct HttpAdapter {
client: reqwest::Client,
base_url: Option<String>,
config: HttpConfig,
}
fn build_http_client(
config: &HttpConfig,
connect_timeout_override_ms: Option<u64>,
) -> reqwest::Client {
let mut builder = reqwest::Client::builder()
.timeout(std::time::Duration::from_millis(config.timeout_ms))
.redirect(if config.follow_redirects {
reqwest::redirect::Policy::limited(10)
} else {
reqwest::redirect::Policy::none()
});
if let Some(ct) = connect_timeout_override_ms.or(config.connect_timeout_ms) {
builder = builder.connect_timeout(std::time::Duration::from_millis(ct));
}
builder.build().expect("failed to build HTTP client")
}
impl Default for HttpAdapter {
fn default() -> Self {
Self::new()
}
}
impl HttpAdapter {
pub fn new() -> Self {
Self::with_config(HttpConfig::default())
}
pub fn with_config(config: HttpConfig) -> Self {
let client = build_http_client(&config, None);
let base_url = config.base_url.clone();
Self {
client,
base_url,
config,
}
}
}
fn format_error_chain(e: &reqwest::Error) -> String {
use std::error::Error;
let mut msg = e.to_string();
let mut src: Option<&(dyn Error + 'static)> = e.source();
while let Some(s) = src {
let layer = s.to_string();
if !layer.is_empty() && !msg.contains(&layer) {
msg.push_str(": ");
msg.push_str(&layer);
}
src = s.source();
}
msg
}
fn is_transient_failure(e: &reqwest::Error) -> bool {
use std::error::Error;
if e.is_timeout() || e.is_connect() {
return true;
}
let mut src: Option<&(dyn Error + 'static)> = e.source();
while let Some(s) = src {
if let Some(io) = s.downcast_ref::<std::io::Error>() {
use std::io::ErrorKind::*;
if matches!(
io.kind(),
ConnectionRefused
| ConnectionReset
| ConnectionAborted
| TimedOut
| NotConnected
| BrokenPipe
) {
return true;
}
}
let low = s.to_string().to_ascii_lowercase();
if low.contains("timed out")
|| low.contains("connection refused")
|| low.contains("connection reset")
|| low.contains("dns error")
|| low.contains("unreachable")
{
return true;
}
src = s.source();
}
false
}
#[derive(Debug, Clone)]
struct OkStatusSpec(Vec<(u16, u16)>);
impl OkStatusSpec {
fn parse(spec: &str) -> Result<Self, String> {
let mut ranges = Vec::new();
for piece in spec.split(',') {
let piece = piece.trim();
if piece.is_empty() {
continue;
}
let (lo, hi) = match piece.split_once('-') {
Some((a, b)) => (a.trim(), b.trim()),
None => (piece, piece),
};
let lo: u16 = lo.parse().map_err(|_| {
format!("ok_status '{spec}': '{piece}' is not a status code or range")
})?;
let hi: u16 = hi.parse().map_err(|_| {
format!("ok_status '{spec}': '{piece}' is not a status code or range")
})?;
if lo > hi {
return Err(format!("ok_status '{spec}': range '{piece}' is inverted"));
}
ranges.push((lo, hi));
}
if ranges.is_empty() {
return Err(format!("ok_status '{spec}': no status codes"));
}
Ok(Self(ranges))
}
fn accepts(&self, status: u16) -> bool {
self.0.iter().any(|&(lo, hi)| (lo..=hi).contains(&status))
}
}
fn classify_reqwest_error(e: &reqwest::Error) -> String {
if e.is_timeout() {
"Timeout".into()
} else if e.is_connect() {
"ConnectionRefused".into()
} else if is_transient_failure(e) {
if format_error_chain(e)
.to_ascii_lowercase()
.contains("timed out")
{
"Timeout".into()
} else {
"ConnectionRefused".into()
}
} else if e.is_request() {
"RequestError".into()
} else {
"HttpError".into()
}
}
impl DriverAdapter for HttpAdapter {
fn name(&self) -> &str {
"http"
}
fn known_op_fields(&self) -> Option<&'static [&'static str]> {
Some(&[
"method",
"content_type",
"uri",
"url",
"body",
"headers",
"request_timeout_ms",
"on_timeout",
"connect_timeout",
"expect_body",
"ok_status",
])
}
fn map_op<'a>(
&'a self,
template: &'a ParsedOp,
parent: std::sync::Arc<dyn nmbrs_runtime::adapter::Kernel>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<Box<dyn OpDispenser>, String>> + Send + 'a>,
> {
Box::pin(async move {
let method = template
.op
.get("method")
.and_then(|v: &serde_json::Value| v.as_str())
.map(|s: &str| s.to_uppercase())
.unwrap_or_else(|| "GET".into());
let content_type = template
.op
.get("content_type")
.and_then(|v: &serde_json::Value| v.as_str())
.unwrap_or("application/json")
.to_string();
let uri_template = template
.op
.get("uri")
.or_else(|| template.op.get("url"))
.and_then(|v| v.as_str())
.map(String::from);
let body_template = template
.op
.get("body")
.and_then(|v| v.as_str())
.map(String::from);
let headers_template = template
.op
.get("headers")
.and_then(|v| v.as_str())
.map(String::from);
let per_op_timeout_ms = template.op.get("request_timeout_ms").and_then(|v| {
v.as_u64()
.or_else(|| v.as_str().and_then(|s| s.parse::<u64>().ok()))
});
let expect_body = template
.op
.get("expect_body")
.and_then(|v: &serde_json::Value| v.as_bool())
.unwrap_or(true);
let on_timeout_accept = template
.op
.get("on_timeout")
.and_then(|v| v.as_str())
.map(|s| s.eq_ignore_ascii_case("accept"))
.unwrap_or(false);
let connect_timeout_ms = template
.op
.get("connect_timeout")
.and_then(|v| v.as_str())
.and_then(|s| nmbrs_runtime::timeval::parse_time_ms(s).ok());
let client = match connect_timeout_ms {
Some(_) => build_http_client(&self.config, connect_timeout_ms),
None => self.client.clone(),
};
let ok_status = match template.op.get("ok_status").and_then(|v| v.as_str()) {
Some(spec) => Some(
OkStatusSpec::parse(spec)
.map_err(|e| format!("op '{}': {e}", template.name))?,
),
None => None,
};
Ok(Box::new(HttpDispenser {
client,
base_url: self.base_url.clone(),
method,
content_type,
canonical_kernel: parent,
uri_template,
body_template,
headers_template,
per_op_timeout_ms,
on_timeout_accept,
expect_body,
ok_status,
}) as Box<dyn OpDispenser>)
})
}
}
struct HttpDispenser {
client: reqwest::Client,
base_url: Option<String>,
method: String,
content_type: String,
canonical_kernel: std::sync::Arc<dyn nmbrs_runtime::adapter::Kernel>,
uri_template: Option<String>,
body_template: Option<String>,
headers_template: Option<String>,
per_op_timeout_ms: Option<u64>,
on_timeout_accept: bool,
expect_body: bool,
ok_status: Option<OkStatusSpec>,
}
impl OpDispenser for HttpDispenser {
fn canonical_kernel(&self) -> Option<&std::sync::Arc<dyn nmbrs_runtime::adapter::Kernel>> {
Some(&self.canonical_kernel)
}
fn execute<'a>(
&'a self,
_cycle: u64,
ctx: &'a nmbrs_runtime::adapter::ExecCtx<'a>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
> {
let wires = ctx.wires;
Box::pin(async move {
let uri_template = self.uri_template.as_deref().ok_or_else(|| {
ExecutionError::Op(AdapterError {
error_name: "missing_field".into(),
message: "HTTP op requires a 'uri' or 'url' field".into(),
retryable: false,
})
})?;
let uri =
nmbrs_runtime::wires::substitute_via_wires(uri_template, wires).map_err(|e| {
ExecutionError::Op(AdapterError {
error_name: "BindError".into(),
message: format!("uri: {e}"),
retryable: false,
})
})?;
let full_url = if let Some(ref base) = self.base_url {
if uri.starts_with("http://") || uri.starts_with("https://") {
uri.clone()
} else {
format!("{}{}", base.trim_end_matches('/'), uri)
}
} else {
uri.clone()
};
let body = match &self.body_template {
Some(t) => Some(
nmbrs_runtime::wires::substitute_via_wires(t, wires).map_err(|e| {
ExecutionError::Op(AdapterError {
error_name: "BindError".into(),
message: format!("body: {e}"),
retryable: false,
})
})?,
),
None => None,
};
let extra_headers: Vec<(String, String)> = match &self.headers_template {
Some(t) => {
let rendered =
nmbrs_runtime::wires::substitute_via_wires(t, wires).map_err(|e| {
ExecutionError::Op(AdapterError {
error_name: "BindError".into(),
message: format!("headers: {e}"),
retryable: false,
})
})?;
rendered
.lines()
.filter_map(|line| {
let mut parts = line.splitn(2, ':');
let name = parts.next()?.trim().to_string();
let value = parts.next()?.trim().to_string();
Some((name, value))
})
.collect()
}
None => Vec::new(),
};
let mut builder = match self.method.as_str() {
"GET" => self.client.get(&full_url),
"POST" => self.client.post(&full_url),
"PUT" => self.client.put(&full_url),
"DELETE" => self.client.delete(&full_url),
"PATCH" => self.client.patch(&full_url),
"HEAD" => self.client.head(&full_url),
other => {
return Err(ExecutionError::Op(AdapterError {
error_name: "InvalidMethod".into(),
message: format!("unsupported HTTP method: {other}"),
retryable: false,
}));
}
};
builder = builder.header("Content-Type", &self.content_type);
for (name, value) in &extra_headers {
builder = builder.header(name.as_str(), value.as_str());
}
if let Some(ms) = self.per_op_timeout_ms {
builder = builder.timeout(std::time::Duration::from_millis(ms));
}
if let Some(body_str) = body {
builder = builder.body(body_str);
}
let request_start = std::time::Instant::now();
let response = match builder.send().await {
Ok(r) => r,
Err(e) => {
if e.is_timeout() && self.on_timeout_accept {
let elapsed_ms = request_start.elapsed().as_millis();
let configured_ms = self
.per_op_timeout_ms
.map(|n| n.to_string())
.unwrap_or_else(|| "client-default".to_string());
nmbrs_runtime::observer::log(
if self.expect_body {
nmbrs_runtime::observer::LogLevel::Warn
} else {
nmbrs_runtime::observer::LogLevel::Debug
},
&format!(
"http: `on_timeout: accept` swallowed a \
request timeout after {elapsed_ms}ms \
(configured per_op_timeout_ms={configured_ms}) \
→ returning Ok(body=None). \
URL={full_url}. \
Value predicates in a downstream `verify:` \
go vacuous (nothing to read); `is: not_null` \
or `min_rows:` still fail, which is how to \
demand a body here."
),
);
return Ok(OpResult {
body: None,
skipped: false,
});
}
let retryable = is_transient_failure(&e);
let scope = if e.is_connect() {
ExecutionError::Adapter
} else {
ExecutionError::Op
};
return Err(scope(AdapterError {
error_name: classify_reqwest_error(&e),
message: format_error_chain(&e),
retryable,
}));
}
};
let status = response.status().as_u16() as i32;
let success = match &self.ok_status {
Some(spec) => spec.accepts(response.status().as_u16()),
None => response.status().is_success(),
};
let content_type_says_json = response
.headers()
.get(reqwest::header::CONTENT_TYPE)
.and_then(|v| v.to_str().ok())
.map(|ct| ct.contains("json"))
.unwrap_or(false);
let body_text = response.text().await.map_err(|e| {
ExecutionError::Op(AdapterError {
error_name: "BodyReadError".into(),
message: format!("failed to read response body: {e}"),
retryable: false,
})
})?;
if success {
let looks_like_json = body_text.trim_start().starts_with(['{', '[']);
let parsed_json = if content_type_says_json || looks_like_json {
serde_json::from_str::<serde_json::Value>(&body_text).ok()
} else {
None
};
let body: Box<dyn ResultBody> = match parsed_json {
Some(v) => Box::new(JsonBody(v)),
None => Box::new(TextBody(body_text)),
};
Ok(OpResult {
body: Some(body),
skipped: false,
})
} else {
Err(ExecutionError::Op(AdapterError {
error_name: format!("HttpStatus{}", status),
message: format!("HTTP {} {}: {}", status, full_url, &body_text),
retryable: (500..600).contains(&status),
}))
}
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use nmbrs_workload::model::ParsedOp;
#[test]
fn default_config() {
let config = HttpConfig::default();
assert_eq!(config.timeout_ms, 30_000);
assert!(config.follow_redirects);
assert!(config.base_url.is_none());
}
#[test]
fn adapter_creates() {
let _adapter = HttpAdapter::new();
}
async fn spawn_stalling_listener() -> u16 {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind 127.0.0.1:0");
let port = listener.local_addr().expect("local_addr").port();
tokio::spawn(async move {
while let Ok((sock, _)) = listener.accept().await {
tokio::spawn(async move {
let _hold = sock;
std::future::pending::<()>().await;
});
}
});
port
}
fn test_kernel() -> std::sync::Arc<dyn nmbrs_runtime::adapter::Kernel> {
std::sync::Arc::new(
polydat::dsl::compile::compile_polydat_interpreter("input cycle: u64\n").unwrap(),
)
}
fn http_op(
method: &str,
uri: &str,
request_timeout_ms: Option<&str>,
on_timeout: Option<&str>,
) -> ParsedOp {
let mut template = ParsedOp::simple("test", "");
template.op.remove("stmt");
template
.op
.insert("method".into(), serde_json::Value::String(method.into()));
template
.op
.insert("uri".into(), serde_json::Value::String(uri.into()));
if let Some(ms) = request_timeout_ms {
template.op.insert(
"request_timeout_ms".into(),
serde_json::Value::String(ms.into()),
);
}
if let Some(v) = on_timeout {
template
.op
.insert("on_timeout".into(), serde_json::Value::String(v.into()));
}
template
}
#[tokio::test]
async fn on_timeout_accept_swallows_request_timeout() {
let port = spawn_stalling_listener().await;
let adapter = HttpAdapter::new();
let template = http_op(
"GET",
&format!("http://127.0.0.1:{port}/"),
Some("100"), Some("accept"), );
let dispenser = adapter
.map_op(&template, test_kernel())
.await
.expect("map_op");
let mut k =
polydat::dsl::compile::compile_polydat_interpreter("input cycle: u64\n").unwrap();
let cw = nmbrs_runtime::wires::CycleWires::new(&mut k);
let pulls = nmbrs_runtime::fixture::ResolvedPulls::empty();
let empty = nmbrs_runtime::adapter::ResolvedFields::new(Vec::new(), Vec::new());
let ctx = nmbrs_runtime::adapter::ExecCtx::with_wires(&empty, &pulls, &cw);
let result = dispenser
.execute(0, &ctx)
.await
.expect("on_timeout: accept should map Timeout → Ok(empty)");
assert!(
result.body.is_none(),
"expected empty-body OpResult after accepted timeout"
);
assert!(
!result.skipped,
"accepted-timeout is a real (not skipped) op result"
);
}
#[tokio::test]
async fn timeout_without_accept_still_errors() {
let port = spawn_stalling_listener().await;
let adapter = HttpAdapter::new();
let template = http_op(
"GET",
&format!("http://127.0.0.1:{port}/"),
Some("100"),
None,
);
let dispenser = adapter
.map_op(&template, test_kernel())
.await
.expect("map_op");
let mut k =
polydat::dsl::compile::compile_polydat_interpreter("input cycle: u64\n").unwrap();
let cw = nmbrs_runtime::wires::CycleWires::new(&mut k);
let pulls = nmbrs_runtime::fixture::ResolvedPulls::empty();
let empty = nmbrs_runtime::adapter::ResolvedFields::new(Vec::new(), Vec::new());
let ctx = nmbrs_runtime::adapter::ExecCtx::with_wires(&empty, &pulls, &cw);
let err = dispenser
.execute(0, &ctx)
.await
.expect_err("default behaviour: client-side timeout → op error");
match err {
ExecutionError::Op(ad) => assert_eq!(
ad.error_name, "Timeout",
"expected error_name='Timeout', got: {ad:?}"
),
other => panic!("expected ExecutionError::Op(Timeout), got {other:?}"),
}
}
#[tokio::test]
async fn successful_json_response_populates_body() {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind 127.0.0.1:0");
let port = listener.local_addr().expect("local_addr").port();
tokio::spawn(async move {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
loop {
let Ok((mut sock, _)) = listener.accept().await else {
break;
};
tokio::spawn(async move {
let mut buf = vec![0u8; 4096];
let _ = sock.read(&mut buf).await;
let body = r#"{"status":200,"value":null,"request":{"type":"exec"}}"#;
let response = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
body.len(),
body,
);
let _ = sock.write_all(response.as_bytes()).await;
});
}
});
let adapter = HttpAdapter::new();
let template = http_op(
"POST",
&format!("http://127.0.0.1:{port}/jolokia/"),
None, None, );
let dispenser = adapter
.map_op(&template, test_kernel())
.await
.expect("map_op");
let mut k =
polydat::dsl::compile::compile_polydat_interpreter("input cycle: u64\n").unwrap();
let cw = nmbrs_runtime::wires::CycleWires::new(&mut k);
let pulls = nmbrs_runtime::fixture::ResolvedPulls::empty();
let empty = nmbrs_runtime::adapter::ResolvedFields::new(Vec::new(), Vec::new());
let ctx = nmbrs_runtime::adapter::ExecCtx::with_wires(&empty, &pulls, &cw);
let result = dispenser
.execute(0, &ctx)
.await
.expect("successful HTTP request should return Ok");
let body = result.body.as_ref().expect(
"successful response with body must populate result.body — \
no body indicates an adapter regression (the only legit \
body=None path is timeout-accept, which this test doesn't \
exercise)",
);
let json = body.to_json();
assert_eq!(
json.get("status").and_then(|v| v.as_u64()),
Some(200),
"body should preserve the server's `status` field; got: {json}"
);
}
}
inventory::submit! {
nmbrs_runtime::adapter::AdapterRegistration {
names: || &["http"],
known_params: || &["base_url", "host", "timeout"],
display_preference: |_params| nmbrs_runtime::adapter::DisplayPreference::Auto,
supported_controls: || &[],
create: |params| Box::pin(async move {
Ok(std::sync::Arc::new(HttpAdapter::with_config(HttpConfig::from_params(¶ms)))
as std::sync::Arc<dyn nmbrs_runtime::adapter::DriverAdapter>)
}),
}
}
inventory::submit! {
nmbrs_runtime::adapter::SharedDriverRegistration {
adapter: "http",
driver: nmbrs_runtime::adapter::DEFAULT_DRIVER_NAME,
share_capability: nmbrs_runtime::resource_pool::ShareCapability::Shared,
resource_key: |params| {
let cfg = HttpConfig::from_params(params);
Ok(nmbrs_runtime::resource_pool::ResourceKey::new("http")
.with("base_url", cfg.base_url.unwrap_or_default())
.with("timeout_ms", cfg.timeout_ms.to_string()))
},
}
}