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