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::{AsyncBufReadExt, AsyncWriteExt, BufReader};
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_quote_context_enabled, quote, quote_for_context, scan_quote_context, QuoteContext,
99};
100pub use state::{
101 get_shell_settings, global_state, reset_global_state, set_shell_option, unset_shell_option,
102 GlobalState, ShellSettings,
103};
104pub use stream::{AsyncIterator, IntoStream, OutputChunk, OutputStream, StreamingRunner};
105pub use trace::trace;
106
107fn resolve_spawn_cwd(cwd: Option<&PathBuf>) -> Option<PathBuf> {
119 if let Some(c) = cwd {
121 return Some(c.clone());
122 }
123
124 match std::env::current_dir() {
127 Ok(_) => None,
128 Err(e) => {
129 let fallback = std::env::var_os("HOME")
130 .or_else(|| std::env::var_os("USERPROFILE"))
131 .map(PathBuf::from)
132 .unwrap_or_else(std::env::temp_dir);
133 trace(
134 "ProcessRunner",
135 &format!(
136 "current_dir() failed ({}); spawning in fallback directory {}",
137 e,
138 fallback.display()
139 ),
140 );
141 if fallback.exists() {
142 Some(fallback)
143 } else {
144 Some(std::env::temp_dir())
145 }
146 }
147 }
148}
149
150#[derive(Debug, thiserror::Error)]
152pub enum Error {
153 #[error("IO error: {0}")]
154 Io(#[from] std::io::Error),
155
156 #[error("Command failed with exit code {code}: {message}")]
157 CommandFailed { code: i32, message: String },
158
159 #[error("Command not found: {0}")]
160 CommandNotFound(String),
161
162 #[error("Parse error: {0}")]
163 ParseError(String),
164
165 #[error("Cancelled")]
166 Cancelled,
167}
168
169pub type Result<T> = std::result::Result<T, Error>;
171
172#[derive(Debug, Clone)]
174pub struct RunOptions {
175 pub mirror: bool,
177 pub capture: bool,
179 pub stdin: StdinOption,
181 pub cwd: Option<PathBuf>,
183 pub env: Option<HashMap<String, String>>,
185 pub interactive: bool,
187 pub shell_operators: bool,
189 pub trace: bool,
191}
192
193impl Default for RunOptions {
194 fn default() -> Self {
195 RunOptions {
196 mirror: true,
197 capture: true,
198 stdin: StdinOption::Inherit,
199 cwd: None,
200 env: None,
201 interactive: false,
202 shell_operators: true,
203 trace: true,
204 }
205 }
206}
207
208#[derive(Debug, Clone)]
210pub enum StdinOption {
211 Inherit,
213 Pipe,
215 Content(String),
217 Null,
219}
220
221pub struct ProcessRunner {
223 command: String,
224 options: RunOptions,
225 child: Option<Child>,
226 result: Option<CommandResult>,
227 started: bool,
228 finished: bool,
229 cancelled: bool,
230 output_tx: Option<mpsc::Sender<StreamChunk>>,
231 output_rx: Option<mpsc::Receiver<StreamChunk>>,
232}
233
234impl ProcessRunner {
235 pub fn new(command: impl Into<String>, options: RunOptions) -> Self {
237 let (tx, rx) = mpsc::channel(1024);
238 ProcessRunner {
239 command: command.into(),
240 options,
241 child: None,
242 result: None,
243 started: false,
244 finished: false,
245 cancelled: false,
246 output_tx: Some(tx),
247 output_rx: Some(rx),
248 }
249 }
250
251 pub async fn start(&mut self) -> Result<()> {
253 if self.started {
254 return Ok(());
255 }
256 self.started = true;
257
258 utils::trace_lazy("ProcessRunner", || {
259 format!("Starting command: {}", self.command)
260 });
261
262 let first_word = if has_shell_escapes(&self.command) {
266 ""
267 } else {
268 self.command.split_whitespace().next().unwrap_or("")
269 };
270 if let Some(result) = self.try_virtual_command(first_word).await {
271 self.result = Some(result);
272 self.finished = true;
273 return Ok(());
274 }
275
276 let _parsed = if self.options.shell_operators && !needs_real_shell(&self.command) {
278 parse_shell_command(&self.command)
279 } else {
280 None
281 };
282
283 let shell = find_available_shell();
285
286 let mut cmd = Command::new(&shell.cmd);
287 for arg in &shell.args {
288 cmd.arg(arg);
289 }
290 cmd.arg(&self.command);
291
292 match &self.options.stdin {
294 StdinOption::Inherit => {
295 cmd.stdin(Stdio::inherit());
296 }
297 StdinOption::Pipe => {
298 cmd.stdin(Stdio::piped());
299 }
300 StdinOption::Content(_) => {
301 cmd.stdin(Stdio::piped());
302 }
303 StdinOption::Null => {
304 cmd.stdin(Stdio::null());
305 }
306 }
307
308 if self.options.capture || self.options.mirror {
310 cmd.stdout(Stdio::piped());
311 cmd.stderr(Stdio::piped());
312 } else {
313 cmd.stdout(Stdio::inherit());
314 cmd.stderr(Stdio::inherit());
315 }
316
317 if let Some(cwd) = resolve_spawn_cwd(self.options.cwd.as_ref()) {
320 cmd.current_dir(cwd);
321 }
322
323 if let Some(ref env_vars) = self.options.env {
325 for (key, value) in env_vars {
326 cmd.env(key, value);
327 }
328 }
329
330 let child = cmd.spawn()?;
332 self.child = Some(child);
333
334 Ok(())
335 }
336
337 pub async fn run(&mut self) -> Result<CommandResult> {
339 self.start().await?;
340
341 if let Some(result) = &self.result {
342 return Ok(result.clone());
343 }
344
345 let mut child = self.child.take().ok_or_else(|| {
346 Error::Io(std::io::Error::new(
347 std::io::ErrorKind::Other,
348 "Process not started",
349 ))
350 })?;
351
352 if let StdinOption::Content(ref content) = self.options.stdin {
354 if let Some(mut stdin) = child.stdin.take() {
355 let content = content.clone();
356 tokio::spawn(async move {
357 let _ = stdin.write_all(content.as_bytes()).await;
358 let _ = stdin.shutdown().await;
359 });
360 }
361 }
362
363 let mut stdout_content = String::new();
365 let mut stderr_content = String::new();
366
367 if let Some(stdout) = child.stdout.take() {
368 let mut reader = BufReader::new(stdout).lines();
369 while let Ok(Some(line)) = reader.next_line().await {
370 if self.options.mirror {
371 println!("{}", line);
372 }
373 stdout_content.push_str(&line);
374 stdout_content.push('\n');
375 }
376 }
377
378 if let Some(stderr) = child.stderr.take() {
379 let mut reader = BufReader::new(stderr).lines();
380 while let Ok(Some(line)) = reader.next_line().await {
381 if self.options.mirror {
382 eprintln!("{}", line);
383 }
384 stderr_content.push_str(&line);
385 stderr_content.push('\n');
386 }
387 }
388
389 let status = child.wait().await?;
390 let code = status.code().unwrap_or(-1);
391
392 let result = CommandResult {
393 stdout: stdout_content,
394 stderr: stderr_content,
395 code,
396 };
397
398 self.result = Some(result.clone());
399 self.finished = true;
400
401 Ok(result)
402 }
403
404 async fn try_virtual_command(&self, cmd_name: &str) -> Option<CommandResult> {
406 if !commands::are_virtual_commands_enabled() {
407 return None;
408 }
409
410 let parts: Vec<&str> = self.command.split_whitespace().collect();
412 let args: Vec<String> = parts.iter().skip(1).map(|s| s.to_string()).collect();
413
414 let ctx = CommandContext {
415 args,
416 stdin: match &self.options.stdin {
417 StdinOption::Content(s) => Some(s.clone()),
418 _ => None,
419 },
420 cwd: self.options.cwd.clone(),
421 env: self.options.env.clone(),
422 output_tx: self.output_tx.clone(),
423 is_cancelled: None,
424 };
425
426 match cmd_name {
427 "echo" => Some(commands::echo(ctx).await),
428 "pwd" => Some(commands::pwd(ctx).await),
429 "cd" => Some(commands::cd(ctx).await),
430 "true" => Some(commands::r#true(ctx).await),
431 "false" => Some(commands::r#false(ctx).await),
432 "sleep" => Some(commands::sleep(ctx).await),
433 "cat" => Some(commands::cat(ctx).await),
434 "ls" => Some(commands::ls(ctx).await),
435 "mkdir" => Some(commands::mkdir(ctx).await),
436 "rm" => Some(commands::rm(ctx).await),
437 "touch" => Some(commands::touch(ctx).await),
438 "cp" => Some(commands::cp(ctx).await),
439 "mv" => Some(commands::mv(ctx).await),
440 "basename" => Some(commands::basename(ctx).await),
441 "dirname" => Some(commands::dirname(ctx).await),
442 "env" => Some(commands::env(ctx).await),
443 "exit" => Some(commands::exit(ctx).await),
444 "which" => Some(commands::which(ctx).await),
445 "yes" => Some(commands::yes(ctx).await),
446 "seq" => Some(commands::seq(ctx).await),
447 "test" => Some(commands::test(ctx).await),
448 _ => None,
449 }
450 }
451
452 pub fn kill(&mut self) -> Result<()> {
454 self.cancelled = true;
455 if let Some(ref mut child) = self.child {
456 child.start_kill()?;
457 }
458 Ok(())
459 }
460
461 pub fn is_finished(&self) -> bool {
463 self.finished
464 }
465
466 pub fn result(&self) -> Option<&CommandResult> {
468 self.result.as_ref()
469 }
470
471 pub fn command(&self) -> &str {
473 &self.command
474 }
475
476 pub fn options(&self) -> &RunOptions {
478 &self.options
479 }
480}
481
482#[derive(Debug, Clone)]
484struct ShellConfig {
485 cmd: String,
486 args: Vec<String>,
487}
488
489fn find_available_shell() -> ShellConfig {
491 let is_windows = cfg!(windows);
492
493 if is_windows {
494 let shells = [
496 ("cmd.exe", vec!["/c"]),
497 ("powershell.exe", vec!["-Command"]),
498 ];
499
500 for (cmd, args) in shells {
501 if which::which(cmd).is_ok() {
502 return ShellConfig {
503 cmd: cmd.to_string(),
504 args: args.into_iter().map(String::from).collect(),
505 };
506 }
507 }
508
509 ShellConfig {
510 cmd: "cmd.exe".to_string(),
511 args: vec!["/c".to_string()],
512 }
513 } else {
514 let shells = [
516 ("/bin/sh", vec!["-c"]),
517 ("/usr/bin/sh", vec!["-c"]),
518 ("/bin/bash", vec!["-c"]),
519 ("sh", vec!["-c"]),
520 ];
521
522 for (cmd, args) in shells {
523 if std::path::Path::new(cmd).exists() || which::which(cmd).is_ok() {
524 return ShellConfig {
525 cmd: cmd.to_string(),
526 args: args.into_iter().map(String::from).collect(),
527 };
528 }
529 }
530
531 ShellConfig {
532 cmd: "/bin/sh".to_string(),
533 args: vec!["-c".to_string()],
534 }
535 }
536}
537
538pub async fn run(command: impl Into<String>) -> Result<CommandResult> {
543 let mut runner = ProcessRunner::new(command, RunOptions::default());
544 runner.run().await
545}
546
547pub use run as execute;
550
551pub async fn exec(command: impl Into<String>, options: RunOptions) -> Result<CommandResult> {
553 let mut runner = ProcessRunner::new(command, options);
554 runner.run().await
555}
556
557pub fn create(command: impl Into<String>, options: RunOptions) -> ProcessRunner {
559 ProcessRunner::new(command, options)
560}
561
562pub fn run_sync(command: impl Into<String>) -> Result<CommandResult> {
564 let rt = tokio::runtime::Runtime::new()?;
565 rt.block_on(run(command))
566}
567
568