1use std::collections::HashMap;
28use std::path::PathBuf;
29use std::process::Stdio;
30use tokio::io::{AsyncReadExt, AsyncWriteExt};
31use tokio::process::Command;
32
33use crate::trace::trace_lazy;
34use crate::{CommandResult, Result, RunOptions, StdinOption};
35
36struct VirtualCommandResult {
37 result: CommandResult,
38 cd_context: Option<crate::commands::cd::CdContext>,
39}
40
41#[derive(Debug, Clone)]
45pub struct Pipeline {
46 commands: Vec<String>,
48 stdin: Option<String>,
50 cwd: Option<PathBuf>,
52 env: Option<HashMap<String, String>>,
54 mirror: bool,
56 capture: bool,
58}
59
60impl Default for Pipeline {
61 fn default() -> Self {
62 Self::new()
63 }
64}
65
66impl Pipeline {
67 pub fn new() -> Self {
69 Pipeline {
70 commands: Vec::new(),
71 stdin: None,
72 cwd: None,
73 env: None,
74 mirror: true,
75 capture: true,
76 }
77 }
78
79 #[allow(clippy::should_implement_trait)]
84 pub fn add(mut self, command: impl Into<String>) -> Self {
85 self.commands.push(command.into());
86 self
87 }
88
89 pub fn stdin(mut self, content: impl Into<String>) -> Self {
91 self.stdin = Some(content.into());
92 self
93 }
94
95 pub fn cwd(mut self, path: impl Into<PathBuf>) -> Self {
97 self.cwd = Some(path.into());
98 self
99 }
100
101 pub fn env(mut self, env: HashMap<String, String>) -> Self {
103 self.env = Some(env);
104 self
105 }
106
107 pub fn mirror_output(mut self, mirror: bool) -> Self {
109 self.mirror = mirror;
110 self
111 }
112
113 pub fn capture_output(mut self, capture: bool) -> Self {
115 self.capture = capture;
116 self
117 }
118
119 pub async fn run(self) -> Result<CommandResult> {
121 if self.commands.is_empty() {
122 return Ok(CommandResult {
123 stdout: String::new(),
124 stderr: "No commands in pipeline".to_string(),
125 code: 1,
126 });
127 }
128
129 trace_lazy("Pipeline", || {
130 format!("Running pipeline with {} commands", self.commands.len())
131 });
132
133 let mut current_stdin = self.stdin.clone();
134 let mut effective_cwd = self.cwd.clone();
135 let mut effective_env = self.env.clone();
136 let mut last_result = CommandResult {
137 stdout: String::new(),
138 stderr: String::new(),
139 code: 0,
140 };
141 let mut accumulated_stderr = String::new();
142
143 for (i, cmd_str) in self.commands.iter().enumerate() {
144 let is_last = i == self.commands.len() - 1;
145
146 trace_lazy("Pipeline", || {
147 format!(
148 "Executing command {}/{}: {}",
149 i + 1,
150 self.commands.len(),
151 cmd_str
152 )
153 });
154
155 let first_word = cmd_str.split_whitespace().next().unwrap_or("");
157 if crate::commands::are_virtual_commands_enabled() {
158 if let Some(result) = self
159 .try_virtual_command(
160 first_word,
161 cmd_str,
162 ¤t_stdin,
163 effective_cwd.as_ref(),
164 effective_env.as_ref(),
165 )
166 .await
167 {
168 let VirtualCommandResult { result, cd_context } = result;
169 if result.code != 0 {
170 return Ok(CommandResult {
171 stdout: result.stdout,
172 stderr: accumulated_stderr + &result.stderr,
173 code: result.code,
174 });
175 }
176 current_stdin = Some(result.stdout.clone());
177 accumulated_stderr.push_str(&result.stderr);
178 if let Some(context) = cd_context {
179 let env = effective_env.get_or_insert_with(|| std::env::vars().collect());
180 env.insert(
181 "OLDPWD".to_string(),
182 context.oldpwd.to_string_lossy().to_string(),
183 );
184 env.insert("PWD".to_string(), context.cwd.to_string_lossy().to_string());
185 effective_cwd = Some(context.cwd);
186 }
187 last_result = result;
188 continue;
189 }
190 }
191
192 let shell = find_available_shell();
194 let mut cmd = Command::new(&shell.cmd);
195 for arg in &shell.args {
196 cmd.arg(arg);
197 }
198 cmd.arg(crate::utils::with_exported_process_context(
199 cmd_str,
200 effective_env.as_ref(),
201 ));
202
203 cmd.stdin(Stdio::piped());
205 cmd.stdout(Stdio::piped());
206 cmd.stderr(Stdio::piped());
207
208 if let Some(cwd) = crate::resolve_spawn_cwd(effective_cwd.as_ref()) {
211 cmd.current_dir(cwd);
212 }
213
214 if let Some(ref env_vars) = effective_env {
216 for (key, value) in env_vars {
217 cmd.env(key, value);
218 }
219 }
220
221 let mut child = cmd.spawn()?;
223
224 if let Some(ref stdin_content) = current_stdin {
226 if let Some(mut stdin) = child.stdin.take() {
227 let content = stdin_content.clone();
228 tokio::spawn(async move {
229 let _ = stdin.write_all(content.as_bytes()).await;
230 let _ = stdin.shutdown().await;
231 });
232 }
233 }
234
235 let mut stdout_content = String::new();
237 if let Some(mut stdout) = child.stdout.take() {
238 stdout.read_to_string(&mut stdout_content).await?;
239 }
240
241 let mut stderr_content = String::new();
243 if let Some(mut stderr) = child.stderr.take() {
244 stderr.read_to_string(&mut stderr_content).await?;
245 }
246
247 if is_last && self.mirror {
249 if !stdout_content.is_empty() {
250 print!("{}", stdout_content);
251 }
252 if !stderr_content.is_empty() {
253 eprint!("{}", stderr_content);
254 }
255 }
256
257 let status = child.wait().await?;
259 let code = status.code().unwrap_or(-1);
260
261 accumulated_stderr.push_str(&stderr_content);
262
263 if code != 0 {
264 return Ok(CommandResult {
265 stdout: stdout_content,
266 stderr: accumulated_stderr,
267 code,
268 });
269 }
270
271 current_stdin = Some(stdout_content.clone());
273 last_result = CommandResult {
274 stdout: stdout_content,
275 stderr: String::new(),
276 code,
277 };
278 }
279
280 Ok(CommandResult {
281 stdout: last_result.stdout,
282 stderr: accumulated_stderr,
283 code: last_result.code,
284 })
285 }
286
287 async fn try_virtual_command(
289 &self,
290 cmd_name: &str,
291 full_cmd: &str,
292 stdin: &Option<String>,
293 cwd: Option<&PathBuf>,
294 env: Option<&HashMap<String, String>>,
295 ) -> Option<VirtualCommandResult> {
296 let parts: Vec<&str> = full_cmd.split_whitespace().collect();
297 let args: Vec<String> = parts.iter().skip(1).map(|s| s.to_string()).collect();
298
299 let ctx = crate::commands::CommandContext {
300 args,
301 stdin: stdin.clone(),
302 cwd: cwd.cloned(),
303 env: env.cloned(),
304 output_tx: None,
305 is_cancelled: None,
306 };
307
308 let (result, cd_context) = match cmd_name {
309 "echo" => (crate::commands::echo(ctx).await, None),
310 "pwd" => (crate::commands::pwd(ctx).await, None),
311 "cd" => crate::commands::cd::resolve_cd(ctx).await,
312 "true" => (crate::commands::r#true(ctx).await, None),
313 "false" => (crate::commands::r#false(ctx).await, None),
314 "sleep" => (crate::commands::sleep(ctx).await, None),
315 "cat" => (crate::commands::cat(ctx).await, None),
316 "ls" => (crate::commands::ls(ctx).await, None),
317 "mkdir" => (crate::commands::mkdir(ctx).await, None),
318 "rm" => (crate::commands::rm(ctx).await, None),
319 "touch" => (crate::commands::touch(ctx).await, None),
320 "cp" => (crate::commands::cp(ctx).await, None),
321 "mv" => (crate::commands::mv(ctx).await, None),
322 "basename" => (crate::commands::basename(ctx).await, None),
323 "dirname" => (crate::commands::dirname(ctx).await, None),
324 "env" => (crate::commands::env(ctx).await, None),
325 "exit" => (crate::commands::exit(ctx).await, None),
326 "which" => (crate::commands::which(ctx).await, None),
327 "yes" => (crate::commands::yes(ctx).await, None),
328 "seq" => (crate::commands::seq(ctx).await, None),
329 "test" => (crate::commands::test(ctx).await, None),
330 _ => return None,
331 };
332 Some(VirtualCommandResult { result, cd_context })
333 }
334}
335
336#[derive(Debug, Clone)]
338struct ShellConfig {
339 cmd: String,
340 args: Vec<String>,
341}
342
343fn find_available_shell() -> ShellConfig {
345 let is_windows = cfg!(windows);
346
347 if is_windows {
348 ShellConfig {
349 cmd: "cmd.exe".to_string(),
350 args: vec!["/c".to_string()],
351 }
352 } else {
353 let shells = [
354 ("/bin/sh", "-c"),
355 ("/usr/bin/sh", "-c"),
356 ("/bin/bash", "-c"),
357 ];
358
359 for (cmd, arg) in shells {
360 if std::path::Path::new(cmd).exists() {
361 return ShellConfig {
362 cmd: cmd.to_string(),
363 args: vec![arg.to_string()],
364 };
365 }
366 }
367
368 ShellConfig {
369 cmd: "/bin/sh".to_string(),
370 args: vec!["-c".to_string()],
371 }
372 }
373}
374
375pub trait PipelineExt {
377 fn pipe(self, command: impl Into<String>) -> PipelineBuilder;
379}
380
381impl PipelineExt for crate::ProcessRunner {
382 fn pipe(self, command: impl Into<String>) -> PipelineBuilder {
383 PipelineBuilder {
384 first: self,
385 additional: vec![command.into()],
386 }
387 }
388}
389
390pub struct PipelineBuilder {
392 first: crate::ProcessRunner,
393 additional: Vec<String>,
394}
395
396impl PipelineBuilder {
397 pub fn pipe(mut self, command: impl Into<String>) -> Self {
399 self.additional.push(command.into());
400 self
401 }
402
403 pub async fn run(mut self) -> Result<CommandResult> {
405 let first_result = self.first.run().await?;
407
408 if first_result.code != 0 {
409 return Ok(first_result);
410 }
411
412 let mut current_stdin = Some(first_result.stdout);
414 let mut accumulated_stderr = first_result.stderr;
415 let mut last_result = CommandResult {
416 stdout: String::new(),
417 stderr: String::new(),
418 code: 0,
419 };
420
421 for cmd_str in &self.additional {
422 let mut runner = crate::ProcessRunner::new(
423 cmd_str.clone(),
424 RunOptions {
425 stdin: StdinOption::Content(current_stdin.take().unwrap_or_default()),
426 mirror: false,
427 capture: true,
428 ..Default::default()
429 },
430 );
431
432 let result = runner.run().await?;
433 accumulated_stderr.push_str(&result.stderr);
434
435 if result.code != 0 {
436 return Ok(CommandResult {
437 stdout: result.stdout,
438 stderr: accumulated_stderr,
439 code: result.code,
440 });
441 }
442
443 current_stdin = Some(result.stdout.clone());
444 last_result = result;
445 }
446
447 Ok(CommandResult {
448 stdout: last_result.stdout,
449 stderr: accumulated_stderr,
450 code: last_result.code,
451 })
452 }
453}