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
67        let output = http.run().await?;
68        let duration_ms = start.elapsed().as_millis() as u64;
69
70        info!(
71            step_kind = "http",
72            method = %self.config.method,
73            url = %self.config.url,
74            status = output.status(),
75            duration_ms,
76            "http step completed"
77        );
78
79        #[cfg(feature = "prometheus")]
80        {
81            use ironflow_core::metric_names::{HTTP_DURATION_SECONDS, HTTP_TOTAL, STATUS_SUCCESS};
82            use metrics::{counter, histogram};
83            counter!(HTTP_TOTAL, "method" => self.config.method.clone(), "status" => STATUS_SUCCESS).increment(1);
84            histogram!(HTTP_DURATION_SECONDS).record(duration_ms as f64 / 1000.0);
85        }
86
87        Ok(StepOutput {
88            output: json!({
89                "status": output.status(),
90                "headers": output.headers(),
91                "body": output.body(),
92            }),
93            duration_ms,
94            cost_usd: Decimal::ZERO,
95            input_tokens: None,
96            cache_read_input_tokens: None,
97            cache_creation_input_tokens: None,
98            output_tokens: None,
99            model: None,
100            debug_messages: None,
101            artifacts: StepArtifacts::default(),
102        })
103    }
104}
105
106#[cfg(test)]
107mod tests {
108    use super::*;
109    use axum::Router;
110    use axum::routing::{delete, get, patch, post, put};
111    use ironflow_core::providers::claude::ClaudeCodeProvider;
112    use ironflow_core::providers::record_replay::RecordReplayProvider;
113    use tokio::net::TcpListener;
114
115    fn create_test_provider() -> Arc<dyn AgentProvider> {
116        let inner = ClaudeCodeProvider::new();
117        Arc::new(RecordReplayProvider::replay(
118            inner,
119            "/tmp/ironflow-fixtures",
120        ))
121    }
122
123    async fn start_test_server() -> String {
124        let app = Router::new()
125            .route("/status/200", get(|| async { "ok" }))
126            .route("/post", post(|| async { "ok" }))
127            .route("/put", put(|| async { "ok" }))
128            .route("/patch", patch(|| async { "ok" }))
129            .route("/delete", delete(|| async { "ok" }))
130            .route("/headers", get(|| async { "ok" }));
131
132        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
133        let port = listener.local_addr().unwrap().port();
134        tokio::spawn(async move {
135            axum::serve(listener, app).await.unwrap();
136        });
137        format!("http://localhost:{port}")
138    }
139
140    #[tokio::test]
141    async fn http_get_method() {
142        let base = start_test_server().await;
143        let config = HttpConfig::get(&format!("{base}/status/200"));
144        let executor = HttpExecutor::new(&config);
145        let provider = create_test_provider();
146
147        let result = executor.execute(&provider).await;
148        assert!(result.is_ok());
149        let output = result.unwrap();
150        assert!(output.output.get("status").is_some());
151        assert!(output.output.get("headers").is_some());
152        assert!(output.output.get("body").is_some());
153    }
154
155    #[tokio::test]
156    async fn http_post_method() {
157        let base = start_test_server().await;
158        let config = HttpConfig::post(&format!("{base}/post"));
159        let executor = HttpExecutor::new(&config);
160        let provider = create_test_provider();
161
162        let result = executor.execute(&provider).await;
163        assert!(result.is_ok());
164    }
165
166    #[tokio::test]
167    async fn http_put_method() {
168        let base = start_test_server().await;
169        let config = HttpConfig::put(&format!("{base}/put"));
170        let executor = HttpExecutor::new(&config);
171        let provider = create_test_provider();
172
173        let result = executor.execute(&provider).await;
174        assert!(result.is_ok());
175    }
176
177    #[tokio::test]
178    async fn http_patch_method() {
179        let base = start_test_server().await;
180        let config = HttpConfig::patch(&format!("{base}/patch"));
181        let executor = HttpExecutor::new(&config);
182        let provider = create_test_provider();
183
184        let result = executor.execute(&provider).await;
185        assert!(result.is_ok());
186    }
187
188    #[tokio::test]
189    async fn http_delete_method() {
190        let base = start_test_server().await;
191        let config = HttpConfig::delete(&format!("{base}/delete"));
192        let executor = HttpExecutor::new(&config);
193        let provider = create_test_provider();
194
195        let result = executor.execute(&provider).await;
196        assert!(result.is_ok());
197    }
198
199    #[tokio::test]
200    async fn http_unsupported_method_returns_error() {
201        let base = start_test_server().await;
202        let mut config = HttpConfig::get(&format!("{base}/status/200"));
203        config.method = "INVALID".to_string();
204        let executor = HttpExecutor::new(&config);
205        let provider = create_test_provider();
206
207        let result = executor.execute(&provider).await;
208        assert!(result.is_err());
209        match result {
210            Err(EngineError::StepConfig(msg)) => {
211                assert!(msg.contains("unsupported HTTP method"));
212            }
213            _ => panic!("expected StepConfig error"),
214        }
215    }
216
217    #[tokio::test]
218    async fn http_with_custom_headers() {
219        let base = start_test_server().await;
220        let config = HttpConfig::get(&format!("{base}/headers"))
221            .header("X-Custom-Header", "test-value")
222            .header("Authorization", "Bearer token");
223        let executor = HttpExecutor::new(&config);
224        let provider = create_test_provider();
225
226        let result = executor.execute(&provider).await;
227        assert!(result.is_ok());
228    }
229
230    #[tokio::test]
231    async fn http_with_json_body() {
232        let base = start_test_server().await;
233        let config =
234            HttpConfig::post(&format!("{base}/post")).json(json!({"key": "value", "number": 42}));
235        let executor = HttpExecutor::new(&config);
236        let provider = create_test_provider();
237
238        let result = executor.execute(&provider).await;
239        assert!(result.is_ok());
240    }
241
242    #[tokio::test]
243    async fn http_step_output_has_structure() {
244        let base = start_test_server().await;
245        let config = HttpConfig::get(&format!("{base}/status/200"));
246        let executor = HttpExecutor::new(&config);
247        let provider = create_test_provider();
248
249        let output = executor.execute(&provider).await.unwrap();
250        assert!(output.output.get("status").is_some());
251        assert!(output.output.get("headers").is_some());
252        assert!(output.output.get("body").is_some());
253        assert_eq!(output.cost_usd, Decimal::ZERO);
254    }
255}