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