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        account_id: None,
77    }
78}
79
80/// Executor for shell steps.
81///
82/// Runs a shell command and captures stdout, stderr, and exit code.
83/// When a [`StepLogSender`] is attached, stdout and stderr are streamed
84/// line-by-line in real time.
85pub struct ShellExecutor<'a> {
86    config: &'a ShellConfig,
87    log_sender: Option<StepLogSender>,
88}
89
90impl<'a> ShellExecutor<'a> {
91    /// Create a new shell executor from a config reference.
92    pub fn new(config: &'a ShellConfig) -> Self {
93        Self {
94            config,
95            log_sender: None,
96        }
97    }
98
99    /// Attach a log sender for real-time line streaming.
100    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    /// Non-streaming execution via [`Shell::run()`].
121    async fn execute_buffered(&self) -> Result<StepOutput, EngineError> {
122        let start = Instant::now();
123
124        let mut shell = Shell::new(&self.config.command);
125        if let Some(secs) = self.config.timeout_secs {
126            shell = shell.timeout(Duration::from_secs(secs));
127        }
128        if let Some(ref dir) = self.config.dir {
129            shell = shell.dir(dir);
130        }
131        for (key, value) in &self.config.env {
132            shell = shell.env(key, value);
133        }
134        if self.config.clean_env {
135            shell = shell.clean_env();
136        }
137
138        // `Shell::run` drops stdout on a non-zero exit, so the option needs its
139        // own capture. A global dry run keeps going through `Shell::run`, which
140        // short-circuits without spawning anything.
141        let (stdout, stderr, exit_code) = if self.config.exit_code_as_output && !is_dry_run() {
142            self.run_capturing().await?
143        } else {
144            let output = shell.run().await?;
145            (
146                output.stdout().to_string(),
147                output.stderr().to_string(),
148                output.exit_code(),
149            )
150        };
151        let duration_ms = start.elapsed().as_millis() as u64;
152
153        info!(
154            step_kind = "shell",
155            command = %self.config.command,
156            exit_code,
157            duration_ms,
158            "shell step completed"
159        );
160
161        self.record_metrics(duration_ms);
162
163        Ok(build_output(&stdout, &stderr, exit_code, duration_ms))
164    }
165
166    /// Run the command to completion and keep stdout, stderr and the exit
167    /// code whatever the code is. Used when `exit_code_as_output` is set.
168    async fn run_capturing(&self) -> Result<(String, String, i32), EngineError> {
169        let mut cmd = Command::new("sh");
170        cmd.arg("-c").arg(&self.config.command);
171        cmd.stdout(Stdio::piped())
172            .stderr(Stdio::piped())
173            .kill_on_drop(true);
174
175        if self.config.clean_env {
176            cmd.env_clear();
177        }
178        if let Some(ref dir) = self.config.dir {
179            cmd.current_dir(dir);
180        }
181        for (key, value) in &self.config.env {
182            cmd.env(key, value);
183        }
184
185        let child = cmd.spawn().map_err(|e| {
186            EngineError::Operation(OperationError::Shell {
187                exit_code: -1,
188                stderr: format!("failed to spawn shell: {e}"),
189            })
190        })?;
191
192        let timeout_dur = self
193            .config
194            .timeout_secs
195            .map(Duration::from_secs)
196            .unwrap_or(DEFAULT_SHELL_TIMEOUT);
197
198        let output = match tokio::time::timeout(timeout_dur, child.wait_with_output()).await {
199            Ok(Ok(output)) => output,
200            Ok(Err(e)) => {
201                return Err(EngineError::Operation(OperationError::Shell {
202                    exit_code: -1,
203                    stderr: format!("failed to wait for shell: {e}"),
204                }));
205            }
206            Err(_) => {
207                return Err(EngineError::Operation(OperationError::Timeout {
208                    step: self.config.command.clone(),
209                    limit: timeout_dur,
210                }));
211            }
212        };
213
214        let exit_code = exit_code_of(output.status);
215        let stdout = truncate_output(&output.stdout, "shell stdout");
216        let stderr = truncate_output(&output.stderr, "shell stderr");
217        Ok((stdout, stderr, exit_code))
218    }
219
220    /// Streaming execution: reads stdout/stderr line-by-line and forwards
221    /// each line to the [`StepLogSender`] in real time.
222    async fn execute_streaming(&self, sender: StepLogSender) -> Result<StepOutput, EngineError> {
223        let start = Instant::now();
224
225        let mut cmd = Command::new("sh");
226        cmd.arg("-c").arg(&self.config.command);
227        cmd.stdout(Stdio::piped())
228            .stderr(Stdio::piped())
229            .kill_on_drop(true);
230
231        if self.config.clean_env {
232            cmd.env_clear();
233        }
234        if let Some(ref dir) = self.config.dir {
235            cmd.current_dir(dir);
236        }
237        for (key, value) in &self.config.env {
238            cmd.env(key, value);
239        }
240
241        let mut child = cmd.spawn().map_err(|e| {
242            EngineError::Operation(OperationError::Shell {
243                exit_code: -1,
244                stderr: format!("failed to spawn shell: {e}"),
245            })
246        })?;
247
248        let stdout_pipe = child.stdout.take().expect("stdout piped");
249        let stderr_pipe = child.stderr.take().expect("stderr piped");
250
251        let stdout_task = spawn(read_and_stream(
252            stdout_pipe,
253            sender.clone(),
254            LogStream::Stdout,
255        ));
256        let stderr_task = spawn(read_and_stream(stderr_pipe, sender, LogStream::Stderr));
257
258        let timeout_dur = self
259            .config
260            .timeout_secs
261            .map(Duration::from_secs)
262            .unwrap_or(DEFAULT_SHELL_TIMEOUT);
263
264        let status = match tokio::time::timeout(timeout_dur, child.wait()).await {
265            Ok(Ok(status)) => status,
266            Ok(Err(e)) => {
267                return Err(EngineError::Operation(OperationError::Shell {
268                    exit_code: -1,
269                    stderr: format!("failed to wait for shell: {e}"),
270                }));
271            }
272            Err(_) => {
273                child.kill().await.ok();
274                return Err(EngineError::Operation(OperationError::Timeout {
275                    step: self.config.command.clone(),
276                    limit: timeout_dur,
277                }));
278            }
279        };
280
281        let raw_stdout = stdout_task.await.unwrap_or_default();
282        let raw_stderr = stderr_task.await.unwrap_or_default();
283
284        let stdout = truncate_output(raw_stdout.as_bytes(), "shell stdout");
285        let stderr = truncate_output(raw_stderr.as_bytes(), "shell stderr");
286
287        let exit_code = exit_code_of(status);
288        let duration_ms = start.elapsed().as_millis() as u64;
289
290        info!(
291            step_kind = "shell",
292            command = %self.config.command,
293            exit_code,
294            duration_ms,
295            streaming = true,
296            "shell step completed"
297        );
298
299        self.record_metrics(duration_ms);
300
301        if exit_code != 0 && !self.config.exit_code_as_output {
302            return Err(EngineError::Operation(OperationError::Shell {
303                exit_code,
304                stderr: stderr.clone(),
305            }));
306        }
307
308        Ok(build_output(&stdout, &stderr, exit_code, duration_ms))
309    }
310
311    #[allow(unused_variables)]
312    fn record_metrics(&self, duration_ms: u64) {
313        #[cfg(feature = "prometheus")]
314        {
315            use ironflow_core::metric_names::{
316                SHELL_DURATION_SECONDS, SHELL_TOTAL, STATUS_SUCCESS,
317            };
318            use metrics::{counter, histogram};
319            counter!(SHELL_TOTAL, "status" => STATUS_SUCCESS).increment(1);
320            histogram!(SHELL_DURATION_SECONDS).record(duration_ms as f64 / 1000.0);
321        }
322    }
323}
324
325#[cfg(test)]
326mod tests {
327    use super::*;
328    use ironflow_core::providers::claude::ClaudeCodeProvider;
329    use ironflow_core::providers::record_replay::RecordReplayProvider;
330
331    fn create_test_provider() -> Arc<dyn AgentProvider> {
332        let inner = ClaudeCodeProvider::new();
333        Arc::new(RecordReplayProvider::replay(
334            inner,
335            "/tmp/ironflow-fixtures",
336        ))
337    }
338
339    #[tokio::test]
340    async fn shell_simple_command() {
341        let config = ShellConfig::new("echo hello");
342        let executor = ShellExecutor::new(&config);
343        let provider = create_test_provider();
344
345        let result = executor.execute(&provider).await;
346        assert!(result.is_ok());
347        let output = result.unwrap();
348        assert_eq!(output.output["exit_code"].as_i64().unwrap(), 0);
349        assert!(output.output["stdout"].as_str().unwrap().contains("hello"));
350    }
351
352    #[tokio::test]
353    async fn shell_nonzero_exit_returns_error() {
354        let config = ShellConfig::new("exit 1");
355        let executor = ShellExecutor::new(&config);
356        let provider = create_test_provider();
357
358        let result = executor.execute(&provider).await;
359        assert!(result.is_err());
360    }
361
362    #[tokio::test]
363    async fn shell_env_variables() {
364        let config = ShellConfig::new("echo $MY_VAR").env("MY_VAR", "test_value");
365        let executor = ShellExecutor::new(&config);
366        let provider = create_test_provider();
367
368        let result = executor.execute(&provider).await;
369        assert!(result.is_ok());
370        let output = result.unwrap();
371        assert!(
372            output.output["stdout"]
373                .as_str()
374                .unwrap()
375                .contains("test_value")
376        );
377    }
378
379    #[tokio::test]
380    async fn shell_step_output_has_structure() {
381        let config = ShellConfig::new("echo test");
382        let executor = ShellExecutor::new(&config);
383        let provider = create_test_provider();
384
385        let output = executor.execute(&provider).await.unwrap();
386        assert!(output.output.get("stdout").is_some());
387        assert!(output.output.get("stderr").is_some());
388        assert!(output.output.get("exit_code").is_some());
389        assert_eq!(output.cost_usd, Decimal::ZERO);
390        assert!(output.duration_ms < 5000);
391    }
392
393    #[tokio::test]
394    async fn shell_command_with_pipe() {
395        let config = ShellConfig::new("echo hello | grep hello");
396        let executor = ShellExecutor::new(&config);
397        let provider = create_test_provider();
398
399        let result = executor.execute(&provider).await;
400        assert!(result.is_ok());
401        let output = result.unwrap();
402        assert_eq!(output.output["exit_code"].as_i64().unwrap(), 0);
403        assert!(output.output["stdout"].as_str().unwrap().contains("hello"));
404    }
405
406    #[tokio::test]
407    async fn shell_streaming_emits_lines() {
408        let config = ShellConfig::new("echo line1 && echo line2");
409        let (sender, mut receiver) = crate::log_sender::channel();
410        let step_sender = StepLogSender::new(
411            sender,
412            uuid::Uuid::now_v7(),
413            uuid::Uuid::now_v7(),
414            "test".to_string(),
415        );
416        let executor = ShellExecutor::new(&config).with_log_sender(step_sender);
417        let provider = create_test_provider();
418
419        let result = executor.execute(&provider).await;
420        assert!(result.is_ok());
421
422        let output = result.unwrap();
423        assert!(output.output["stdout"].as_str().unwrap().contains("line1"));
424        assert!(output.output["stdout"].as_str().unwrap().contains("line2"));
425
426        let mut lines = Vec::new();
427        while let Ok(line) = receiver.try_recv() {
428            lines.push(line);
429        }
430        assert!(lines.len() >= 2);
431        assert_eq!(lines[0].stream, LogStream::Stdout);
432        assert_eq!(lines[0].line, "line1");
433        assert_eq!(lines[1].line, "line2");
434    }
435
436    #[tokio::test]
437    async fn shell_streaming_captures_stderr() {
438        let config = ShellConfig::new("echo err >&2");
439        let (sender, mut receiver) = crate::log_sender::channel();
440        let step_sender = StepLogSender::new(
441            sender,
442            uuid::Uuid::now_v7(),
443            uuid::Uuid::now_v7(),
444            "test".to_string(),
445        );
446        let executor = ShellExecutor::new(&config).with_log_sender(step_sender);
447        let provider = create_test_provider();
448
449        let result = executor.execute(&provider).await;
450        assert!(result.is_ok());
451
452        let mut stderr_lines = Vec::new();
453        while let Ok(line) = receiver.try_recv() {
454            if line.stream == LogStream::Stderr {
455                stderr_lines.push(line);
456            }
457        }
458        assert!(!stderr_lines.is_empty());
459        assert_eq!(stderr_lines[0].line, "err");
460    }
461
462    fn streaming_sender() -> StepLogSender {
463        let (sender, _receiver) = crate::log_sender::channel();
464        StepLogSender::new(
465            sender,
466            uuid::Uuid::now_v7(),
467            uuid::Uuid::now_v7(),
468            "test".to_string(),
469        )
470    }
471
472    #[tokio::test]
473    async fn shell_exit_code_as_output_buffered_keeps_code_and_streams() {
474        let config = ShellConfig::new("echo out; echo err >&2; exit 3").exit_code_as_output();
475        let executor = ShellExecutor::new(&config);
476        let provider = create_test_provider();
477
478        let output = executor
479            .execute(&provider)
480            .await
481            .expect("a non-zero exit is an output");
482        assert_eq!(output.output["exit_code"], 3);
483        assert_eq!(output.stdout().trim(), "out");
484        assert_eq!(output.stderr().trim(), "err");
485        assert!(!output.is_success());
486        assert_eq!(output.exit_code(), Some(3));
487    }
488
489    #[tokio::test]
490    async fn shell_exit_code_as_output_streaming_keeps_code() {
491        let config = ShellConfig::new("echo out; echo err >&2; exit 3").exit_code_as_output();
492        let executor = ShellExecutor::new(&config).with_log_sender(streaming_sender());
493        let provider = create_test_provider();
494
495        let output = executor
496            .execute(&provider)
497            .await
498            .expect("a non-zero exit is an output");
499        assert_eq!(output.output["exit_code"], 3);
500        assert_eq!(output.stdout(), "out");
501        assert_eq!(output.stderr(), "err");
502        assert!(!output.is_success());
503        assert_eq!(output.exit_code(), Some(3));
504    }
505
506    #[tokio::test]
507    async fn shell_exit_code_as_output_zero_exit_is_success() {
508        let config = ShellConfig::new("echo fine").exit_code_as_output();
509        let provider = create_test_provider();
510
511        let buffered = ShellExecutor::new(&config)
512            .execute(&provider)
513            .await
514            .expect("exit 0 succeeds");
515        assert!(buffered.is_success());
516        assert_eq!(buffered.exit_code(), Some(0));
517
518        let streaming = ShellExecutor::new(&config)
519            .with_log_sender(streaming_sender())
520            .execute(&provider)
521            .await
522            .expect("exit 0 succeeds");
523        assert!(streaming.is_success());
524    }
525
526    #[tokio::test]
527    async fn shell_exit_code_as_output_still_errors_on_timeout() {
528        let config = ShellConfig::new("sleep 5")
529            .timeout_secs(1)
530            .exit_code_as_output();
531        let provider = create_test_provider();
532
533        let buffered = ShellExecutor::new(&config).execute(&provider).await;
534        assert!(matches!(
535            buffered,
536            Err(EngineError::Operation(OperationError::Timeout { .. }))
537        ));
538
539        let streaming = ShellExecutor::new(&config)
540            .with_log_sender(streaming_sender())
541            .execute(&provider)
542            .await;
543        assert!(matches!(
544            streaming,
545            Err(EngineError::Operation(OperationError::Timeout { .. }))
546        ));
547    }
548
549    #[tokio::test]
550    async fn shell_without_exit_code_as_output_nonzero_still_errors() {
551        let config = ShellConfig::new("echo out; exit 3");
552        let provider = create_test_provider();
553
554        let buffered = ShellExecutor::new(&config).execute(&provider).await;
555        assert!(matches!(
556            buffered,
557            Err(EngineError::Operation(OperationError::Shell {
558                exit_code: 3,
559                ..
560            }))
561        ));
562
563        let streaming = ShellExecutor::new(&config)
564            .with_log_sender(streaming_sender())
565            .execute(&provider)
566            .await;
567        assert!(matches!(
568            streaming,
569            Err(EngineError::Operation(OperationError::Shell {
570                exit_code: 3,
571                ..
572            }))
573        ));
574    }
575
576    #[tokio::test]
577    async fn shell_streaming_nonzero_exit_returns_error() {
578        let config = ShellConfig::new("exit 42");
579        let (sender, _receiver) = crate::log_sender::channel();
580        let step_sender = StepLogSender::new(
581            sender,
582            uuid::Uuid::now_v7(),
583            uuid::Uuid::now_v7(),
584            "test".to_string(),
585        );
586        let executor = ShellExecutor::new(&config).with_log_sender(step_sender);
587        let provider = create_test_provider();
588
589        let result = executor.execute(&provider).await;
590        assert!(result.is_err());
591    }
592}