nmbrs_adapter_http/lib.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! HTTP adapter: executes operations as HTTP requests.
5//!
6//! Op template fields map to HTTP request components:
7//! - `method` — GET, POST, PUT, DELETE, PATCH, HEAD (default: GET)
8//! - `uri` or `url` — the request URL (required)
9//! - `body` — request body (for POST/PUT/PATCH)
10//! - `content_type` — Content-Type header (default: application/json)
11//! - `headers` — additional headers as "Name: Value" lines
12//! - `ok_status` — expected status codes (default: 200-299)
13//!
14//! Example workload:
15//! ```yaml
16//! bindings: |
17//! user_id := mod(hash(cycle), 1000000)
18//! ops:
19//! read:
20//! method: GET
21//! uri: "http://localhost:8080/api/users/{user_id}"
22//! write:
23//! method: POST
24//! uri: "http://localhost:8080/api/users"
25//! body: '{"id": {user_id}, "name": "user_{user_id}"}'
26//! content_type: application/json
27//! ```
28
29use nmbrs_runtime::adapter::{
30 AdapterError, DriverAdapter, ExecutionError, JsonBody, OpDispenser, OpResult, ResultBody,
31 TextBody,
32};
33use nmbrs_workload::model::ParsedOp;
34
35/// Configuration for the HTTP adapter.
36pub struct HttpConfig {
37 /// Base URL prefix prepended to relative URIs.
38 pub base_url: Option<String>,
39 /// Default timeout per request in milliseconds.
40 pub timeout_ms: u64,
41 /// Timeout for ESTABLISHING the TCP connection (the connect
42 /// phase), in ms. `None` leaves reqwest on the OS default (~tens
43 /// of seconds for an unreachable host). Distinct from `timeout_ms`
44 /// — the whole-request deadline — which cannot bound a connection
45 /// that never opens. Set via the `connect_timeout` param.
46 pub connect_timeout_ms: Option<u64>,
47 /// Whether to follow redirects.
48 pub follow_redirects: bool,
49}
50
51impl Default for HttpConfig {
52 fn default() -> Self {
53 Self {
54 base_url: None,
55 timeout_ms: 30_000,
56 connect_timeout_ms: None,
57 follow_redirects: true,
58 }
59 }
60}
61
62impl HttpConfig {
63 /// Construct a config from CLI/workload params.
64 pub fn from_params(params: &std::collections::HashMap<String, String>) -> Self {
65 Self {
66 base_url: params.get("base_url").or(params.get("host")).cloned(),
67 timeout_ms: params
68 .get("timeout")
69 .and_then(|s| s.parse().ok())
70 .unwrap_or(30_000),
71 // No client-wide `connect_timeout` from workload params: that
72 // name is already the CQL cluster-connect timeout at the workload
73 // root, and one value can't be both. HTTP takes `connect_timeout`
74 // as a PER-OP field instead (see `map_op`), so a Jolokia op can
75 // fail-fast without touching the CQL connect budget.
76 connect_timeout_ms: None,
77 follow_redirects: true,
78 }
79 }
80}
81
82/// The HTTP adapter: executes ops as HTTP requests.
83pub struct HttpAdapter {
84 client: reqwest::Client,
85 base_url: Option<String>,
86 /// Retained so `map_op` can rebuild a client with a PER-OP
87 /// `connect_timeout` — reqwest's `.connect_timeout()` is client-wide,
88 /// not settable per request.
89 config: HttpConfig,
90}
91
92/// Build a reqwest client from the adapter config, optionally overriding the
93/// connect-timeout for a single op. reqwest's `.connect_timeout()` is a
94/// client-wide setting, so a per-op value needs its own client — built once
95/// at map_op (per op template), never per request.
96fn build_http_client(
97 config: &HttpConfig,
98 connect_timeout_override_ms: Option<u64>,
99) -> reqwest::Client {
100 let mut builder = reqwest::Client::builder()
101 .timeout(std::time::Duration::from_millis(config.timeout_ms))
102 .redirect(if config.follow_redirects {
103 reqwest::redirect::Policy::limited(10)
104 } else {
105 reqwest::redirect::Policy::none()
106 });
107 if let Some(ct) = connect_timeout_override_ms.or(config.connect_timeout_ms) {
108 builder = builder.connect_timeout(std::time::Duration::from_millis(ct));
109 }
110 builder.build().expect("failed to build HTTP client")
111}
112
113impl Default for HttpAdapter {
114 fn default() -> Self {
115 Self::new()
116 }
117}
118
119impl HttpAdapter {
120 /// Create with default config.
121 pub fn new() -> Self {
122 Self::with_config(HttpConfig::default())
123 }
124
125 /// Create with explicit config.
126 pub fn with_config(config: HttpConfig) -> Self {
127 let client = build_http_client(&config, None);
128 let base_url = config.base_url.clone();
129 Self {
130 client,
131 base_url,
132 config,
133 }
134 }
135}
136
137/// Classify a reqwest error into an error name for the error router.
138/// Format a reqwest error together with its `.source()` chain. reqwest's
139/// own `Display` shows only the top layer (`error sending request for url
140/// (…)`); the ACTUAL cause — connection refused, connect timeout, dns
141/// failure — lives in the sources. Append each distinct layer so the
142/// operator sees WHAT failed, not just that something did.
143fn format_error_chain(e: &reqwest::Error) -> String {
144 use std::error::Error;
145 let mut msg = e.to_string();
146 let mut src: Option<&(dyn Error + 'static)> = e.source();
147 while let Some(s) = src {
148 let layer = s.to_string();
149 if !layer.is_empty() && !msg.contains(&layer) {
150 msg.push_str(": ");
151 msg.push_str(&layer);
152 }
153 src = s.source();
154 }
155 msg
156}
157
158/// True when the failure is a transient connection-phase problem —
159/// connect refused/reset/timeout, or a request timeout — worth retrying.
160/// reqwest's top-level `is_connect()` / `is_timeout()` report false when
161/// the real cause is buried under a generic "error sending request", so
162/// also walk the `.source()` chain for an io connect/timeout error.
163fn is_transient_failure(e: &reqwest::Error) -> bool {
164 use std::error::Error;
165 if e.is_timeout() || e.is_connect() {
166 return true;
167 }
168 let mut src: Option<&(dyn Error + 'static)> = e.source();
169 while let Some(s) = src {
170 if let Some(io) = s.downcast_ref::<std::io::Error>() {
171 use std::io::ErrorKind::*;
172 if matches!(
173 io.kind(),
174 ConnectionRefused
175 | ConnectionReset
176 | ConnectionAborted
177 | TimedOut
178 | NotConnected
179 | BrokenPipe
180 ) {
181 return true;
182 }
183 }
184 let low = s.to_string().to_ascii_lowercase();
185 if low.contains("timed out")
186 || low.contains("connection refused")
187 || low.contains("connection reset")
188 || low.contains("dns error")
189 || low.contains("unreachable")
190 {
191 return true;
192 }
193 src = s.source();
194 }
195 false
196}
197
198/// Parsed `ok_status` spec: comma-separated status codes and
199/// inclusive ranges (`"200-299,404"`). The adapter's doc has
200/// always promised this field; SRD-30's unknown-field guard is
201/// what surfaced that it was never wired.
202#[derive(Debug, Clone)]
203struct OkStatusSpec(Vec<(u16, u16)>);
204
205impl OkStatusSpec {
206 fn parse(spec: &str) -> Result<Self, String> {
207 let mut ranges = Vec::new();
208 for piece in spec.split(',') {
209 let piece = piece.trim();
210 if piece.is_empty() {
211 continue;
212 }
213 let (lo, hi) = match piece.split_once('-') {
214 Some((a, b)) => (a.trim(), b.trim()),
215 None => (piece, piece),
216 };
217 let lo: u16 = lo.parse().map_err(|_| {
218 format!("ok_status '{spec}': '{piece}' is not a status code or range")
219 })?;
220 let hi: u16 = hi.parse().map_err(|_| {
221 format!("ok_status '{spec}': '{piece}' is not a status code or range")
222 })?;
223 if lo > hi {
224 return Err(format!("ok_status '{spec}': range '{piece}' is inverted"));
225 }
226 ranges.push((lo, hi));
227 }
228 if ranges.is_empty() {
229 return Err(format!("ok_status '{spec}': no status codes"));
230 }
231 Ok(Self(ranges))
232 }
233
234 fn accepts(&self, status: u16) -> bool {
235 self.0.iter().any(|&(lo, hi)| (lo..=hi).contains(&status))
236 }
237}
238
239fn classify_reqwest_error(e: &reqwest::Error) -> String {
240 if e.is_timeout() {
241 "Timeout".into()
242 } else if e.is_connect() {
243 "ConnectionRefused".into()
244 } else if is_transient_failure(e) {
245 // Connect/timeout cause buried under a generic request error —
246 // name it by the underlying reason so `errors:` policies and the
247 // operator both see a connection problem, not bare "RequestError".
248 if format_error_chain(e)
249 .to_ascii_lowercase()
250 .contains("timed out")
251 {
252 "Timeout".into()
253 } else {
254 "ConnectionRefused".into()
255 }
256 } else if e.is_request() {
257 "RequestError".into()
258 } else {
259 "HttpError".into()
260 }
261}
262
263impl DriverAdapter for HttpAdapter {
264 fn name(&self) -> &str {
265 "http"
266 }
267
268 /// HTTP adapter reads a closed vocabulary of op fields:
269 /// request-shape (`method`, `uri` / `url`), body framing
270 /// (`content_type`, `body`), and header overrides
271 /// (`headers`). Declaring the list opts this adapter into
272 /// SRD 30's unknown-field guard — typos like `bdoy:` or
273 /// misplaced core directives surface at init time rather
274 /// than silently becoming ResolvedFields the adapter never
275 /// looks at.
276 fn known_op_fields(&self) -> Option<&'static [&'static str]> {
277 // `request_timeout_ms` (not `timeout_ms`) to avoid
278 // colliding with the polling wrapper's `timeout_ms`,
279 // which is the loop-level deadline. The HTTP adapter's
280 // value is a single-request budget.
281 //
282 // `on_timeout` is the SRD-74-style modifier that turns
283 // an HTTP-client-side timeout into a non-error empty
284 // result. Pairs with a short `request_timeout_ms` to
285 // express "fire and yield" — the server keeps doing
286 // its work whether or not the client is still
287 // listening (the canonical use case is Cassandra's
288 // synchronous `forceKeyspaceCompaction`, where the
289 // poll layer is the actual waiter / observer).
290 Some(&[
291 "method",
292 "content_type",
293 "uri",
294 "url",
295 "body",
296 "headers",
297 "request_timeout_ms",
298 "on_timeout",
299 "connect_timeout",
300 "expect_body",
301 "ok_status",
302 ])
303 }
304
305 fn map_op<'a>(
306 &'a self,
307 template: &'a ParsedOp,
308 parent: std::sync::Arc<dyn nmbrs_runtime::adapter::Kernel>,
309 ) -> std::pin::Pin<
310 Box<dyn std::future::Future<Output = Result<Box<dyn OpDispenser>, String>> + Send + 'a>,
311 > {
312 Box::pin(async move {
313 // Extract static method from template (default GET)
314 let method = template
315 .op
316 .get("method")
317 .and_then(|v: &serde_json::Value| v.as_str())
318 .map(|s: &str| s.to_uppercase())
319 .unwrap_or_else(|| "GET".into());
320
321 // Extract content type (default application/json)
322 let content_type = template
323 .op
324 .get("content_type")
325 .and_then(|v: &serde_json::Value| v.as_str())
326 .unwrap_or("application/json")
327 .to_string();
328
329 // SRD-68 Push 5: snapshot the per-cycle field templates at
330 // map_op. Each is rendered through `substitute_via_wires`
331 // at execute — the generic Polydat API resolves bind points by
332 // name, no synthesis-layer ResolvedFields involvement.
333 // `url` is an alias for `uri`; honour whichever appears.
334 let uri_template = template
335 .op
336 .get("uri")
337 .or_else(|| template.op.get("url"))
338 .and_then(|v| v.as_str())
339 .map(String::from);
340 let body_template = template
341 .op
342 .get("body")
343 .and_then(|v| v.as_str())
344 .map(String::from);
345 let headers_template = template
346 .op
347 .get("headers")
348 .and_then(|v| v.as_str())
349 .map(String::from);
350 // Per-op timeout override. Cassandra's
351 // `forceKeyspaceCompaction` JMX op is synchronous (blocks
352 // for the entire compaction); the default 30s client
353 // timeout is far too short for any real table size. This
354 // field lets workloads opt into a longer per-request
355 // budget without raising the adapter-wide default.
356 // Named `request_timeout_ms` (not `timeout_ms`) so it
357 // doesn't collide with the polling wrapper's loop-level
358 // `timeout_ms`.
359 let per_op_timeout_ms = template.op.get("request_timeout_ms").and_then(|v| {
360 v.as_u64()
361 .or_else(|| v.as_str().and_then(|s| s.parse::<u64>().ok()))
362 });
363
364 // `on_timeout: accept` is the fire-and-yield modifier
365 // (SRD-74-style). When a request_timeout_ms is set and
366 // the HTTP client trips it, the adapter would normally
367 // return a `Timeout` ExecutionError that fails the
368 // op. With `accept`, that specific outcome converts to
369 // a successful `OpResult` with no body — the server
370 // is presumed to still be doing the work; the polling
371 // layer is what observes its completion.
372 //
373 // Errors that are NOT `is_timeout()` (connection
374 // refused, body read errors, non-2xx responses) still
375 // surface normally. The modifier only translates
376 // client-side request-timeout firings.
377 // `expect_body: false` DECLARES that a body-less success is a normal
378 // outcome for this op — the fire-and-forget trigger whose work the
379 // poll layer observes, or any 204. The accept-timeout diagnostic
380 // below exists to explain a SURPRISE; an op that has said it expects
381 // no body is not surprised, and on a 256-phase sweep that warning is
382 // pure noise repeated once per tier. Declaring it drops the line to
383 // Debug rather than removing it, so `--log-level=debug` can still
384 // recover the timing.
385 let expect_body = template
386 .op
387 .get("expect_body")
388 .and_then(|v: &serde_json::Value| v.as_bool())
389 .unwrap_or(true);
390 let on_timeout_accept = template
391 .op
392 .get("on_timeout")
393 .and_then(|v| v.as_str())
394 .map(|s| s.eq_ignore_ascii_case("accept"))
395 .unwrap_or(false);
396
397 // Per-op connect timeout. reqwest's `.connect_timeout()` is
398 // client-wide, so an op that sets `connect_timeout` (a duration
399 // spec-string like `15s`, or a bare number = fractional seconds) gets
400 // its OWN client built with that value. Distinct from
401 // `request_timeout_ms` (the response deadline): this bounds the TCP
402 // CONNECT phase, so an unreachable endpoint fails fast and `retries:`
403 // kicks in instead of hanging on the OS default (~tens of seconds).
404 let connect_timeout_ms = template
405 .op
406 .get("connect_timeout")
407 .and_then(|v| v.as_str())
408 .and_then(|s| nmbrs_runtime::timeval::parse_time_ms(s).ok());
409 let client = match connect_timeout_ms {
410 Some(_) => build_http_client(&self.config, connect_timeout_ms),
411 None => self.client.clone(),
412 };
413
414 // `ok_status` — which response statuses count as success
415 // for THIS op (`"200-299,404"`). Default: reqwest's
416 // is_success (2xx). The canonical use is idempotent
417 // teardown, where 404 on an absent resource is the no-op
418 // outcome, not an error.
419 let ok_status = match template.op.get("ok_status").and_then(|v| v.as_str()) {
420 Some(spec) => Some(
421 OkStatusSpec::parse(spec)
422 .map_err(|e| format!("op '{}': {e}", template.name))?,
423 ),
424 None => None,
425 };
426
427 Ok(Box::new(HttpDispenser {
428 client,
429 base_url: self.base_url.clone(),
430 method,
431 content_type,
432 canonical_kernel: parent,
433 uri_template,
434 body_template,
435 headers_template,
436 per_op_timeout_ms,
437 on_timeout_accept,
438 expect_body,
439 ok_status,
440 }) as Box<dyn OpDispenser>)
441 })
442 }
443}
444
445/// Op dispenser for the HTTP adapter. Pre-analyzes method and content type
446/// at init time; resolves URI and body from wires per-cycle.
447struct HttpDispenser {
448 client: reqwest::Client,
449 base_url: Option<String>,
450 method: String,
451 content_type: String,
452 /// SRD-68 invariant I-3: dispenser-owned canonical Polydat Kernel.
453 canonical_kernel: std::sync::Arc<dyn nmbrs_runtime::adapter::Kernel>,
454 /// Cycle-time templates rendered through `substitute_via_wires`.
455 /// `uri` is mandatory; `body` and `headers` are optional.
456 uri_template: Option<String>,
457 body_template: Option<String>,
458 headers_template: Option<String>,
459 /// Optional per-op request timeout override. When set, the
460 /// builder applies `.timeout(...)` on the request — bypassing
461 /// the adapter's client-wide default. Use for long-running
462 /// JMX/REST calls (e.g. Jolokia synchronous
463 /// `forceKeyspaceCompaction`) that legitimately take many
464 /// minutes.
465 per_op_timeout_ms: Option<u64>,
466 /// When `true`, a request-timeout firing translates to a
467 /// successful empty-body `OpResult` instead of a
468 /// `Timeout` op error. Pair with a short
469 /// `per_op_timeout_ms` to express "fire and yield":
470 /// submit the request, give the server a brief window to
471 /// respond, but don't fail the phase if it doesn't —
472 /// the server keeps working server-side regardless of
473 /// whether the client is still listening.
474 on_timeout_accept: bool,
475 /// False when the workload declared `expect_body: false`.
476 expect_body: bool,
477 /// Op-declared success statuses (`ok_status:`); `None` means
478 /// the 2xx default.
479 ok_status: Option<OkStatusSpec>,
480}
481
482impl OpDispenser for HttpDispenser {
483 fn canonical_kernel(&self) -> Option<&std::sync::Arc<dyn nmbrs_runtime::adapter::Kernel>> {
484 Some(&self.canonical_kernel)
485 }
486
487 fn execute<'a>(
488 &'a self,
489 _cycle: u64,
490 ctx: &'a nmbrs_runtime::adapter::ExecCtx<'a>,
491 ) -> std::pin::Pin<
492 Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
493 > {
494 let wires = ctx.wires;
495 Box::pin(async move {
496 let uri_template = self.uri_template.as_deref().ok_or_else(|| {
497 ExecutionError::Op(AdapterError {
498 error_name: "missing_field".into(),
499 message: "HTTP op requires a 'uri' or 'url' field".into(),
500 retryable: false,
501 })
502 })?;
503
504 // SRD-68 Push 5: render each per-cycle template via the
505 // generic wires API. Bind-point resolution failures are
506 // returned as op errors so the error router decides.
507 let uri =
508 nmbrs_runtime::wires::substitute_via_wires(uri_template, wires).map_err(|e| {
509 ExecutionError::Op(AdapterError {
510 error_name: "BindError".into(),
511 message: format!("uri: {e}"),
512 retryable: false,
513 })
514 })?;
515
516 let full_url = if let Some(ref base) = self.base_url {
517 if uri.starts_with("http://") || uri.starts_with("https://") {
518 uri.clone()
519 } else {
520 format!("{}{}", base.trim_end_matches('/'), uri)
521 }
522 } else {
523 uri.clone()
524 };
525
526 let body = match &self.body_template {
527 Some(t) => Some(
528 nmbrs_runtime::wires::substitute_via_wires(t, wires).map_err(|e| {
529 ExecutionError::Op(AdapterError {
530 error_name: "BindError".into(),
531 message: format!("body: {e}"),
532 retryable: false,
533 })
534 })?,
535 ),
536 None => None,
537 };
538
539 // Parse additional headers from the rendered headers
540 // field. Per-line `Name: Value` entries.
541 let extra_headers: Vec<(String, String)> = match &self.headers_template {
542 Some(t) => {
543 let rendered =
544 nmbrs_runtime::wires::substitute_via_wires(t, wires).map_err(|e| {
545 ExecutionError::Op(AdapterError {
546 error_name: "BindError".into(),
547 message: format!("headers: {e}"),
548 retryable: false,
549 })
550 })?;
551 rendered
552 .lines()
553 .filter_map(|line| {
554 let mut parts = line.splitn(2, ':');
555 let name = parts.next()?.trim().to_string();
556 let value = parts.next()?.trim().to_string();
557 Some((name, value))
558 })
559 .collect()
560 }
561 None => Vec::new(),
562 };
563
564 let mut builder = match self.method.as_str() {
565 "GET" => self.client.get(&full_url),
566 "POST" => self.client.post(&full_url),
567 "PUT" => self.client.put(&full_url),
568 "DELETE" => self.client.delete(&full_url),
569 "PATCH" => self.client.patch(&full_url),
570 "HEAD" => self.client.head(&full_url),
571 other => {
572 return Err(ExecutionError::Op(AdapterError {
573 error_name: "InvalidMethod".into(),
574 message: format!("unsupported HTTP method: {other}"),
575 retryable: false,
576 }));
577 }
578 };
579
580 builder = builder.header("Content-Type", &self.content_type);
581
582 for (name, value) in &extra_headers {
583 builder = builder.header(name.as_str(), value.as_str());
584 }
585
586 // Per-op timeout override. When unset, reqwest falls
587 // back to the adapter's client-wide default
588 // (`timeout=` in workload params, 30s otherwise).
589 if let Some(ms) = self.per_op_timeout_ms {
590 builder = builder.timeout(std::time::Duration::from_millis(ms));
591 }
592
593 if let Some(body_str) = body {
594 builder = builder.body(body_str);
595 }
596
597 let request_start = std::time::Instant::now();
598 let response = match builder.send().await {
599 Ok(r) => r,
600 Err(e) => {
601 // `on_timeout: accept` converts ONLY the
602 // client-side request-timeout firing into a
603 // benign empty-body success. Every other error
604 // path (connection refused, request build
605 // failure, body read failure) still surfaces
606 // — the modifier is narrowly scoped to the
607 // "fire and yield" pattern where the
608 // expectation is that the request reached the
609 // server but the server is doing long work
610 // synchronously, and the polling layer is the
611 // canonical waiter / observer.
612 if e.is_timeout() && self.on_timeout_accept {
613 // Diagnostic: this branch is the ONLY way
614 // the HTTP adapter returns `body: None`.
615 // Value predicates downstream go vacuous
616 // rather than failing, so without this log
617 // an accepted timeout would be entirely
618 // silent — and "the call timed out but the
619 // server is still working" is exactly the
620 // state an operator needs to see. Surfacing
621 // the accept here makes the chain obvious
622 // in session.log without changing the
623 // success-shape semantics.
624 let elapsed_ms = request_start.elapsed().as_millis();
625 let configured_ms = self
626 .per_op_timeout_ms
627 .map(|n| n.to_string())
628 .unwrap_or_else(|| "client-default".to_string());
629 nmbrs_runtime::observer::log(
630 // An op that declared `expect_body: false` has said a
631 // body-less success is its normal outcome, so this is
632 // not news - Debug, not Warn. Without that declaration
633 // it stays a warning: a silently swallowed timeout IS
634 // worth seeing.
635 if self.expect_body {
636 nmbrs_runtime::observer::LogLevel::Warn
637 } else {
638 nmbrs_runtime::observer::LogLevel::Debug
639 },
640 &format!(
641 "http: `on_timeout: accept` swallowed a \
642 request timeout after {elapsed_ms}ms \
643 (configured per_op_timeout_ms={configured_ms}) \
644 → returning Ok(body=None). \
645 URL={full_url}. \
646 Value predicates in a downstream `verify:` \
647 go vacuous (nothing to read); `is: not_null` \
648 or `min_rows:` still fail, which is how to \
649 demand a body here."
650 ),
651 );
652 return Ok(OpResult {
653 body: None,
654 skipped: false,
655 });
656 }
657 let retryable = is_transient_failure(&e);
658 let scope = if e.is_connect() {
659 ExecutionError::Adapter
660 } else {
661 ExecutionError::Op
662 };
663 return Err(scope(AdapterError {
664 error_name: classify_reqwest_error(&e),
665 message: format_error_chain(&e),
666 retryable,
667 }));
668 }
669 };
670
671 let status = response.status().as_u16() as i32;
672 let success = match &self.ok_status {
673 Some(spec) => spec.accepts(response.status().as_u16()),
674 None => response.status().is_success(),
675 };
676 // Capture content-type before consuming the response
677 // body so we can pick the right `ResultBody` shape.
678 // `application/json` (or any `…/json` subtype like
679 // `application/vnd.api+json`) parses into a `JsonBody`
680 // — verify-blocks can then address nested fields
681 // (`field: status, eq: "200"`) instead of substring
682 // matching on the raw text.
683 let content_type_says_json = response
684 .headers()
685 .get(reqwest::header::CONTENT_TYPE)
686 .and_then(|v| v.to_str().ok())
687 .map(|ct| ct.contains("json"))
688 .unwrap_or(false);
689 let body_text = response.text().await.map_err(|e| {
690 ExecutionError::Op(AdapterError {
691 error_name: "BodyReadError".into(),
692 message: format!("failed to read response body: {e}"),
693 retryable: false,
694 })
695 })?;
696
697 if success {
698 // Promote to `JsonBody` whenever the body parses
699 // as JSON — not just when the server bothered to
700 // set the right Content-Type. Jolokia 1.x and
701 // various JMX bridges return JSON with a
702 // `text/plain` (or missing) content type;
703 // requiring the header would make verify blocks
704 // unable to address nested fields (`field:
705 // status, eq: "200"` → `<not-json>` even though
706 // the body literally is JSON).
707 //
708 // Gate the parse attempt on a cheap prefix check
709 // (`{` / `[` after whitespace) so we don't
710 // serde_json::from_str scan arbitrary text
711 // bodies that happen to start with a digit or a
712 // quoted string. That keeps "parse a scalar like
713 // "42" into a JSON number" — a real risk for
714 // plain-text endpoints — from happening.
715 let looks_like_json = body_text.trim_start().starts_with(['{', '[']);
716 let parsed_json = if content_type_says_json || looks_like_json {
717 serde_json::from_str::<serde_json::Value>(&body_text).ok()
718 } else {
719 None
720 };
721 let body: Box<dyn ResultBody> = match parsed_json {
722 Some(v) => Box::new(JsonBody(v)),
723 None => Box::new(TextBody(body_text)),
724 };
725 Ok(OpResult {
726 body: Some(body),
727 skipped: false,
728 })
729 } else {
730 Err(ExecutionError::Op(AdapterError {
731 error_name: format!("HttpStatus{}", status),
732 message: format!("HTTP {} {}: {}", status, full_url, &body_text),
733 retryable: (500..600).contains(&status),
734 }))
735 }
736 })
737 }
738}
739
740#[cfg(test)]
741mod tests {
742 use super::*;
743 use nmbrs_workload::model::ParsedOp;
744
745 #[test]
746 fn default_config() {
747 let config = HttpConfig::default();
748 assert_eq!(config.timeout_ms, 30_000);
749 assert!(config.follow_redirects);
750 assert!(config.base_url.is_none());
751 }
752
753 #[test]
754 fn adapter_creates() {
755 let _adapter = HttpAdapter::new();
756 }
757
758 /// Spin up a TCP listener that accepts a single connection
759 /// and then sleeps forever — the canonical "server is busy
760 /// doing the long thing, won't answer" shape that
761 /// `on_timeout: accept` exists to handle. Returns the
762 /// listener's bound port; the caller addresses
763 /// `http://127.0.0.1:<port>/`.
764 async fn spawn_stalling_listener() -> u16 {
765 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
766 .await
767 .expect("bind 127.0.0.1:0");
768 let port = listener.local_addr().expect("local_addr").port();
769 tokio::spawn(async move {
770 // Accept connections in a loop so a test that
771 // retries doesn't deadlock on the second attempt.
772 // Each accepted socket is just held — never read,
773 // never written — until the test process tears it
774 // down.
775 while let Ok((sock, _)) = listener.accept().await {
776 // Stash the socket on the task heap so the OS
777 // doesn't drop the connection and let the client
778 // see EOF instead of a timeout.
779 tokio::spawn(async move {
780 let _hold = sock;
781 std::future::pending::<()>().await;
782 });
783 }
784 });
785 port
786 }
787
788 fn test_kernel() -> std::sync::Arc<dyn nmbrs_runtime::adapter::Kernel> {
789 std::sync::Arc::new(
790 polydat::dsl::compile::compile_polydat_interpreter("input cycle: u64\n").unwrap(),
791 )
792 }
793
794 /// Build an HTTP op-template programmatically. Bypasses
795 /// the full workload parser to keep the test focused on
796 /// the adapter's per-op behaviour.
797 fn http_op(
798 method: &str,
799 uri: &str,
800 request_timeout_ms: Option<&str>,
801 on_timeout: Option<&str>,
802 ) -> ParsedOp {
803 let mut template = ParsedOp::simple("test", "");
804 template.op.remove("stmt");
805 template
806 .op
807 .insert("method".into(), serde_json::Value::String(method.into()));
808 template
809 .op
810 .insert("uri".into(), serde_json::Value::String(uri.into()));
811 if let Some(ms) = request_timeout_ms {
812 template.op.insert(
813 "request_timeout_ms".into(),
814 serde_json::Value::String(ms.into()),
815 );
816 }
817 if let Some(v) = on_timeout {
818 template
819 .op
820 .insert("on_timeout".into(), serde_json::Value::String(v.into()));
821 }
822 template
823 }
824
825 /// With `on_timeout: accept`, a request-timeout firing
826 /// converts to a successful empty-body `OpResult`. The
827 /// canonical use case is Cassandra's synchronous
828 /// `forceKeyspaceCompaction` (server keeps working
829 /// regardless of whether the client is still listening);
830 /// here we stand up a TcpListener that accepts the
831 /// connection but never answers, which produces the same
832 /// client-side reqwest error.
833 #[tokio::test]
834 async fn on_timeout_accept_swallows_request_timeout() {
835 let port = spawn_stalling_listener().await;
836 let adapter = HttpAdapter::new();
837 let template = http_op(
838 "GET",
839 &format!("http://127.0.0.1:{port}/"),
840 Some("100"), // 100ms request timeout
841 Some("accept"), // swallow client-side timeout
842 );
843 let dispenser = adapter
844 .map_op(&template, test_kernel())
845 .await
846 .expect("map_op");
847
848 let mut k =
849 polydat::dsl::compile::compile_polydat_interpreter("input cycle: u64\n").unwrap();
850 let cw = nmbrs_runtime::wires::CycleWires::new(&mut k);
851 let pulls = nmbrs_runtime::fixture::ResolvedPulls::empty();
852 let empty = nmbrs_runtime::adapter::ResolvedFields::new(Vec::new(), Vec::new());
853 let ctx = nmbrs_runtime::adapter::ExecCtx::with_wires(&empty, &pulls, &cw);
854
855 let result = dispenser
856 .execute(0, &ctx)
857 .await
858 .expect("on_timeout: accept should map Timeout → Ok(empty)");
859 assert!(
860 result.body.is_none(),
861 "expected empty-body OpResult after accepted timeout"
862 );
863 assert!(
864 !result.skipped,
865 "accepted-timeout is a real (not skipped) op result"
866 );
867 }
868
869 /// Without `on_timeout: accept`, the same stalling-listener
870 /// scenario produces a `Timeout` op error. Pins the
871 /// negative case so the accept-branch can't accidentally
872 /// regress to swallowing every error category.
873 #[tokio::test]
874 async fn timeout_without_accept_still_errors() {
875 let port = spawn_stalling_listener().await;
876 let adapter = HttpAdapter::new();
877 let template = http_op(
878 "GET",
879 &format!("http://127.0.0.1:{port}/"),
880 Some("100"),
881 None,
882 );
883 let dispenser = adapter
884 .map_op(&template, test_kernel())
885 .await
886 .expect("map_op");
887
888 let mut k =
889 polydat::dsl::compile::compile_polydat_interpreter("input cycle: u64\n").unwrap();
890 let cw = nmbrs_runtime::wires::CycleWires::new(&mut k);
891 let pulls = nmbrs_runtime::fixture::ResolvedPulls::empty();
892 let empty = nmbrs_runtime::adapter::ResolvedFields::new(Vec::new(), Vec::new());
893 let ctx = nmbrs_runtime::adapter::ExecCtx::with_wires(&empty, &pulls, &cw);
894
895 let err = dispenser
896 .execute(0, &ctx)
897 .await
898 .expect_err("default behaviour: client-side timeout → op error");
899 match err {
900 ExecutionError::Op(ad) => assert_eq!(
901 ad.error_name, "Timeout",
902 "expected error_name='Timeout', got: {ad:?}"
903 ),
904 other => panic!("expected ExecutionError::Op(Timeout), got {other:?}"),
905 }
906 }
907
908 /// Happy-path regression test: when the server returns a
909 /// well-formed JSON body, the adapter's `OpResult.body` is
910 /// `Some(JsonBody(...))` — not `None`. This pins the
911 /// invariant that the HTTP adapter never returns
912 /// `body: None` for a successful request (the only None
913 /// path is the timeout-accept branch tested elsewhere). A
914 /// regression here would surface as
915 /// `<no body returned by op>` in validation diagnostics
916 /// even though the server replied normally.
917 #[tokio::test]
918 async fn successful_json_response_populates_body() {
919 // Spin up a one-shot HTTP server that replies with a
920 // Jolokia-shaped JSON body.
921 let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
922 .await
923 .expect("bind 127.0.0.1:0");
924 let port = listener.local_addr().expect("local_addr").port();
925 tokio::spawn(async move {
926 use tokio::io::{AsyncReadExt, AsyncWriteExt};
927 loop {
928 let Ok((mut sock, _)) = listener.accept().await else {
929 break;
930 };
931 tokio::spawn(async move {
932 let mut buf = vec![0u8; 4096];
933 let _ = sock.read(&mut buf).await;
934 let body = r#"{"status":200,"value":null,"request":{"type":"exec"}}"#;
935 let response = format!(
936 "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
937 body.len(),
938 body,
939 );
940 let _ = sock.write_all(response.as_bytes()).await;
941 });
942 }
943 });
944
945 let adapter = HttpAdapter::new();
946 let template = http_op(
947 "POST",
948 &format!("http://127.0.0.1:{port}/jolokia/"),
949 None, // no per-op timeout
950 None, // no on_timeout
951 );
952 let dispenser = adapter
953 .map_op(&template, test_kernel())
954 .await
955 .expect("map_op");
956
957 let mut k =
958 polydat::dsl::compile::compile_polydat_interpreter("input cycle: u64\n").unwrap();
959 let cw = nmbrs_runtime::wires::CycleWires::new(&mut k);
960 let pulls = nmbrs_runtime::fixture::ResolvedPulls::empty();
961 let empty = nmbrs_runtime::adapter::ResolvedFields::new(Vec::new(), Vec::new());
962 let ctx = nmbrs_runtime::adapter::ExecCtx::with_wires(&empty, &pulls, &cw);
963
964 let result = dispenser
965 .execute(0, &ctx)
966 .await
967 .expect("successful HTTP request should return Ok");
968 let body = result.body.as_ref().expect(
969 "successful response with body must populate result.body — \
970 no body indicates an adapter regression (the only legit \
971 body=None path is timeout-accept, which this test doesn't \
972 exercise)",
973 );
974 let json = body.to_json();
975 assert_eq!(
976 json.get("status").and_then(|v| v.as_u64()),
977 Some(200),
978 "body should preserve the server's `status` field; got: {json}"
979 );
980 }
981}
982
983// =========================================================================
984// Adapter Registration (inventory-based, link-time)
985// =========================================================================
986
987inventory::submit! {
988 nmbrs_runtime::adapter::AdapterRegistration {
989 names: || &["http"],
990 known_params: || &["base_url", "host", "timeout"],
991 display_preference: |_params| nmbrs_runtime::adapter::DisplayPreference::Auto,
992 supported_controls: || &[],
993 create: |params| Box::pin(async move {
994 Ok(std::sync::Arc::new(HttpAdapter::with_config(HttpConfig::from_params(¶ms)))
995 as std::sync::Arc<dyn nmbrs_runtime::adapter::DriverAdapter>)
996 }),
997 }
998}
999
1000// SRD-35 Push C: HTTP adapter declares itself
1001// pool-shareable. The reqwest `Client` is documented
1002// thread-safe and pools connections internally; sharing
1003// one `HttpAdapter` across all phases that target the
1004// same `(base_url, timeout)` combination eliminates the
1005// per-phase TLS handshake / connection-establish storm.
1006//
1007// `base_url` and `timeout` are instance-shaping (the same
1008// reqwest client serves every request that uses them);
1009// per-call URL paths and method overrides come in via the
1010// op-template layer and don't affect the resource key.
1011inventory::submit! {
1012 nmbrs_runtime::adapter::SharedDriverRegistration {
1013 adapter: "http",
1014 driver: nmbrs_runtime::adapter::DEFAULT_DRIVER_NAME,
1015 share_capability: nmbrs_runtime::resource_pool::ShareCapability::Shared,
1016 resource_key: |params| {
1017 let cfg = HttpConfig::from_params(params);
1018 Ok(nmbrs_runtime::resource_pool::ResourceKey::new("http")
1019 .with("base_url", cfg.base_url.unwrap_or_default())
1020 .with("timeout_ms", cfg.timeout_ms.to_string()))
1021 },
1022 }
1023}