Skip to main content

ironflow_engine/executor/
shell.rs

1//! Shell step executor.
2
3use 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
31/// Read lines from an async reader, emit each line to the sender, and
32/// accumulate the full output as a single `String`.
33async 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
50/// Exit code of a finished process: the real code, `-signal` when a signal
51/// killed it, `-1` otherwise.
52fn exit_code_of(status: ExitStatus) -> i32 {
53    status
54        .code()
55        .or_else(|| status.signal().map(|signal| -signal))
56        .unwrap_or(-1)
57}
58
59/// Build the output JSON shared by the buffered and streaming paths.
60fn 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
79/// Executor for shell steps.
80///
81/// Runs a shell command and captures stdout, stderr, and exit code.
82/// When a [`StepLogSender`] is attached, stdout and stderr are streamed
83/// line-by-line in real time.
84pub struct ShellExecutor<'a> {
85    config: &'a ShellConfig,
86    log_sender: Option<StepLogSender>,
87}
88
89impl<'a> ShellExecutor<'a> {
90    /// Create a new shell executor from a config reference.
91    pub fn new(config: &'a ShellConfig) -> Self {
92        Self {
93            config,
94            log_sender: None,
95        }
96    }
97
98    /// Attach a log sender for real-time line streaming.
99    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    /// Non-streaming execution via [`Shell::run()`].
120    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        // `Shell::run` drops stdout on a non-zero exit, so the option needs its
138        // own capture. A global dry run keeps going through `Shell::run`, which
139        // short-circuits without spawning anything.
140        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    /// Run the command to completion and keep stdout, stderr and the exit
166    /// code whatever the code is. Used when `exit_code_as_output` is set.
167    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    /// Streaming execution: reads stdout/stderr line-by-line and forwards
220    /// each line to the [`StepLogSender`] in real time.
221    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}