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 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 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 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 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 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}