Skip to main content

ironflow_core/operations/
http.rs

1//! Http operation - perform HTTP requests with timeout and header control.
2//!
3//! The [`Http`] builder sends an HTTP request via [`reqwest`], captures the
4//! response, and returns an [`HttpOutput`] on success. It implements
5//! [`IntoFuture`] so you can `await` it directly:
6//!
7//! ```no_run
8//! use ironflow_core::operations::http::Http;
9//!
10//! # async fn example() -> Result<(), ironflow_core::error::OperationError> {
11//! let output = Http::get("https://httpbin.org/get").await?;
12//! println!("status: {}", output.status());
13//! # Ok(())
14//! # }
15//! ```
16
17use std::collections::HashMap;
18use std::env::{self, VarError};
19use std::future::{Future, IntoFuture};
20use std::pin::Pin;
21use std::time::{Duration, Instant};
22
23use reqwest::redirect::Policy;
24use reqwest::{Client, Method};
25use serde::de::DeserializeOwned;
26use serde_json::Value;
27use std::sync::LazyLock;
28use tokio::time;
29use tracing::{debug, warn};
30use url::Url;
31
32use crate::retry::RetryPolicy;
33use crate::ssrf::{self, AllowedHosts, GuardedResolver};
34use crate::trace_context::WorkflowTraceContext;
35
36/// Default timeout for HTTP requests (30 seconds).
37const DEFAULT_HTTP_TIMEOUT: Duration = Duration::from_secs(30);
38
39use crate::error::OperationError;
40#[cfg(feature = "prometheus")]
41use crate::metric_names;
42use crate::utils::MAX_OUTPUT_SIZE;
43
44/// Environment variable listing, comma-separated, the hosts every [`Http`] request of
45/// the deployment may reach even when they are internal.
46const ALLOWED_HOSTS_ENV: &str = "IRONFLOW_HTTP_ALLOWED_HOSTS";
47
48static ENV_ALLOWED_HOSTS: LazyLock<AllowedHosts> =
49    LazyLock::new(|| match env::var(ALLOWED_HOSTS_ENV) {
50        Ok(list) => AllowedHosts::parse_list(&list),
51        Err(VarError::NotPresent) => AllowedHosts::default(),
52        Err(err) => {
53            warn!(error = %err, "{ALLOWED_HOSTS_ENV} ignored: no internal host is allowed");
54            AllowedHosts::default()
55        }
56    });
57
58/// Client for hosts that may be internal: no SSRF guard, environment proxies honored.
59static HTTP_CLIENT: LazyLock<Client> = LazyLock::new(|| {
60    Client::builder()
61        .redirect(Policy::none())
62        .build()
63        .expect("failed to build HTTP client")
64});
65
66/// Client for every other host. Proxies are ignored: a proxy resolves the target
67/// itself, out of reach of [`GuardedResolver`].
68static GUARDED_HTTP_CLIENT: LazyLock<Client> = LazyLock::new(|| {
69    Client::builder()
70        .redirect(Policy::none())
71        .no_proxy()
72        .dns_resolver(GuardedResolver::default())
73        .build()
74        .expect("failed to build HTTP client")
75});
76
77/// Builder for executing an HTTP request.
78///
79/// Supports method, URL, headers, body (JSON or text), and timeout.
80/// The response body is captured as a string, with optional typed
81/// JSON deserialization via [`HttpOutput::json`].
82///
83/// Unlike [`Shell`](crate::operations::shell::Shell), `Http` does **not**
84/// fail on non-2xx status codes - use [`HttpOutput::is_success`] to check.
85/// Only transport-level errors (DNS, timeout, connection refused) produce
86/// an [`OperationError::Http`].
87///
88/// # Examples
89///
90/// ```no_run
91/// use std::time::Duration;
92/// use ironflow_core::operations::http::Http;
93///
94/// # async fn example() -> Result<(), ironflow_core::error::OperationError> {
95/// let output = Http::post("https://httpbin.org/post")
96///     .header("Authorization", "Bearer token123")
97///     .json(serde_json::json!({"key": "value"}))
98///     .timeout(Duration::from_secs(30))
99///     .await?;
100///
101/// println!("status: {}, body: {}", output.status(), output.body());
102/// # Ok(())
103/// # }
104/// ```
105#[must_use = "an Http request does nothing until .run() or .await is called"]
106pub struct Http {
107    method: Method,
108    url: String,
109    headers: HashMap<String, String>,
110    body: Option<HttpBody>,
111    timeout: Option<Duration>,
112    max_response_size: usize,
113    dry_run: Option<bool>,
114    retry_policy: Option<RetryPolicy>,
115    allowed_hosts: AllowedHosts,
116}
117
118enum HttpBody {
119    Text(String),
120    Json(Value),
121}
122
123impl Http {
124    /// Create a request builder with an arbitrary HTTP method.
125    ///
126    /// # Panics
127    ///
128    /// Panics if `url` is empty.
129    pub fn new(method: Method, url: &str) -> Self {
130        let trimmed = url.trim();
131        assert!(!trimmed.is_empty(), "url must not be empty");
132        assert!(
133            trimmed.starts_with("http://") || trimmed.starts_with("https://"),
134            "url must use http:// or https:// scheme, got: {trimmed}"
135        );
136        Self {
137            method,
138            url: trimmed.to_string(),
139            headers: HashMap::new(),
140            body: None,
141            timeout: Some(DEFAULT_HTTP_TIMEOUT),
142            max_response_size: MAX_OUTPUT_SIZE,
143            dry_run: None,
144            retry_policy: None,
145            allowed_hosts: AllowedHosts::default(),
146        }
147    }
148
149    /// Create a GET request builder.
150    ///
151    /// # Examples
152    ///
153    /// ```no_run
154    /// use ironflow_core::operations::http::Http;
155    ///
156    /// # async fn example() -> Result<(), ironflow_core::error::OperationError> {
157    /// let output = Http::get("https://httpbin.org/get").await?;
158    /// # Ok(())
159    /// # }
160    /// ```
161    pub fn get(url: &str) -> Self {
162        Self::new(Method::GET, url)
163    }
164
165    /// Create a POST request builder.
166    pub fn post(url: &str) -> Self {
167        Self::new(Method::POST, url)
168    }
169
170    /// Create a PUT request builder.
171    pub fn put(url: &str) -> Self {
172        Self::new(Method::PUT, url)
173    }
174
175    /// Create a PATCH request builder.
176    pub fn patch(url: &str) -> Self {
177        Self::new(Method::PATCH, url)
178    }
179
180    /// Create a DELETE request builder.
181    pub fn delete(url: &str) -> Self {
182        Self::new(Method::DELETE, url)
183    }
184
185    /// Add a header to the request.
186    ///
187    /// Can be called multiple times to set several headers.
188    pub fn header(mut self, key: &str, value: &str) -> Self {
189        self.headers.insert(key.to_string(), value.to_string());
190        self
191    }
192
193    /// Set a JSON body.
194    ///
195    /// `Content-Type: application/json` is added automatically by reqwest.
196    /// Takes ownership of the [`Value`] to avoid cloning.
197    pub fn json(mut self, value: Value) -> Self {
198        self.body = Some(HttpBody::Json(value));
199        self
200    }
201
202    /// Set a plain text body.
203    pub fn text(mut self, body: &str) -> Self {
204        self.body = Some(HttpBody::Text(body.to_string()));
205        self
206    }
207
208    /// Override the timeout for the request.
209    ///
210    /// If the request does not complete within this duration, an
211    /// [`OperationError::Http`] is returned. Defaults to 30 seconds.
212    pub fn timeout(mut self, timeout: Duration) -> Self {
213        self.timeout = Some(timeout);
214        self
215    }
216
217    /// Allow this request to reach `host` even when it is, or resolves to, a private,
218    /// loopback, link-local or cloud metadata address.
219    ///
220    /// Without it, such a target fails with [`OperationError::Http`] before anything is
221    /// sent. Pass the host as it appears in the URL (`"billing.internal"`, `"10.0.0.5"`,
222    /// `"::1"`); the match ignores case and IPv6 brackets. Hosts allowed for the whole
223    /// deployment go in the `IRONFLOW_HTTP_ALLOWED_HOSTS` environment variable, a
224    /// comma-separated list read once at the first request.
225    ///
226    /// # Examples
227    ///
228    /// ```no_run
229    /// use ironflow_core::operations::http::Http;
230    ///
231    /// # async fn example() -> Result<(), ironflow_core::error::OperationError> {
232    /// let output = Http::get("http://billing.internal:8080/invoices")
233    ///     .allow_host("billing.internal")
234    ///     .await?;
235    /// # Ok(())
236    /// # }
237    /// ```
238    pub fn allow_host(mut self, host: &str) -> Self {
239        self.allowed_hosts.add(host);
240        self
241    }
242
243    /// Set the maximum allowed response body size in bytes.
244    ///
245    /// If the response body exceeds this limit, an [`OperationError::Http`] is
246    /// returned. Defaults to 10 MiB.
247    pub fn max_response_size(mut self, bytes: usize) -> Self {
248        self.max_response_size = bytes;
249        self
250    }
251
252    /// Retry the request up to `max_retries` times on transient failures.
253    ///
254    /// Uses default exponential backoff settings (200ms initial, 2x multiplier,
255    /// 30s cap). For custom backoff parameters, use [`retry_policy`](Http::retry_policy).
256    ///
257    /// Only transient errors are retried: transport errors (DNS, timeout,
258    /// connection refused) and responses with status 5xx or 429. Client errors
259    /// (4xx except 429) and SSRF blocks are never retried.
260    ///
261    /// # Panics
262    ///
263    /// Panics if `max_retries` is `0`.
264    ///
265    /// # Examples
266    ///
267    /// ```no_run
268    /// use ironflow_core::operations::http::Http;
269    ///
270    /// # async fn example() -> Result<(), ironflow_core::error::OperationError> {
271    /// let output = Http::get("https://api.example.com/data")
272    ///     .retry(3)
273    ///     .await?;
274    /// # Ok(())
275    /// # }
276    /// ```
277    pub fn retry(mut self, max_retries: u32) -> Self {
278        self.retry_policy = Some(RetryPolicy::new(max_retries));
279        self
280    }
281
282    /// Set a custom [`RetryPolicy`] for this request.
283    ///
284    /// Allows full control over backoff duration, multiplier, and max delay.
285    /// See [`RetryPolicy`] for details.
286    ///
287    /// # Examples
288    ///
289    /// ```no_run
290    /// use std::time::Duration;
291    /// use ironflow_core::operations::http::Http;
292    /// use ironflow_core::retry::RetryPolicy;
293    ///
294    /// # async fn example() -> Result<(), ironflow_core::error::OperationError> {
295    /// let output = Http::get("https://api.example.com/data")
296    ///     .retry_policy(
297    ///         RetryPolicy::new(5)
298    ///             .backoff(Duration::from_millis(500))
299    ///             .max_backoff(Duration::from_secs(60))
300    ///             .multiplier(3.0)
301    ///     )
302    ///     .await?;
303    /// # Ok(())
304    /// # }
305    /// ```
306    pub fn retry_policy(mut self, policy: RetryPolicy) -> Self {
307        self.retry_policy = Some(policy);
308        self
309    }
310
311    /// Attach a [`WorkflowTraceContext`] to this request.
312    ///
313    /// When set, the `traceparent` header is automatically injected into
314    /// the request using the context's [`to_traceparent`](WorkflowTraceContext::to_traceparent)
315    /// value. This enables distributed tracing correlation with downstream
316    /// services.
317    ///
318    /// # Examples
319    ///
320    /// ```no_run
321    /// use ironflow_core::operations::http::Http;
322    /// use ironflow_core::trace_context::WorkflowTraceContext;
323    ///
324    /// # async fn example() -> Result<(), ironflow_core::error::OperationError> {
325    /// let ctx = WorkflowTraceContext::new_root();
326    /// let output = Http::get("https://api.example.com/data")
327    ///     .trace_context(&ctx)
328    ///     .await?;
329    /// # Ok(())
330    /// # }
331    /// ```
332    pub fn trace_context(self, ctx: &WorkflowTraceContext) -> Self {
333        self.header("traceparent", &ctx.to_traceparent())
334    }
335
336    /// Enable or disable dry-run mode for this specific operation.
337    ///
338    /// When dry-run is active, the request is logged but not sent.
339    /// A synthetic [`HttpOutput`] is returned with status 200, empty body,
340    /// and 0ms duration.
341    ///
342    /// If not set, falls back to the global dry-run setting
343    /// (see [`set_dry_run`](crate::dry_run::set_dry_run)).
344    pub fn dry_run(mut self, enabled: bool) -> Self {
345        self.dry_run = Some(enabled);
346        self
347    }
348
349    /// Execute the HTTP request.
350    ///
351    /// If a [`retry_policy`](Http::retry_policy) is configured, transient
352    /// failures (transport errors, 5xx, 429) are retried with exponential
353    /// backoff. Non-retryable errors and successful responses are returned
354    /// immediately.
355    ///
356    /// # Errors
357    ///
358    /// Returns [`OperationError::Http`] if the request fails at the transport
359    /// layer (network error, DNS failure, timeout) or if the response body
360    /// cannot be read. Non-2xx status codes are **not** treated as errors.
361    ///
362    /// Also returns [`OperationError::Http`], before anything is sent and without
363    /// retries, if the URL cannot be parsed, or if its host is, or resolves to, a
364    /// private, loopback, link-local or cloud metadata address that
365    /// [`allow_host`](Http::allow_host) or `IRONFLOW_HTTP_ALLOWED_HOSTS` does not allow.
366    #[tracing::instrument(name = "http", skip_all, fields(method = %self.method, url = %self.url))]
367    pub async fn run(self) -> Result<HttpOutput, OperationError> {
368        if crate::dry_run::effective_dry_run(self.dry_run) {
369            debug!(method = %self.method, url = %self.url, "[dry-run] http request skipped");
370            return Ok(HttpOutput {
371                status: 200,
372                headers: HashMap::new(),
373                body: String::new(),
374                duration_ms: 0,
375            });
376        }
377
378        let url = Url::parse(&self.url).map_err(|e| OperationError::Http {
379            status: None,
380            message: format!("invalid URL {}: {e}", self.url),
381        })?;
382        let client = if self.allowed_hosts.contains_url_host(&url)
383            || ENV_ALLOWED_HOSTS.contains_url_host(&url)
384        {
385            &*HTTP_CLIENT
386        } else {
387            ssrf::check_url(&url)
388                .await
389                .map_err(|blocked| OperationError::Http {
390                    status: None,
391                    message: blocked.to_string(),
392                })?;
393            &*GUARDED_HTTP_CLIENT
394        };
395
396        let result = self.execute_once(client).await;
397
398        let policy = match &self.retry_policy {
399            Some(p) => p,
400            None => return result,
401        };
402
403        // If the first attempt succeeded with a non-retryable status, return it.
404        // If it failed with a non-retryable error, return it.
405        match &result {
406            Ok(output) if !crate::retry::is_retryable_status(output.status) => return result,
407            Err(err) if !crate::retry::is_retryable(err) => return result,
408            _ => {}
409        }
410
411        let mut last_result = result;
412
413        for attempt in 0..policy.max_retries {
414            let delay = policy.delay_for_attempt(attempt);
415            warn!(
416                attempt = attempt + 1,
417                max_retries = policy.max_retries,
418                delay_ms = delay.as_millis() as u64,
419                "retrying http request"
420            );
421            time::sleep(delay).await;
422
423            last_result = self.execute_once(client).await;
424
425            match &last_result {
426                Ok(output) if !crate::retry::is_retryable_status(output.status) => {
427                    return last_result;
428                }
429                Err(err) if !crate::retry::is_retryable(err) => return last_result,
430                _ => {}
431            }
432        }
433
434        last_result
435    }
436
437    /// Execute a single HTTP request attempt (no retry logic).
438    async fn execute_once(&self, client: &Client) -> Result<HttpOutput, OperationError> {
439        debug!(method = %self.method, url = %self.url, "executing http request");
440        let start = Instant::now();
441
442        #[cfg(feature = "prometheus")]
443        let method_label = self.method.to_string();
444
445        let mut builder = client.request(self.method.clone(), &self.url);
446
447        if let Some(timeout) = self.timeout {
448            builder = builder.timeout(timeout);
449        }
450
451        for (k, v) in &self.headers {
452            builder = builder.header(k.as_str(), v.as_str());
453        }
454
455        match &self.body {
456            Some(HttpBody::Json(v)) => {
457                builder = builder.json(v);
458            }
459            Some(HttpBody::Text(t)) => {
460                builder = builder.body(t.clone());
461            }
462            None => {}
463        }
464
465        let response = match builder.send().await {
466            Ok(resp) => resp,
467            Err(e) => {
468                #[cfg(feature = "prometheus")]
469                {
470                    metrics::counter!(metric_names::HTTP_TOTAL, "method" => method_label, "status" => metric_names::STATUS_ERROR).increment(1);
471                }
472                return Err(OperationError::Http {
473                    status: None,
474                    // A DNS answer that changed since `check_url` (rebinding) is refused
475                    // by the resolver, deep in the source chain.
476                    message: match ssrf::find_blocked(&e) {
477                        Some(blocked) => blocked.to_string(),
478                        None => format!("request failed: {e}"),
479                    },
480                });
481            }
482        };
483
484        let status = response.status().as_u16();
485        let headers: HashMap<String, String> = response
486            .headers()
487            .iter()
488            .map(|(k, v)| {
489                let val = match v.to_str() {
490                    Ok(s) => s.to_string(),
491                    Err(_) => {
492                        debug!(header = %k, "non-UTF-8 header value, replacing with empty string");
493                        String::new()
494                    }
495                };
496                (k.to_string(), val)
497            })
498            .collect();
499        let max_response_size = self.max_response_size;
500        let response_too_large = |size: usize, limit: usize| OperationError::Http {
501            status: Some(status),
502            message: format!(
503                "response body too large: {size} bytes exceeds limit of {limit} bytes"
504            ),
505        };
506
507        if let Some(cl) = response.content_length() {
508            let content_length = usize::try_from(cl).unwrap_or(usize::MAX);
509            if content_length > max_response_size {
510                return Err(response_too_large(content_length, max_response_size));
511            }
512        }
513
514        let mut body_bytes = Vec::new();
515        let mut response = response;
516        loop {
517            match response.chunk().await {
518                Ok(Some(chunk)) => {
519                    if body_bytes.len() + chunk.len() > max_response_size {
520                        return Err(response_too_large(
521                            body_bytes.len() + chunk.len(),
522                            max_response_size,
523                        ));
524                    }
525                    body_bytes.extend_from_slice(&chunk);
526                }
527                Ok(None) => break,
528                Err(e) => {
529                    return Err(OperationError::Http {
530                        status: Some(status),
531                        message: format!("failed to read response body: {e}"),
532                    });
533                }
534            }
535        }
536
537        let body = String::from_utf8_lossy(&body_bytes).into_owned();
538        let duration_ms = start.elapsed().as_millis() as u64;
539
540        debug!(
541            status,
542            body_len = body.len(),
543            duration_ms,
544            "http request completed"
545        );
546
547        #[cfg(feature = "prometheus")]
548        {
549            let status_label = status.to_string();
550            metrics::counter!(metric_names::HTTP_TOTAL, "method" => method_label, "status" => status_label).increment(1);
551            metrics::histogram!(metric_names::HTTP_DURATION_SECONDS)
552                .record(duration_ms as f64 / 1000.0);
553        }
554
555        Ok(HttpOutput {
556            status,
557            headers,
558            body,
559            duration_ms,
560        })
561    }
562}
563
564impl IntoFuture for Http {
565    type Output = Result<HttpOutput, OperationError>;
566    type IntoFuture = Pin<Box<dyn Future<Output = Self::Output> + Send>>;
567
568    fn into_future(self) -> Self::IntoFuture {
569        Box::pin(self.run())
570    }
571}
572
573/// Output of a completed HTTP request.
574///
575/// Contains the status code, response headers, body, and duration.
576#[derive(Debug)]
577pub struct HttpOutput {
578    status: u16,
579    headers: HashMap<String, String>,
580    body: String,
581    duration_ms: u64,
582}
583
584impl HttpOutput {
585    /// Return the HTTP status code (e.g. `200`, `404`).
586    pub fn status(&self) -> u16 {
587        self.status
588    }
589
590    /// Return the response headers as a string map.
591    pub fn headers(&self) -> &HashMap<String, String> {
592        &self.headers
593    }
594
595    /// Return the response body as text.
596    pub fn body(&self) -> &str {
597        &self.body
598    }
599
600    /// Deserialize the response body as JSON into the given type `T`.
601    ///
602    /// # Errors
603    ///
604    /// Returns [`OperationError::Deserialize`] if parsing fails.
605    pub fn json<T: DeserializeOwned>(&self) -> Result<T, OperationError> {
606        serde_json::from_str(&self.body).map_err(OperationError::deserialize::<T>)
607    }
608
609    /// Return the wall-clock duration of the request in milliseconds.
610    pub fn duration_ms(&self) -> u64 {
611        self.duration_ms
612    }
613
614    /// Return `true` if the status code is in the 2xx range.
615    pub fn is_success(&self) -> bool {
616        (200..300).contains(&self.status)
617    }
618}
619
620#[cfg(test)]
621mod tests;