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;
12use ironflow_store::entities::StepKind;
13
14use crate::config::HttpConfig;
15use crate::error::EngineError;
16
17use super::{StepArtifacts, StepExecutor, StepOutput};
18
19pub struct HttpExecutor<'a> {
23 config: &'a HttpConfig,
24}
25
26impl<'a> HttpExecutor<'a> {
27 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 account_id: None,
103 })
104 }
105}
106
107#[cfg(test)]
108mod tests {
109 use super::*;
110 use axum::Router;
111 use axum::routing::{delete, get, patch, post, put};
112 use ironflow_core::providers::claude::ClaudeCodeProvider;
113 use ironflow_core::providers::record_replay::RecordReplayProvider;
114 use tokio::net::TcpListener;
115
116 fn create_test_provider() -> Arc<dyn AgentProvider> {
117 let inner = ClaudeCodeProvider::new();
118 Arc::new(RecordReplayProvider::replay(
119 inner,
120 "/tmp/ironflow-fixtures",
121 ))
122 }
123
124 async fn start_test_server() -> String {
125 let app = Router::new()
126 .route("/status/200", get(|| async { "ok" }))
127 .route("/post", post(|| async { "ok" }))
128 .route("/put", put(|| async { "ok" }))
129 .route("/patch", patch(|| async { "ok" }))
130 .route("/delete", delete(|| async { "ok" }))
131 .route("/headers", get(|| async { "ok" }));
132
133 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
134 let port = listener.local_addr().unwrap().port();
135 tokio::spawn(async move {
136 axum::serve(listener, app).await.unwrap();
137 });
138 format!("http://localhost:{port}")
139 }
140
141 #[tokio::test]
142 async fn http_get_method() {
143 let base = start_test_server().await;
144 let config = HttpConfig::get(&format!("{base}/status/200"));
145 let executor = HttpExecutor::new(&config);
146 let provider = create_test_provider();
147
148 let result = executor.execute(&provider).await;
149 assert!(result.is_ok());
150 let output = result.unwrap();
151 assert!(output.output.get("status").is_some());
152 assert!(output.output.get("headers").is_some());
153 assert!(output.output.get("body").is_some());
154 }
155
156 #[tokio::test]
157 async fn http_post_method() {
158 let base = start_test_server().await;
159 let config = HttpConfig::post(&format!("{base}/post"));
160 let executor = HttpExecutor::new(&config);
161 let provider = create_test_provider();
162
163 let result = executor.execute(&provider).await;
164 assert!(result.is_ok());
165 }
166
167 #[tokio::test]
168 async fn http_put_method() {
169 let base = start_test_server().await;
170 let config = HttpConfig::put(&format!("{base}/put"));
171 let executor = HttpExecutor::new(&config);
172 let provider = create_test_provider();
173
174 let result = executor.execute(&provider).await;
175 assert!(result.is_ok());
176 }
177
178 #[tokio::test]
179 async fn http_patch_method() {
180 let base = start_test_server().await;
181 let config = HttpConfig::patch(&format!("{base}/patch"));
182 let executor = HttpExecutor::new(&config);
183 let provider = create_test_provider();
184
185 let result = executor.execute(&provider).await;
186 assert!(result.is_ok());
187 }
188
189 #[tokio::test]
190 async fn http_delete_method() {
191 let base = start_test_server().await;
192 let config = HttpConfig::delete(&format!("{base}/delete"));
193 let executor = HttpExecutor::new(&config);
194 let provider = create_test_provider();
195
196 let result = executor.execute(&provider).await;
197 assert!(result.is_ok());
198 }
199
200 #[tokio::test]
201 async fn http_unsupported_method_returns_error() {
202 let base = start_test_server().await;
203 let mut config = HttpConfig::get(&format!("{base}/status/200"));
204 config.method = "INVALID".to_string();
205 let executor = HttpExecutor::new(&config);
206 let provider = create_test_provider();
207
208 let result = executor.execute(&provider).await;
209 assert!(result.is_err());
210 match result {
211 Err(EngineError::StepConfig(msg)) => {
212 assert!(msg.contains("unsupported HTTP method"));
213 }
214 _ => panic!("expected StepConfig error"),
215 }
216 }
217
218 #[tokio::test]
219 async fn http_with_custom_headers() {
220 let base = start_test_server().await;
221 let config = HttpConfig::get(&format!("{base}/headers"))
222 .header("X-Custom-Header", "test-value")
223 .header("Authorization", "Bearer token");
224 let executor = HttpExecutor::new(&config);
225 let provider = create_test_provider();
226
227 let result = executor.execute(&provider).await;
228 assert!(result.is_ok());
229 }
230
231 #[tokio::test]
232 async fn http_with_json_body() {
233 let base = start_test_server().await;
234 let config =
235 HttpConfig::post(&format!("{base}/post")).json(json!({"key": "value", "number": 42}));
236 let executor = HttpExecutor::new(&config);
237 let provider = create_test_provider();
238
239 let result = executor.execute(&provider).await;
240 assert!(result.is_ok());
241 }
242
243 #[tokio::test]
244 async fn http_step_output_has_structure() {
245 let base = start_test_server().await;
246 let config = HttpConfig::get(&format!("{base}/status/200"));
247 let executor = HttpExecutor::new(&config);
248 let provider = create_test_provider();
249
250 let output = executor.execute(&provider).await.unwrap();
251 assert!(output.output.get("status").is_some());
252 assert!(output.output.get("headers").is_some());
253 assert!(output.output.get("body").is_some());
254 assert_eq!(output.cost_usd, Decimal::ZERO);
255 }
256}