Skip to main content

ironflow_engine/executor/
http.rs

1//! HTTP step executor.
2
3use std::sync::Arc;
4use std::time::{Duration, Instant};
5
6use rust_decimal::Decimal;
7use serde_json::json;
8use tracing::info;
9
10use ironflow_core::operations::http::Http;
11use ironflow_core::provider::AgentProvider;
12use ironflow_store::entities::StepKind;
13
14use crate::config::HttpConfig;
15use crate::error::EngineError;
16
17use super::{StepArtifacts, StepExecutor, StepOutput};
18
19/// Executor for HTTP steps.
20///
21/// Sends an HTTP request and captures the response status, headers, and body.
22pub struct HttpExecutor<'a> {
23    config: &'a HttpConfig,
24}
25
26impl<'a> HttpExecutor<'a> {
27    /// Create a new HTTP executor from a config reference.
28    pub fn new(config: &'a HttpConfig) -> Self {
29        Self { config }
30    }
31}
32
33impl StepExecutor for HttpExecutor<'_> {
34    fn kind(&self) -> StepKind {
35        StepKind::Http
36    }
37
38    async fn execute(&self, _provider: &Arc<dyn AgentProvider>) -> Result<StepOutput, EngineError> {
39        let start = Instant::now();
40
41        let mut http = match self.config.method.to_uppercase().as_str() {
42            "GET" => Http::get(&self.config.url),
43            "POST" => Http::post(&self.config.url),
44            "PUT" => Http::put(&self.config.url),
45            "PATCH" => Http::patch(&self.config.url),
46            "DELETE" => Http::delete(&self.config.url),
47            other => {
48                return Err(EngineError::StepConfig(format!(
49                    "unsupported HTTP method: {other}"
50                )));
51            }
52        };
53
54        for (name, value) in &self.config.headers {
55            http = http.header(name, value);
56        }
57        if let Some(ref ctx) = self.config.trace_context {
58            http = http.trace_context(ctx);
59        }
60        if let Some(ref body) = self.config.body {
61            http = http.json(body.clone());
62        }
63        if let Some(secs) = self.config.timeout_secs {
64            http = http.timeout(Duration::from_secs(secs));
65        }
66        for host in &self.config.allowed_hosts {
67            http = http.allow_host(host);
68        }
69
70        let output = http.run().await?;
71        let duration_ms = start.elapsed().as_millis() as u64;
72
73        info!(
74            step_kind = "http",
75            method = %self.config.method,
76            url = %self.config.url,
77            status = output.status(),
78            duration_ms,
79            "http step completed"
80        );
81
82        #[cfg(feature = "prometheus")]
83        {
84            use ironflow_core::metric_names::{HTTP_DURATION_SECONDS, HTTP_TOTAL, STATUS_SUCCESS};
85            use metrics::{counter, histogram};
86            counter!(HTTP_TOTAL, "method" => self.config.method.clone(), "status" => STATUS_SUCCESS).increment(1);
87            histogram!(HTTP_DURATION_SECONDS).record(duration_ms as f64 / 1000.0);
88        }
89
90        Ok(StepOutput {
91            output: json!({
92                "status": output.status(),
93                "headers": output.headers(),
94                "body": output.body(),
95            }),
96            duration_ms,
97            cost_usd: Decimal::ZERO,
98            input_tokens: None,
99            cache_read_input_tokens: None,
100            cache_creation_input_tokens: None,
101            output_tokens: None,
102            model: None,
103            debug_messages: None,
104            artifacts: StepArtifacts::default(),
105            account_id: None,
106            environment_id: None,
107        })
108    }
109}
110
111#[cfg(test)]
112mod tests {
113    use super::*;
114    use axum::Router;
115    use axum::routing::{delete, get, patch, post, put};
116    use ironflow_core::providers::claude::ClaudeCodeProvider;
117    use ironflow_core::providers::record_replay::RecordReplayProvider;
118    use tokio::net::TcpListener;
119
120    fn create_test_provider() -> Arc<dyn AgentProvider> {
121        let inner = ClaudeCodeProvider::new();
122        Arc::new(RecordReplayProvider::replay(
123            inner,
124            "/tmp/ironflow-fixtures",
125        ))
126    }
127
128    async fn start_test_server() -> String {
129        let app = Router::new()
130            .route("/status/200", get(|| async { "ok" }))
131            .route("/post", post(|| async { "ok" }))
132            .route("/put", put(|| async { "ok" }))
133            .route("/patch", patch(|| async { "ok" }))
134            .route("/delete", delete(|| async { "ok" }))
135            .route("/headers", get(|| async { "ok" }));
136
137        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
138        let port = listener.local_addr().unwrap().port();
139        tokio::spawn(async move {
140            axum::serve(listener, app).await.unwrap();
141        });
142        format!("http://localhost:{port}")
143    }
144
145    #[tokio::test]
146    async fn http_internal_host_refused_without_allow_host() {
147        let base = start_test_server().await;
148        let config = HttpConfig::get(&format!("{base}/status/200"));
149        let executor = HttpExecutor::new(&config);
150        let provider = create_test_provider();
151
152        let err = executor.execute(&provider).await.unwrap_err();
153        assert!(
154            err.to_string()
155                .contains("URL host localhost resolves to a blocked IP address"),
156            "{err}"
157        );
158    }
159
160    #[tokio::test]
161    async fn http_get_method() {
162        let base = start_test_server().await;
163        let config = HttpConfig::get(&format!("{base}/status/200")).allow_host("localhost");
164        let executor = HttpExecutor::new(&config);
165        let provider = create_test_provider();
166
167        let result = executor.execute(&provider).await;
168        assert!(result.is_ok());
169        let output = result.unwrap();
170        assert!(output.output.get("status").is_some());
171        assert!(output.output.get("headers").is_some());
172        assert!(output.output.get("body").is_some());
173    }
174
175    #[tokio::test]
176    async fn http_post_method() {
177        let base = start_test_server().await;
178        let config = HttpConfig::post(&format!("{base}/post")).allow_host("localhost");
179        let executor = HttpExecutor::new(&config);
180        let provider = create_test_provider();
181
182        let result = executor.execute(&provider).await;
183        assert!(result.is_ok());
184    }
185
186    #[tokio::test]
187    async fn http_put_method() {
188        let base = start_test_server().await;
189        let config = HttpConfig::put(&format!("{base}/put")).allow_host("localhost");
190        let executor = HttpExecutor::new(&config);
191        let provider = create_test_provider();
192
193        let result = executor.execute(&provider).await;
194        assert!(result.is_ok());
195    }
196
197    #[tokio::test]
198    async fn http_patch_method() {
199        let base = start_test_server().await;
200        let config = HttpConfig::patch(&format!("{base}/patch")).allow_host("localhost");
201        let executor = HttpExecutor::new(&config);
202        let provider = create_test_provider();
203
204        let result = executor.execute(&provider).await;
205        assert!(result.is_ok());
206    }
207
208    #[tokio::test]
209    async fn http_delete_method() {
210        let base = start_test_server().await;
211        let config = HttpConfig::delete(&format!("{base}/delete")).allow_host("localhost");
212        let executor = HttpExecutor::new(&config);
213        let provider = create_test_provider();
214
215        let result = executor.execute(&provider).await;
216        assert!(result.is_ok());
217    }
218
219    #[tokio::test]
220    async fn http_unsupported_method_returns_error() {
221        let base = start_test_server().await;
222        let mut config = HttpConfig::get(&format!("{base}/status/200"));
223        config.method = "INVALID".to_string();
224        let executor = HttpExecutor::new(&config);
225        let provider = create_test_provider();
226
227        let result = executor.execute(&provider).await;
228        assert!(result.is_err());
229        match result {
230            Err(EngineError::StepConfig(msg)) => {
231                assert!(msg.contains("unsupported HTTP method"));
232            }
233            _ => panic!("expected StepConfig error"),
234        }
235    }
236
237    #[tokio::test]
238    async fn http_with_custom_headers() {
239        let base = start_test_server().await;
240        let config = HttpConfig::get(&format!("{base}/headers"))
241            .header("X-Custom-Header", "test-value")
242            .header("Authorization", "Bearer token")
243            .allow_host("localhost");
244        let executor = HttpExecutor::new(&config);
245        let provider = create_test_provider();
246
247        let result = executor.execute(&provider).await;
248        assert!(result.is_ok());
249    }
250
251    #[tokio::test]
252    async fn http_with_json_body() {
253        let base = start_test_server().await;
254        let config = HttpConfig::post(&format!("{base}/post"))
255            .json(json!({"key": "value", "number": 42}))
256            .allow_host("localhost");
257        let executor = HttpExecutor::new(&config);
258        let provider = create_test_provider();
259
260        let result = executor.execute(&provider).await;
261        assert!(result.is_ok());
262    }
263
264    #[tokio::test]
265    async fn http_step_output_has_structure() {
266        let base = start_test_server().await;
267        let config = HttpConfig::get(&format!("{base}/status/200")).allow_host("localhost");
268        let executor = HttpExecutor::new(&config);
269        let provider = create_test_provider();
270
271        let output = executor.execute(&provider).await.unwrap();
272        assert!(output.output.get("status").is_some());
273        assert!(output.output.get("headers").is_some());
274        assert!(output.output.get("body").is_some());
275        assert_eq!(output.cost_usd, Decimal::ZERO);
276    }
277}