ironflow_engine/executor/
http.rs1use 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
18pub struct HttpExecutor<'a> {
22 config: &'a HttpConfig,
23}
24
25impl<'a> HttpExecutor<'a> {
26 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}