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;