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