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