Skip to main content

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(&params)))
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}