ironflow_engine/executor/
shell.rs1use std::os::unix::process::ExitStatusExt;
4use std::process::{ExitStatus, Stdio};
5use std::sync::Arc;
6use std::time::{Duration, Instant};
7
8use rust_decimal::Decimal;
9use serde_json::json;
10use tokio::io::{AsyncBufReadExt, BufReader};
11use tokio::process::Command;
12use tokio::spawn;
13use tracing::info;
14
15use ironflow_core::dry_run::is_dry_run;
16use ironflow_core::error::OperationError;
17use ironflow_core::operations::shell::Shell;
18use ironflow_core::provider::AgentProvider;
19use ironflow_core::utils::truncate_output;
20use ironflow_store::entities::StepKind;
21
22use crate::config::ShellConfig;
23use crate::error::EngineError;
24use crate::log_sender::StepLogSender;
25use crate::notify::LogStream;
26
27use super::{StepArtifacts, StepExecutor, StepOutput};
28
29const DEFAULT_SHELL_TIMEOUT: Duration = Duration::from_secs(300);
30
31async fn read_and_stream<R: tokio::io::AsyncRead + Unpin>(
34 reader: R,
35 sender: StepLogSender,
36 stream: LogStream,
37) -> String {
38 let mut lines = BufReader::new(reader).lines();
39 let mut collected = String::new();
40 while let Ok(Some(line)) = lines.next_line().await {
41 sender.emit(stream, &line);
42 if !collected.is_empty() {
43 collected.push('\n');
44 }
45 collected.push_str(&line);
46 }
47 collected
48}
49
50fn exit_code_of(status: ExitStatus) -> i32 {
53 status
54 .code()
55 .or_else(|| status.signal().map(|signal| -signal))
56 .unwrap_or(-1)
57}
58
59fn build_output(stdout: &str, stderr: &str, exit_code: i32, duration_ms: u64) -> StepOutput {
61 StepOutput {
62 output: json!({
63 "stdout": stdout,
64 "stderr": stderr,
65 "exit_code": exit_code,
66 }),
67 duration_ms,
68 cost_usd: Decimal::ZERO,
69 input_tokens: None,
70 cache_read_input_tokens: None,
71 cache_creation_input_tokens: None,
72 output_tokens: None,
73 model: None,
74 debug_messages: None,
75 artifacts: StepArtifacts::default(),
76 }
77}
78
79pub struct ShellExecutor<'a> {
85 config: &'a ShellConfig,
86 log_sender: Option<StepLogSender>,
87}
88
89impl<'a> ShellExecutor<'a> {
90 pub fn new(config: &'a ShellConfig) -> Self {
92 Self {
93 config,
94 log_sender: None,
95 }
96 }
97
98 pub fn with_log_sender(mut self, sender: StepLogSender) -> Self {
100 self.log_sender = Some(sender);
101 self
102 }
103}
104
105impl StepExecutor for ShellExecutor<'_> {
106 fn kind(&self) -> StepKind {
107 StepKind::Shell
108 }
109
110 async fn execute(&self, _provider: &Arc<dyn AgentProvider>) -> Result<StepOutput, EngineError> {
111 match self.log_sender {
112 Some(ref sender) => self.execute_streaming(sender.clone()).await,
113 None => self.execute_buffered().await,
114 }
115 }
116}
117
118impl ShellExecutor<'_> {
119 async fn execute_buffered(&self) -> Result<StepOutput, EngineError> {
121 let start = Instant::now();
122
123 let mut shell = Shell::new(&self.config.command);
124 if let Some(secs) = self.config.timeout_secs {
125 shell = shell.timeout(Duration::from_secs(secs));
126 }
127 if let Some(ref dir) = self.config.dir {
128 shell = shell.dir(dir);
129 }
130 for (key, value) in &self.config.env {
131 shell = shell.env(key, value);
132 }
133 if self.config.clean_env {
134 shell = shell.clean_env();
135 }
136
137 let (stdout, stderr, exit_code) = if self.config.exit_code_as_output && !is_dry_run() {
141 self.run_capturing().await?
142 } else {
143 let output = shell.run().await?;
144 (
145 output.stdout().to_string(),
146 output.stderr().to_string(),
147 output.exit_code(),
148 )
149 };
150 let duration_ms = start.elapsed().as_millis() as u64;
151
152 info!(
153 step_kind = "shell",
154 command = %self.config.command,
155 exit_code,
156 duration_ms,
157 "shell step completed"
158 );
159
160 self.record_metrics(duration_ms);
161
162 Ok(build_output(&stdout, &stderr, exit_code, duration_ms))
163 }
164
165 async fn run_capturing(&self) -> Result<(String, String, i32), EngineError> {
168 let mut cmd = Command::new("sh");
169 cmd.arg("-c").arg(&self.config.command);
170 cmd.stdout(Stdio::piped())
171 .stderr(Stdio::piped())
172 .kill_on_drop(true);
173
174 if self.config.clean_env {
175 cmd.env_clear();
176 }
177 if let Some(ref dir) = self.config.dir {
178 cmd.current_dir(dir);
179 }
180 for (key, value) in &self.config.env {
181 cmd.env(key, value);
182 }
183
184 let child = cmd.spawn().map_err(|e| {
185 EngineError::Operation(OperationError::Shell {
186 exit_code: -1,
187 stderr: format!("failed to spawn shell: {e}"),
188 })
189 })?;
190
191 let timeout_dur = self
192 .config
193 .timeout_secs
194 .map(Duration::from_secs)
195 .unwrap_or(DEFAULT_SHELL_TIMEOUT);
196
197 let output = match tokio::time::timeout(timeout_dur, child.wait_with_output()).await {
198 Ok(Ok(output)) => output,
199 Ok(Err(e)) => {
200 return Err(EngineError::Operation(OperationError::Shell {
201 exit_code: -1,
202 stderr: format!("failed to wait for shell: {e}"),
203 }));
204 }
205 Err(_) => {
206 return Err(EngineError::Operation(OperationError::Timeout {
207 step: self.config.command.clone(),
208 limit: timeout_dur,
209 }));
210 }
211 };
212
213 let exit_code = exit_code_of(output.status);
214 let stdout = truncate_output(&output.stdout, "shell stdout");
215 let stderr = truncate_output(&output.stderr, "shell stderr");
216 Ok((stdout, stderr, exit_code))
217 }
218
219 async fn execute_streaming(&self, sender: StepLogSender) -> Result<StepOutput, EngineError> {
222 let start = Instant::now();
223
224 let mut cmd = Command::new("sh");
225 cmd.arg("-c").arg(&self.config.command);
226 cmd.stdout(Stdio::piped())
227 .stderr(Stdio::piped())
228 .kill_on_drop(true);
229
230 if self.config.clean_env {
231 cmd.env_clear();
232 }
233 if let Some(ref dir) = self.config.dir {
234 cmd.current_dir(dir);
235 }
236 for (key, value) in &self.config.env {
237 cmd.env(key, value);
238 }
239
240 let mut child = cmd.spawn().map_err(|e| {
241 EngineError::Operation(OperationError::Shell {
242 exit_code: -1,
243 stderr: format!("failed to spawn shell: {e}"),
244 })
245 })?;
246
247 let stdout_pipe = child.stdout.take().expect("stdout piped");
248 let stderr_pipe = child.stderr.take().expect("stderr piped");
249
250 let stdout_task = spawn(read_and_stream(
251 stdout_pipe,
252 sender.clone(),
253 LogStream::Stdout,
254 ));
255 let stderr_task = spawn(read_and_stream(stderr_pipe, sender, LogStream::Stderr));
256
257 let timeout_dur = self
258 .config
259 .timeout_secs
260 .map(Duration::from_secs)
261 .unwrap_or(DEFAULT_SHELL_TIMEOUT);
262
263 let status = match tokio::time::timeout(timeout_dur, child.wait()).await {
264 Ok(Ok(status)) => status,
265 Ok(Err(e)) => {
266 return Err(EngineError::Operation(OperationError::Shell {
267 exit_code: -1,
268 stderr: format!("failed to wait for shell: {e}"),
269 }));
270 }
271 Err(_) => {
272 child.kill().await.ok();
273 return Err(EngineError::Operation(OperationError::Timeout {
274 step: self.config.command.clone(),
275 limit: timeout_dur,
276 }));
277 }
278 };
279
280 let raw_stdout = stdout_task.await.unwrap_or_default();
281 let raw_stderr = stderr_task.await.unwrap_or_default();
282
283 let stdout = truncate_output(raw_stdout.as_bytes(), "shell stdout");
284 let stderr = truncate_output(raw_stderr.as_bytes(), "shell stderr");
285
286 let exit_code = exit_code_of(status);
287 let duration_ms = start.elapsed().as_millis() as u64;
288
289 info!(
290 step_kind = "shell",
291 command = %self.config.command,
292 exit_code,
293 duration_ms,
294 streaming = true,
295 "shell step completed"
296 );
297
298 self.record_metrics(duration_ms);
299
300 if exit_code != 0 && !self.config.exit_code_as_output {
301 return Err(EngineError::Operation(OperationError::Shell {
302 exit_code,
303 stderr: stderr.clone(),
304 }));
305 }
306
307 Ok(build_output(&stdout, &stderr, exit_code, duration_ms))
308 }
309
310 #[allow(unused_variables)]
311 fn record_metrics(&self, duration_ms: u64) {
312 #[cfg(feature = "prometheus")]
313 {
314 use ironflow_core::metric_names::{
315 SHELL_DURATION_SECONDS, SHELL_TOTAL, STATUS_SUCCESS,
316 };
317 use metrics::{counter, histogram};
318 counter!(SHELL_TOTAL, "status" => STATUS_SUCCESS).increment(1);
319 histogram!(SHELL_DURATION_SECONDS).record(duration_ms as f64 / 1000.0);
320 }
321 }
322}
323
324#[cfg(test)]
325mod tests {
326 use super::*;
327 use ironflow_core::providers::claude::ClaudeCodeProvider;
328 use ironflow_core::providers::record_replay::RecordReplayProvider;
329
330 fn create_test_provider() -> Arc<dyn AgentProvider> {
331 let inner = ClaudeCodeProvider::new();
332 Arc::new(RecordReplayProvider::replay(
333 inner,
334 "/tmp/ironflow-fixtures",
335 ))
336 }
337
338 #[tokio::test]
339 async fn shell_simple_command() {
340 let config = ShellConfig::new("echo hello");
341 let executor = ShellExecutor::new(&config);
342 let provider = create_test_provider();
343
344 let result = executor.execute(&provider).await;
345 assert!(result.is_ok());
346 let output = result.unwrap();
347 assert_eq!(output.output["exit_code"].as_i64().unwrap(), 0);
348 assert!(output.output["stdout"].as_str().unwrap().contains("hello"));
349 }
350
351 #[tokio::test]
352 async fn shell_nonzero_exit_returns_error() {
353 let config = ShellConfig::new("exit 1");
354 let executor = ShellExecutor::new(&config);
355 let provider = create_test_provider();
356
357 let result = executor.execute(&provider).await;
358 assert!(result.is_err());
359 }
360
361 #[tokio::test]
362 async fn shell_env_variables() {
363 let config = ShellConfig::new("echo $MY_VAR").env("MY_VAR", "test_value");
364 let executor = ShellExecutor::new(&config);
365 let provider = create_test_provider();
366
367 let result = executor.execute(&provider).await;
368 assert!(result.is_ok());
369 let output = result.unwrap();
370 assert!(
371 output.output["stdout"]
372 .as_str()
373 .unwrap()
374 .contains("test_value")
375 );
376 }
377
378 #[tokio::test]
379 async fn shell_step_output_has_structure() {
380 let config = ShellConfig::new("echo test");
381 let executor = ShellExecutor::new(&config);
382 let provider = create_test_provider();
383
384 let output = executor.execute(&provider).await.unwrap();
385 assert!(output.output.get("stdout").is_some());
386 assert!(output.output.get("stderr").is_some());
387 assert!(output.output.get("exit_code").is_some());
388 assert_eq!(output.cost_usd, Decimal::ZERO);
389 assert!(output.duration_ms < 5000);
390 }
391
392 #[tokio::test]
393 async fn shell_command_with_pipe() {
394 let config = ShellConfig::new("echo hello | grep hello");
395 let executor = ShellExecutor::new(&config);
396 let provider = create_test_provider();
397
398 let result = executor.execute(&provider).await;
399 assert!(result.is_ok());
400 let output = result.unwrap();
401 assert_eq!(output.output["exit_code"].as_i64().unwrap(), 0);
402 assert!(output.output["stdout"].as_str().unwrap().contains("hello"));
403 }
404
405 #[tokio::test]
406 async fn shell_streaming_emits_lines() {
407 let config = ShellConfig::new("echo line1 && echo line2");
408 let (sender, mut receiver) = crate::log_sender::channel();
409 let step_sender = StepLogSender::new(
410 sender,
411 uuid::Uuid::now_v7(),
412 uuid::Uuid::now_v7(),
413 "test".to_string(),
414 );
415 let executor = ShellExecutor::new(&config).with_log_sender(step_sender);
416 let provider = create_test_provider();
417
418 let result = executor.execute(&provider).await;
419 assert!(result.is_ok());
420
421 let output = result.unwrap();
422 assert!(output.output["stdout"].as_str().unwrap().contains("line1"));
423 assert!(output.output["stdout"].as_str().unwrap().contains("line2"));
424
425 let mut lines = Vec::new();
426 while let Ok(line) = receiver.try_recv() {
427 lines.push(line);
428 }
429 assert!(lines.len() >= 2);
430 assert_eq!(lines[0].stream, LogStream::Stdout);
431 assert_eq!(lines[0].line, "line1");
432 assert_eq!(lines[1].line, "line2");
433 }
434
435 #[tokio::test]
436 async fn shell_streaming_captures_stderr() {
437 let config = ShellConfig::new("echo err >&2");
438 let (sender, mut receiver) = crate::log_sender::channel();
439 let step_sender = StepLogSender::new(
440 sender,
441 uuid::Uuid::now_v7(),
442 uuid::Uuid::now_v7(),
443 "test".to_string(),
444 );
445 let executor = ShellExecutor::new(&config).with_log_sender(step_sender);
446 let provider = create_test_provider();
447
448 let result = executor.execute(&provider).await;
449 assert!(result.is_ok());
450
451 let mut stderr_lines = Vec::new();
452 while let Ok(line) = receiver.try_recv() {
453 if line.stream == LogStream::Stderr {
454 stderr_lines.push(line);
455 }
456 }
457 assert!(!stderr_lines.is_empty());
458 assert_eq!(stderr_lines[0].line, "err");
459 }
460
461 fn streaming_sender() -> StepLogSender {
462 let (sender, _receiver) = crate::log_sender::channel();
463 StepLogSender::new(
464 sender,
465 uuid::Uuid::now_v7(),
466 uuid::Uuid::now_v7(),
467 "test".to_string(),
468 )
469 }
470
471 #[tokio::test]
472 async fn shell_exit_code_as_output_buffered_keeps_code_and_streams() {
473 let config = ShellConfig::new("echo out; echo err >&2; exit 3").exit_code_as_output();
474 let executor = ShellExecutor::new(&config);
475 let provider = create_test_provider();
476
477 let output = executor
478 .execute(&provider)
479 .await
480 .expect("a non-zero exit is an output");
481 assert_eq!(output.output["exit_code"], 3);
482 assert_eq!(output.stdout().trim(), "out");
483 assert_eq!(output.stderr().trim(), "err");
484 assert!(!output.is_success());
485 assert_eq!(output.exit_code(), Some(3));
486 }
487
488 #[tokio::test]
489 async fn shell_exit_code_as_output_streaming_keeps_code() {
490 let config = ShellConfig::new("echo out; echo err >&2; exit 3").exit_code_as_output();
491 let executor = ShellExecutor::new(&config).with_log_sender(streaming_sender());
492 let provider = create_test_provider();
493
494 let output = executor
495 .execute(&provider)
496 .await
497 .expect("a non-zero exit is an output");
498 assert_eq!(output.output["exit_code"], 3);
499 assert_eq!(output.stdout(), "out");
500 assert_eq!(output.stderr(), "err");
501 assert!(!output.is_success());
502 assert_eq!(output.exit_code(), Some(3));
503 }
504
505 #[tokio::test]
506 async fn shell_exit_code_as_output_zero_exit_is_success() {
507 let config = ShellConfig::new("echo fine").exit_code_as_output();
508 let provider = create_test_provider();
509
510 let buffered = ShellExecutor::new(&config)
511 .execute(&provider)
512 .await
513 .expect("exit 0 succeeds");
514 assert!(buffered.is_success());
515 assert_eq!(buffered.exit_code(), Some(0));
516
517 let streaming = ShellExecutor::new(&config)
518 .with_log_sender(streaming_sender())
519 .execute(&provider)
520 .await
521 .expect("exit 0 succeeds");
522 assert!(streaming.is_success());
523 }
524
525 #[tokio::test]
526 async fn shell_exit_code_as_output_still_errors_on_timeout() {
527 let config = ShellConfig::new("sleep 5")
528 .timeout_secs(1)
529 .exit_code_as_output();
530 let provider = create_test_provider();
531
532 let buffered = ShellExecutor::new(&config).execute(&provider).await;
533 assert!(matches!(
534 buffered,
535 Err(EngineError::Operation(OperationError::Timeout { .. }))
536 ));
537
538 let streaming = ShellExecutor::new(&config)
539 .with_log_sender(streaming_sender())
540 .execute(&provider)
541 .await;
542 assert!(matches!(
543 streaming,
544 Err(EngineError::Operation(OperationError::Timeout { .. }))
545 ));
546 }
547
548 #[tokio::test]
549 async fn shell_without_exit_code_as_output_nonzero_still_errors() {
550 let config = ShellConfig::new("echo out; exit 3");
551 let provider = create_test_provider();
552
553 let buffered = ShellExecutor::new(&config).execute(&provider).await;
554 assert!(matches!(
555 buffered,
556 Err(EngineError::Operation(OperationError::Shell {
557 exit_code: 3,
558 ..
559 }))
560 ));
561
562 let streaming = ShellExecutor::new(&config)
563 .with_log_sender(streaming_sender())
564 .execute(&provider)
565 .await;
566 assert!(matches!(
567 streaming,
568 Err(EngineError::Operation(OperationError::Shell {
569 exit_code: 3,
570 ..
571 }))
572 ));
573 }
574
575 #[tokio::test]
576 async fn shell_streaming_nonzero_exit_returns_error() {
577 let config = ShellConfig::new("exit 42");
578 let (sender, _receiver) = crate::log_sender::channel();
579 let step_sender = StepLogSender::new(
580 sender,
581 uuid::Uuid::now_v7(),
582 uuid::Uuid::now_v7(),
583 "test".to_string(),
584 );
585 let executor = ShellExecutor::new(&config).with_log_sender(step_sender);
586 let provider = create_test_provider();
587
588 let result = executor.execute(&provider).await;
589 assert!(result.is_err());
590 }
591}