1use 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 account_id: None,
77 environment_id: None,
78 }
79}
80
81pub struct ShellExecutor<'a> {
87 config: &'a ShellConfig,
88 log_sender: Option<StepLogSender>,
89}
90
91impl<'a> ShellExecutor<'a> {
92 pub fn new(config: &'a ShellConfig) -> Self {
94 Self {
95 config,
96 log_sender: None,
97 }
98 }
99
100 pub fn with_log_sender(mut self, sender: StepLogSender) -> Self {
102 self.log_sender = Some(sender);
103 self
104 }
105}
106
107impl StepExecutor for ShellExecutor<'_> {
108 fn kind(&self) -> StepKind {
109 StepKind::Shell
110 }
111
112 async fn execute(&self, _provider: &Arc<dyn AgentProvider>) -> Result<StepOutput, EngineError> {
113 match self.log_sender {
114 Some(ref sender) => self.execute_streaming(sender.clone()).await,
115 None => self.execute_buffered().await,
116 }
117 }
118}
119
120impl ShellExecutor<'_> {
121 fn command(&self) -> Command {
124 match self.config.args {
125 Some(ref args) => {
126 let mut cmd = Command::new(&self.config.command);
127 cmd.args(args);
128 cmd
129 }
130 None => {
131 let mut cmd = Command::new("sh");
132 cmd.arg("-c").arg(&self.config.command);
133 cmd
134 }
135 }
136 }
137
138 fn display(&self) -> String {
140 match self.config.args {
141 Some(ref args) => [self.config.command.as_str()]
142 .into_iter()
143 .chain(args.iter().map(String::as_str))
144 .collect::<Vec<_>>()
145 .join(" "),
146 None => self.config.command.clone(),
147 }
148 }
149
150 async fn execute_buffered(&self) -> Result<StepOutput, EngineError> {
152 let start = Instant::now();
153
154 let mut shell = match self.config.args {
155 Some(ref args) => Shell::exec(
156 &self.config.command,
157 &args.iter().map(String::as_str).collect::<Vec<_>>(),
158 ),
159 None => Shell::new(&self.config.command),
160 };
161 if let Some(secs) = self.config.timeout_secs {
162 shell = shell.timeout(Duration::from_secs(secs));
163 }
164 if let Some(ref dir) = self.config.dir {
165 shell = shell.dir(dir);
166 }
167 for (key, value) in &self.config.env {
168 shell = shell.env(key, value);
169 }
170 if self.config.clean_env {
171 shell = shell.clean_env();
172 }
173
174 let (stdout, stderr, exit_code) = if self.config.exit_code_as_output && !is_dry_run() {
178 self.run_capturing().await?
179 } else {
180 let output = shell.run().await?;
181 (
182 output.stdout().to_string(),
183 output.stderr().to_string(),
184 output.exit_code(),
185 )
186 };
187 let duration_ms = start.elapsed().as_millis() as u64;
188
189 info!(
190 step_kind = "shell",
191 command = %self.display(),
192 exit_code,
193 duration_ms,
194 "shell step completed"
195 );
196
197 self.record_metrics(duration_ms);
198
199 Ok(build_output(&stdout, &stderr, exit_code, duration_ms))
200 }
201
202 async fn run_capturing(&self) -> Result<(String, String, i32), EngineError> {
205 let mut cmd = self.command();
206 cmd.stdout(Stdio::piped())
207 .stderr(Stdio::piped())
208 .kill_on_drop(true);
209
210 if self.config.clean_env {
211 cmd.env_clear();
212 }
213 if let Some(ref dir) = self.config.dir {
214 cmd.current_dir(dir);
215 }
216 for (key, value) in &self.config.env {
217 cmd.env(key, value);
218 }
219
220 let child = cmd.spawn().map_err(|e| {
221 EngineError::Operation(OperationError::Shell {
222 exit_code: -1,
223 stderr: format!("failed to spawn shell: {e}"),
224 })
225 })?;
226
227 let timeout_dur = self
228 .config
229 .timeout_secs
230 .map(Duration::from_secs)
231 .unwrap_or(DEFAULT_SHELL_TIMEOUT);
232
233 let output = match tokio::time::timeout(timeout_dur, child.wait_with_output()).await {
234 Ok(Ok(output)) => output,
235 Ok(Err(e)) => {
236 return Err(EngineError::Operation(OperationError::Shell {
237 exit_code: -1,
238 stderr: format!("failed to wait for shell: {e}"),
239 }));
240 }
241 Err(_) => {
242 return Err(EngineError::Operation(OperationError::Timeout {
243 step: self.display(),
244 limit: timeout_dur,
245 }));
246 }
247 };
248
249 let exit_code = exit_code_of(output.status);
250 let stdout = truncate_output(&output.stdout, "shell stdout");
251 let stderr = truncate_output(&output.stderr, "shell stderr");
252 Ok((stdout, stderr, exit_code))
253 }
254
255 async fn execute_streaming(&self, sender: StepLogSender) -> Result<StepOutput, EngineError> {
258 let start = Instant::now();
259
260 let mut cmd = self.command();
261 cmd.stdout(Stdio::piped())
262 .stderr(Stdio::piped())
263 .kill_on_drop(true);
264
265 if self.config.clean_env {
266 cmd.env_clear();
267 }
268 if let Some(ref dir) = self.config.dir {
269 cmd.current_dir(dir);
270 }
271 for (key, value) in &self.config.env {
272 cmd.env(key, value);
273 }
274
275 let mut child = cmd.spawn().map_err(|e| {
276 EngineError::Operation(OperationError::Shell {
277 exit_code: -1,
278 stderr: format!("failed to spawn shell: {e}"),
279 })
280 })?;
281
282 let stdout_pipe = child.stdout.take().expect("stdout piped");
283 let stderr_pipe = child.stderr.take().expect("stderr piped");
284
285 let stdout_task = spawn(read_and_stream(
286 stdout_pipe,
287 sender.clone(),
288 LogStream::Stdout,
289 ));
290 let stderr_task = spawn(read_and_stream(stderr_pipe, sender, LogStream::Stderr));
291
292 let timeout_dur = self
293 .config
294 .timeout_secs
295 .map(Duration::from_secs)
296 .unwrap_or(DEFAULT_SHELL_TIMEOUT);
297
298 let status = match tokio::time::timeout(timeout_dur, child.wait()).await {
299 Ok(Ok(status)) => status,
300 Ok(Err(e)) => {
301 return Err(EngineError::Operation(OperationError::Shell {
302 exit_code: -1,
303 stderr: format!("failed to wait for shell: {e}"),
304 }));
305 }
306 Err(_) => {
307 child.kill().await.ok();
308 return Err(EngineError::Operation(OperationError::Timeout {
309 step: self.display(),
310 limit: timeout_dur,
311 }));
312 }
313 };
314
315 let raw_stdout = stdout_task.await.unwrap_or_default();
316 let raw_stderr = stderr_task.await.unwrap_or_default();
317
318 let stdout = truncate_output(raw_stdout.as_bytes(), "shell stdout");
319 let stderr = truncate_output(raw_stderr.as_bytes(), "shell stderr");
320
321 let exit_code = exit_code_of(status);
322 let duration_ms = start.elapsed().as_millis() as u64;
323
324 info!(
325 step_kind = "shell",
326 command = %self.display(),
327 exit_code,
328 duration_ms,
329 streaming = true,
330 "shell step completed"
331 );
332
333 self.record_metrics(duration_ms);
334
335 if exit_code != 0 && !self.config.exit_code_as_output {
336 return Err(EngineError::Operation(OperationError::Shell {
337 exit_code,
338 stderr: stderr.clone(),
339 }));
340 }
341
342 Ok(build_output(&stdout, &stderr, exit_code, duration_ms))
343 }
344
345 #[allow(unused_variables)]
346 fn record_metrics(&self, duration_ms: u64) {
347 #[cfg(feature = "prometheus")]
348 {
349 use ironflow_core::metric_names::{
350 SHELL_DURATION_SECONDS, SHELL_TOTAL, STATUS_SUCCESS,
351 };
352 use metrics::{counter, histogram};
353 counter!(SHELL_TOTAL, "status" => STATUS_SUCCESS).increment(1);
354 histogram!(SHELL_DURATION_SECONDS).record(duration_ms as f64 / 1000.0);
355 }
356 }
357}
358
359#[cfg(test)]
360mod tests {
361 use super::*;
362 use std::path::Path;
363
364 use ironflow_core::providers::claude::ClaudeCodeProvider;
365 use ironflow_core::providers::record_replay::RecordReplayProvider;
366 use tempfile::tempdir;
367 use uuid::Uuid;
368
369 fn create_test_provider() -> Arc<dyn AgentProvider> {
370 let inner = ClaudeCodeProvider::new();
371 Arc::new(RecordReplayProvider::replay(
372 inner,
373 "/tmp/ironflow-fixtures",
374 ))
375 }
376
377 #[tokio::test]
378 async fn shell_simple_command() {
379 let config = ShellConfig::new("echo hello");
380 let executor = ShellExecutor::new(&config);
381 let provider = create_test_provider();
382
383 let result = executor.execute(&provider).await;
384 assert!(result.is_ok());
385 let output = result.unwrap();
386 assert_eq!(output.output["exit_code"].as_i64().unwrap(), 0);
387 assert!(output.output["stdout"].as_str().unwrap().contains("hello"));
388 }
389
390 #[tokio::test]
391 async fn shell_nonzero_exit_returns_error() {
392 let config = ShellConfig::new("exit 1");
393 let executor = ShellExecutor::new(&config);
394 let provider = create_test_provider();
395
396 let result = executor.execute(&provider).await;
397 assert!(result.is_err());
398 }
399
400 #[tokio::test]
401 async fn shell_env_variables() {
402 let config = ShellConfig::new("echo $MY_VAR").env("MY_VAR", "test_value");
403 let executor = ShellExecutor::new(&config);
404 let provider = create_test_provider();
405
406 let result = executor.execute(&provider).await;
407 assert!(result.is_ok());
408 let output = result.unwrap();
409 assert!(
410 output.output["stdout"]
411 .as_str()
412 .unwrap()
413 .contains("test_value")
414 );
415 }
416
417 #[tokio::test]
418 async fn shell_step_output_has_structure() {
419 let config = ShellConfig::new("echo test");
420 let executor = ShellExecutor::new(&config);
421 let provider = create_test_provider();
422
423 let output = executor.execute(&provider).await.unwrap();
424 assert!(output.output.get("stdout").is_some());
425 assert!(output.output.get("stderr").is_some());
426 assert!(output.output.get("exit_code").is_some());
427 assert_eq!(output.cost_usd, Decimal::ZERO);
428 assert!(output.duration_ms < 5000);
429 }
430
431 #[tokio::test]
432 async fn shell_command_with_pipe() {
433 let config = ShellConfig::new("echo hello | grep hello");
434 let executor = ShellExecutor::new(&config);
435 let provider = create_test_provider();
436
437 let result = executor.execute(&provider).await;
438 assert!(result.is_ok());
439 let output = result.unwrap();
440 assert_eq!(output.output["exit_code"].as_i64().unwrap(), 0);
441 assert!(output.output["stdout"].as_str().unwrap().contains("hello"));
442 }
443
444 #[tokio::test]
445 async fn shell_streaming_emits_lines() {
446 let config = ShellConfig::new("echo line1 && echo line2");
447 let (sender, mut receiver) = crate::log_sender::channel();
448 let step_sender = StepLogSender::new(
449 sender,
450 uuid::Uuid::now_v7(),
451 uuid::Uuid::now_v7(),
452 "test".to_string(),
453 );
454 let executor = ShellExecutor::new(&config).with_log_sender(step_sender);
455 let provider = create_test_provider();
456
457 let result = executor.execute(&provider).await;
458 assert!(result.is_ok());
459
460 let output = result.unwrap();
461 assert!(output.output["stdout"].as_str().unwrap().contains("line1"));
462 assert!(output.output["stdout"].as_str().unwrap().contains("line2"));
463
464 let mut lines = Vec::new();
465 while let Ok(line) = receiver.try_recv() {
466 lines.push(line);
467 }
468 assert!(lines.len() >= 2);
469 assert_eq!(lines[0].stream, LogStream::Stdout);
470 assert_eq!(lines[0].line, "line1");
471 assert_eq!(lines[1].line, "line2");
472 }
473
474 #[tokio::test]
475 async fn shell_streaming_captures_stderr() {
476 let config = ShellConfig::new("echo err >&2");
477 let (sender, mut receiver) = crate::log_sender::channel();
478 let step_sender = StepLogSender::new(
479 sender,
480 uuid::Uuid::now_v7(),
481 uuid::Uuid::now_v7(),
482 "test".to_string(),
483 );
484 let executor = ShellExecutor::new(&config).with_log_sender(step_sender);
485 let provider = create_test_provider();
486
487 let result = executor.execute(&provider).await;
488 assert!(result.is_ok());
489
490 let mut stderr_lines = Vec::new();
491 while let Ok(line) = receiver.try_recv() {
492 if line.stream == LogStream::Stderr {
493 stderr_lines.push(line);
494 }
495 }
496 assert!(!stderr_lines.is_empty());
497 assert_eq!(stderr_lines[0].line, "err");
498 }
499
500 fn streaming_sender() -> StepLogSender {
501 let (sender, _receiver) = crate::log_sender::channel();
502 StepLogSender::new(
503 sender,
504 uuid::Uuid::now_v7(),
505 uuid::Uuid::now_v7(),
506 "test".to_string(),
507 )
508 }
509
510 #[tokio::test]
511 async fn shell_exit_code_as_output_buffered_keeps_code_and_streams() {
512 let config = ShellConfig::new("echo out; echo err >&2; exit 3").exit_code_as_output();
513 let executor = ShellExecutor::new(&config);
514 let provider = create_test_provider();
515
516 let output = executor
517 .execute(&provider)
518 .await
519 .expect("a non-zero exit is an output");
520 assert_eq!(output.output["exit_code"], 3);
521 assert_eq!(output.stdout().trim(), "out");
522 assert_eq!(output.stderr().trim(), "err");
523 assert!(!output.is_success());
524 assert_eq!(output.exit_code(), Some(3));
525 }
526
527 #[tokio::test]
528 async fn shell_exit_code_as_output_streaming_keeps_code() {
529 let config = ShellConfig::new("echo out; echo err >&2; exit 3").exit_code_as_output();
530 let executor = ShellExecutor::new(&config).with_log_sender(streaming_sender());
531 let provider = create_test_provider();
532
533 let output = executor
534 .execute(&provider)
535 .await
536 .expect("a non-zero exit is an output");
537 assert_eq!(output.output["exit_code"], 3);
538 assert_eq!(output.stdout(), "out");
539 assert_eq!(output.stderr(), "err");
540 assert!(!output.is_success());
541 assert_eq!(output.exit_code(), Some(3));
542 }
543
544 #[tokio::test]
545 async fn shell_exit_code_as_output_zero_exit_is_success() {
546 let config = ShellConfig::new("echo fine").exit_code_as_output();
547 let provider = create_test_provider();
548
549 let buffered = ShellExecutor::new(&config)
550 .execute(&provider)
551 .await
552 .expect("exit 0 succeeds");
553 assert!(buffered.is_success());
554 assert_eq!(buffered.exit_code(), Some(0));
555
556 let streaming = ShellExecutor::new(&config)
557 .with_log_sender(streaming_sender())
558 .execute(&provider)
559 .await
560 .expect("exit 0 succeeds");
561 assert!(streaming.is_success());
562 }
563
564 #[tokio::test]
565 async fn shell_exit_code_as_output_still_errors_on_timeout() {
566 let config = ShellConfig::new("sleep 5")
567 .timeout_secs(1)
568 .exit_code_as_output();
569 let provider = create_test_provider();
570
571 let buffered = ShellExecutor::new(&config).execute(&provider).await;
572 assert!(matches!(
573 buffered,
574 Err(EngineError::Operation(OperationError::Timeout { .. }))
575 ));
576
577 let streaming = ShellExecutor::new(&config)
578 .with_log_sender(streaming_sender())
579 .execute(&provider)
580 .await;
581 assert!(matches!(
582 streaming,
583 Err(EngineError::Operation(OperationError::Timeout { .. }))
584 ));
585 }
586
587 #[tokio::test]
588 async fn shell_without_exit_code_as_output_nonzero_still_errors() {
589 let config = ShellConfig::new("echo out; exit 3");
590 let provider = create_test_provider();
591
592 let buffered = ShellExecutor::new(&config).execute(&provider).await;
593 assert!(matches!(
594 buffered,
595 Err(EngineError::Operation(OperationError::Shell {
596 exit_code: 3,
597 ..
598 }))
599 ));
600
601 let streaming = ShellExecutor::new(&config)
602 .with_log_sender(streaming_sender())
603 .execute(&provider)
604 .await;
605 assert!(matches!(
606 streaming,
607 Err(EngineError::Operation(OperationError::Shell {
608 exit_code: 3,
609 ..
610 }))
611 ));
612 }
613
614 const HOSTILE_ARG: &str = "x'; echo injected; echo '$(echo sub) `echo tick` ${HOME} \\n";
616
617 #[tokio::test]
618 async fn exec_passes_each_argument_verbatim_buffered() {
619 let config = ShellConfig::exec("printf", &["%s|%s", HOSTILE_ARG, "two words"]);
620 let provider = create_test_provider();
621
622 let output = ShellExecutor::new(&config)
623 .execute(&provider)
624 .await
625 .expect("printf runs");
626
627 assert_eq!(output.stdout(), format!("{HOSTILE_ARG}|two words"));
628 assert_eq!(output.exit_code(), Some(0));
629 }
630
631 #[tokio::test]
632 async fn exec_passes_each_argument_verbatim_streaming() {
633 let config = ShellConfig::exec("printf", &["%s\n", HOSTILE_ARG]);
634 let (sender, mut receiver) = crate::log_sender::channel();
635 let step_sender =
636 StepLogSender::new(sender, Uuid::now_v7(), Uuid::now_v7(), "test".to_string());
637 let provider = create_test_provider();
638
639 let output = ShellExecutor::new(&config)
640 .with_log_sender(step_sender)
641 .execute(&provider)
642 .await
643 .expect("printf runs");
644
645 assert_eq!(output.stdout(), HOSTILE_ARG);
646 let line = receiver.try_recv().expect("one streamed line");
647 assert_eq!(line.line, HOSTILE_ARG);
648 }
649
650 #[tokio::test]
651 async fn exec_with_exit_code_as_output_keeps_the_code_and_args() {
652 let config =
653 ShellConfig::exec("sh", &["-c", "printf %s \"$1\"; exit 3", "sh", HOSTILE_ARG])
654 .exit_code_as_output();
655 let provider = create_test_provider();
656
657 let buffered = ShellExecutor::new(&config)
658 .execute(&provider)
659 .await
660 .expect("a non-zero exit is an output");
661 assert_eq!(buffered.exit_code(), Some(3));
662 assert_eq!(buffered.stdout(), HOSTILE_ARG);
663
664 let streaming = ShellExecutor::new(&config)
665 .with_log_sender(streaming_sender())
666 .execute(&provider)
667 .await
668 .expect("a non-zero exit is an output");
669 assert_eq!(streaming.exit_code(), Some(3));
670 assert_eq!(streaming.stdout(), HOSTILE_ARG);
671 }
672
673 #[tokio::test]
674 async fn exec_nonzero_exit_returns_error() {
675 let config = ShellConfig::exec("false", &[]);
676 let provider = create_test_provider();
677
678 let buffered = ShellExecutor::new(&config).execute(&provider).await;
679 assert!(matches!(
680 buffered,
681 Err(EngineError::Operation(OperationError::Shell {
682 exit_code: 1,
683 ..
684 }))
685 ));
686
687 let streaming = ShellExecutor::new(&config)
688 .with_log_sender(streaming_sender())
689 .execute(&provider)
690 .await;
691 assert!(matches!(
692 streaming,
693 Err(EngineError::Operation(OperationError::Shell {
694 exit_code: 1,
695 ..
696 }))
697 ));
698 }
699
700 #[tokio::test]
701 async fn exec_of_a_missing_program_is_a_spawn_error() {
702 let config = ShellConfig::exec("ironflow-no-such-program", &["arg"]);
703 let provider = create_test_provider();
704
705 for result in [
706 ShellExecutor::new(&config).execute(&provider).await,
707 ShellExecutor::new(&config)
708 .with_log_sender(streaming_sender())
709 .execute(&provider)
710 .await,
711 ShellExecutor::new(&config.clone().exit_code_as_output())
712 .execute(&provider)
713 .await,
714 ] {
715 match result {
716 Err(EngineError::Operation(OperationError::Shell { exit_code, stderr })) => {
717 assert_eq!(exit_code, -1);
718 assert!(stderr.starts_with("failed to spawn"), "stderr: {stderr}");
719 }
720 other => panic!("expected a spawn error, got {other:?}"),
721 }
722 }
723 }
724
725 #[tokio::test]
726 async fn exec_honours_dir_and_env() {
727 let dir = tempdir().expect("temp dir");
728 let config = ShellConfig::exec("sh", &["-c", "pwd; printf %s \"$GREETING\""])
729 .dir(dir.path().to_str().expect("utf-8 path"))
730 .env("GREETING", "hi; echo no");
731 let provider = create_test_provider();
732
733 let output = ShellExecutor::new(&config)
734 .execute(&provider)
735 .await
736 .expect("sh runs");
737
738 let canonical = dir.path().canonicalize().expect("canonical dir");
739 let mut lines = output.stdout().lines();
740 assert_eq!(
741 Path::new(lines.next().expect("pwd line"))
742 .canonicalize()
743 .expect("canonical pwd"),
744 canonical
745 );
746 assert_eq!(lines.next(), Some("hi; echo no"));
747 }
748
749 #[tokio::test]
750 async fn shell_streaming_nonzero_exit_returns_error() {
751 let config = ShellConfig::new("exit 42");
752 let (sender, _receiver) = crate::log_sender::channel();
753 let step_sender = StepLogSender::new(
754 sender,
755 uuid::Uuid::now_v7(),
756 uuid::Uuid::now_v7(),
757 "test".to_string(),
758 );
759 let executor = ShellExecutor::new(&config).with_log_sender(step_sender);
760 let provider = create_test_provider();
761
762 let result = executor.execute(&provider).await;
763 assert!(result.is_err());
764 }
765}