1pub mod ansi;
66pub mod events;
67#[doc(hidden)]
68pub mod macros;
69pub mod pipeline;
70pub mod quote;
71pub mod state;
72pub mod stream;
73pub mod terminal;
74pub mod trace;
75
76pub mod commands;
78pub mod shell_parser;
79pub mod utils;
80
81use std::collections::HashMap;
82use std::path::PathBuf;
83use std::process::Stdio;
84use tokio::io::{AsyncRead, AsyncReadExt, AsyncWriteExt};
85use tokio::process::{Child, Command};
86use tokio::sync::mpsc;
87
88pub use commands::{CommandContext, StreamChunk};
89pub use shell_parser::{needs_real_shell, parse_shell_command, ParsedCommand};
90pub use utils::{CommandResult, VirtualUtils};
91
92pub use ansi::{AnsiConfig, AnsiUtils};
94pub use events::{EventData, EventType, StreamEmitter};
95pub use pipeline::{Pipeline, PipelineBuilder, PipelineExt};
96pub use quote::{
97 escape_for_double_quotes, escape_for_single_quotes, has_shell_escapes,
98 is_pre_quoted_passthrough_enabled, is_quote_context_enabled, quote, quote_for_context,
99 scan_quote_context, QuoteContext,
100};
101pub use state::{
102 get_shell_settings, global_state, reset_global_state, set_shell_option, unset_shell_option,
103 GlobalState, ShellSettings,
104};
105pub use stream::{AsyncIterator, IntoStream, OutputChunk, OutputStream, StreamingRunner};
106pub use trace::trace;
107
108#[derive(Clone, Copy)]
109enum ChildOutput {
110 Stdout,
111 Stderr,
112}
113
114async fn collect_child_output<R>(
120 reader: Option<R>,
121 mirror: bool,
122 target: ChildOutput,
123) -> std::io::Result<Vec<u8>>
124where
125 R: AsyncRead + Unpin,
126{
127 let Some(mut reader) = reader else {
128 return Ok(Vec::new());
129 };
130 let mut collected = Vec::new();
131 let mut buffer = [0_u8; 8192];
132
133 loop {
134 let count = reader.read(&mut buffer).await?;
135 if count == 0 {
136 break;
137 }
138
139 let chunk = &buffer[..count];
140 collected.extend_from_slice(chunk);
141 if mirror {
142 match target {
143 ChildOutput::Stdout => {
144 let mut output = std::io::stdout().lock();
145 let _ = std::io::Write::write_all(&mut output, chunk);
146 let _ = std::io::Write::flush(&mut output);
147 }
148 ChildOutput::Stderr => {
149 let mut output = std::io::stderr().lock();
150 let _ = std::io::Write::write_all(&mut output, chunk);
151 let _ = std::io::Write::flush(&mut output);
152 }
153 }
154 }
155 }
156
157 Ok(collected)
158}
159
160fn fallback_cwd() -> PathBuf {
161 std::env::var_os("HOME")
162 .or_else(|| std::env::var_os("USERPROFILE"))
163 .map(PathBuf::from)
164 .filter(|path| path.is_dir())
165 .unwrap_or_else(std::env::temp_dir)
166}
167
168fn resolve_spawn_cwd(cwd: Option<&PathBuf>) -> Option<PathBuf> {
180 if let Some(c) = cwd {
182 return Some(c.clone());
183 }
184
185 match std::env::current_dir() {
188 Ok(_) => None,
189 Err(e) => {
190 let fallback = fallback_cwd();
191 trace(
192 "ProcessRunner",
193 &format!(
194 "current_dir() failed ({}); spawning in fallback directory {}",
195 e,
196 fallback.display()
197 ),
198 );
199 Some(fallback)
200 }
201 }
202}
203
204#[derive(Debug, thiserror::Error)]
206pub enum Error {
207 #[error("IO error: {0}")]
208 Io(#[from] std::io::Error),
209
210 #[error("Command failed with exit code {code}: {message}")]
211 CommandFailed { code: i32, message: String },
212
213 #[error("Command not found: {0}")]
214 CommandNotFound(String),
215
216 #[error("Parse error: {0}")]
217 ParseError(String),
218
219 #[error("Cancelled")]
220 Cancelled,
221}
222
223impl Error {
224 pub fn command_failed(code: i32, message: impl Into<String>) -> Self {
226 Error::CommandFailed {
227 code,
228 message: message.into(),
229 }
230 }
231
232 pub fn code(&self) -> Option<i32> {
238 match self {
239 Error::CommandFailed { code, .. } => Some(*code),
240 Error::CommandNotFound(_) => Some(127),
243 Error::Io(error) => match error.kind() {
244 std::io::ErrorKind::NotFound => Some(127),
245 std::io::ErrorKind::PermissionDenied => Some(126),
246 _ => None,
247 },
248 Error::Cancelled => Some(130),
250 Error::ParseError(_) => None,
251 }
252 }
253
254 pub fn exit_code(&self) -> Option<i32> {
260 self.code()
261 }
262}
263
264pub type Result<T> = std::result::Result<T, Error>;
266
267#[derive(Debug, Clone)]
269pub struct RunOptions {
270 pub mirror: bool,
272 pub capture: bool,
274 pub stdin: StdinOption,
276 pub cwd: Option<PathBuf>,
278 pub env: Option<HashMap<String, String>>,
280 pub interactive: bool,
282 pub shell_operators: bool,
284 pub trace: bool,
286}
287
288impl Default for RunOptions {
289 fn default() -> Self {
290 RunOptions {
291 mirror: true,
292 capture: true,
293 stdin: StdinOption::Inherit,
294 cwd: None,
295 env: None,
296 interactive: false,
297 shell_operators: true,
298 trace: true,
299 }
300 }
301}
302
303#[derive(Debug, Clone)]
305pub enum StdinOption {
306 Inherit,
308 Pipe,
310 Content(String),
312 Null,
314}
315
316pub struct ProcessRunner {
318 command: String,
319 options: RunOptions,
320 child: Option<Child>,
321 result: Option<CommandResult>,
322 started: bool,
323 finished: bool,
324 cancelled: bool,
325 output_tx: Option<mpsc::Sender<StreamChunk>>,
326 #[allow(dead_code)]
331 output_rx: Option<mpsc::Receiver<StreamChunk>>,
332}
333
334impl ProcessRunner {
335 pub fn new(command: impl Into<String>, options: RunOptions) -> Self {
337 let (tx, rx) = mpsc::channel(1024);
338 ProcessRunner {
339 command: command.into(),
340 options,
341 child: None,
342 result: None,
343 started: false,
344 finished: false,
345 cancelled: false,
346 output_tx: Some(tx),
347 output_rx: Some(rx),
348 }
349 }
350
351 pub async fn start(&mut self) -> Result<()> {
353 if self.started {
354 return Ok(());
355 }
356 self.started = true;
357
358 utils::trace_lazy("ProcessRunner", || {
359 format!("Starting command: {}", self.command)
360 });
361
362 let first_word = if has_shell_escapes(&self.command) || needs_real_shell(&self.command) {
371 ""
372 } else {
373 self.command.split_whitespace().next().unwrap_or("")
374 };
375 if let Some(result) = self.try_virtual_command(first_word).await {
376 self.result = Some(result);
377 self.finished = true;
378 return Ok(());
379 }
380 let _parsed = if self.options.shell_operators && !needs_real_shell(&self.command) {
382 parse_shell_command(&self.command)
383 } else {
384 None
385 };
386
387 let shell = find_available_shell();
389
390 let mut cmd = Command::new(&shell.cmd);
391 for arg in &shell.args {
392 cmd.arg(arg);
393 }
394 utils::append_shell_command(&mut cmd, &self.command, self.options.env.as_ref());
395
396 match &self.options.stdin {
398 StdinOption::Inherit => {
399 cmd.stdin(Stdio::inherit());
400 }
401 StdinOption::Pipe => {
402 cmd.stdin(Stdio::piped());
403 }
404 StdinOption::Content(_) => {
405 cmd.stdin(Stdio::piped());
406 }
407 StdinOption::Null => {
408 cmd.stdin(Stdio::null());
409 }
410 }
411
412 if self.options.capture || self.options.mirror {
414 cmd.stdout(Stdio::piped());
415 cmd.stderr(Stdio::piped());
416 } else {
417 cmd.stdout(Stdio::inherit());
418 cmd.stderr(Stdio::inherit());
419 }
420
421 if let Some(cwd) = resolve_spawn_cwd(self.options.cwd.as_ref()) {
424 cmd.current_dir(cwd);
425 }
426
427 if let Some(ref env_vars) = self.options.env {
429 for (key, value) in env_vars {
430 cmd.env(key, value);
431 }
432 }
433
434 let child = cmd.spawn()?;
436 self.child = Some(child);
437
438 Ok(())
439 }
440
441 pub async fn run(&mut self) -> Result<CommandResult> {
443 self.start().await?;
444
445 if let Some(result) = &self.result {
446 return Ok(result.clone());
447 }
448
449 let mut child = self
450 .child
451 .take()
452 .ok_or_else(|| Error::Io(std::io::Error::other("Process not started")))?;
453
454 if let StdinOption::Content(ref content) = self.options.stdin {
456 if let Some(mut stdin) = child.stdin.take() {
457 let content = content.clone();
458 tokio::spawn(async move {
459 let _ = stdin.write_all(content.as_bytes()).await;
460 let _ = stdin.shutdown().await;
461 });
462 }
463 }
464
465 let stdout = child.stdout.take();
469 let stderr = child.stderr.take();
470 let collected = tokio::try_join!(
471 collect_child_output(stdout, self.options.mirror, ChildOutput::Stdout),
472 collect_child_output(stderr, self.options.mirror, ChildOutput::Stderr),
473 );
474 let (stdout, stderr) = match collected {
475 Ok(output) => output,
476 Err(error) => {
477 let _ = child.start_kill();
480 let _ = child.wait().await;
481 return Err(error.into());
482 }
483 };
484
485 let status = child.wait().await?;
486 let code = status.code().unwrap_or(-1);
487
488 let result = CommandResult {
489 stdout: String::from_utf8_lossy(&stdout).into_owned(),
490 stderr: String::from_utf8_lossy(&stderr).into_owned(),
491 code,
492 };
493
494 self.result = Some(result.clone());
495 self.finished = true;
496
497 Ok(result)
498 }
499
500 async fn try_virtual_command(&self, cmd_name: &str) -> Option<CommandResult> {
502 if !commands::are_virtual_commands_enabled() {
503 return None;
504 }
505
506 if cmd_name.is_empty() {
511 return None;
512 }
513
514 let words = shell_parser::split_command_words(&self.command);
518 let args: Vec<String> = words.into_iter().skip(1).collect();
519
520 let ctx = CommandContext {
521 args,
522 stdin: match &self.options.stdin {
523 StdinOption::Content(s) => Some(s.clone()),
524 _ => None,
525 },
526 cwd: self.options.cwd.clone(),
527 env: self.options.env.clone(),
528 output_tx: self.output_tx.clone(),
529 is_cancelled: None,
530 };
531
532 match cmd_name {
533 "echo" => Some(commands::echo(ctx).await),
534 "pwd" => Some(commands::pwd(ctx).await),
535 "cd" => Some(commands::cd::resolve_cd(ctx).await.0),
536 "true" => Some(commands::r#true(ctx).await),
537 "false" => Some(commands::r#false(ctx).await),
538 "sleep" => Some(commands::sleep(ctx).await),
539 "cat" => Some(commands::cat(ctx).await),
540 "ls" => Some(commands::ls(ctx).await),
541 "mkdir" => Some(commands::mkdir(ctx).await),
542 "rm" => Some(commands::rm(ctx).await),
543 "touch" => Some(commands::touch(ctx).await),
544 "cp" => Some(commands::cp(ctx).await),
545 "mv" => Some(commands::mv(ctx).await),
546 "basename" => Some(commands::basename(ctx).await),
547 "dirname" => Some(commands::dirname(ctx).await),
548 "env" => Some(commands::env(ctx).await),
549 "exit" => Some(commands::exit(ctx).await),
550 "which" => Some(commands::which(ctx).await),
551 "yes" => Some(commands::yes(ctx).await),
552 "seq" => Some(commands::seq(ctx).await),
553 "test" => Some(commands::test(ctx).await),
554 _ => None,
555 }
556 }
557
558 pub fn kill(&mut self) -> Result<()> {
560 self.cancelled = true;
561 if let Some(ref mut child) = self.child {
562 child.start_kill()?;
563 }
564 Ok(())
565 }
566
567 pub fn is_finished(&self) -> bool {
569 self.finished
570 }
571
572 pub fn result(&self) -> Option<&CommandResult> {
574 self.result.as_ref()
575 }
576
577 pub fn command(&self) -> &str {
579 &self.command
580 }
581
582 pub fn options(&self) -> &RunOptions {
584 &self.options
585 }
586}
587
588#[derive(Debug, Clone)]
590struct ShellConfig {
591 cmd: String,
592 args: Vec<String>,
593}
594
595fn find_available_shell() -> ShellConfig {
597 let is_windows = cfg!(windows);
598
599 if is_windows {
600 let shells = [
602 ("cmd.exe", vec!["/c"]),
603 ("powershell.exe", vec!["-Command"]),
604 ];
605
606 for (cmd, args) in shells {
607 if which::which(cmd).is_ok() {
608 return ShellConfig {
609 cmd: cmd.to_string(),
610 args: args.into_iter().map(String::from).collect(),
611 };
612 }
613 }
614
615 ShellConfig {
616 cmd: "cmd.exe".to_string(),
617 args: vec!["/c".to_string()],
618 }
619 } else {
620 let shells = [
622 ("/bin/sh", vec!["-c"]),
623 ("/usr/bin/sh", vec!["-c"]),
624 ("/bin/bash", vec!["-c"]),
625 ("sh", vec!["-c"]),
626 ];
627
628 for (cmd, args) in shells {
629 if std::path::Path::new(cmd).exists() || which::which(cmd).is_ok() {
630 return ShellConfig {
631 cmd: cmd.to_string(),
632 args: args.into_iter().map(String::from).collect(),
633 };
634 }
635 }
636
637 ShellConfig {
638 cmd: "/bin/sh".to_string(),
639 args: vec!["-c".to_string()],
640 }
641 }
642}
643
644pub async fn run(command: impl Into<String>) -> Result<CommandResult> {
649 let mut runner = ProcessRunner::new(command, RunOptions::default());
650 runner.run().await
651}
652
653pub use run as execute;
656
657pub async fn exec(command: impl Into<String>, options: RunOptions) -> Result<CommandResult> {
659 let mut runner = ProcessRunner::new(command, options);
660 runner.run().await
661}
662
663pub fn create(command: impl Into<String>, options: RunOptions) -> ProcessRunner {
665 ProcessRunner::new(command, options)
666}
667
668pub fn run_sync(command: impl Into<String>) -> Result<CommandResult> {
670 let rt = tokio::runtime::Runtime::new()?;
671 rt.block_on(run(command))
672}
673
674