cosh_tools/bash/bsh.rs
1//! Bash executor with security pattern validation and async streaming output.
2//!
3//! Provides [`run`] and [`spawn_bash`] for executing bash commands as child
4//! processes, returning output as an async stream of chunks. Commands are
5//! validated against a set of dangerous patterns (recursive rm, fork bombs,
6//! disk destruction, etc.) before execution.
7
8use async_stream::stream;
9#[cfg(any(unix, windows))]
10use portable_pty::CommandBuilder;
11#[cfg(any(unix, windows))]
12use portable_pty::PtySize;
13#[cfg(any(unix, windows))]
14use portable_pty::native_pty_system;
15use regex::Regex;
16use std::path::Path;
17use std::pin::Pin;
18use std::process::Stdio;
19use std::sync::LazyLock;
20
21#[cfg(any(unix, windows))]
22use std::io::Read;
23#[cfg(unix)]
24use std::sync::Arc;
25#[cfg(unix)]
26use std::sync::atomic::{AtomicBool, Ordering};
27use std::time::Duration;
28#[cfg(unix)]
29use std::time::Instant;
30use tokio::io::AsyncReadExt;
31use tokio::io::{Error, ErrorKind};
32use tokio_stream::Stream;
33use tokio_stream::StreamExt;
34
35#[cfg(any(unix, windows))]
36use tokio::sync::mpsc;
37
38/// Returns the list of critical bash patterns that are checked before execution.
39///
40/// # Panics
41///
42/// Panics if any of the internal regular expressions fail to compile (they are
43/// static and guaranteed valid on first access).
44#[allow(clippy::unwrap_used)]
45pub fn critical_bash_patterns() -> &'static [Regex] {
46 static PATTERNS: LazyLock<Vec<Regex>> = LazyLock::new(|| {
47 vec![
48 // Recursive destruction.
49 Regex::new(r"(?i)\brm\s+-[a-z]*[rRfF][a-z]*\s+\/").unwrap(),
50 Regex::new(r"(?i)\bsudo\s+rm\b").unwrap(),
51 Regex::new(r"(?i)\bchmod\s+-R\s+[0-7]+\s+\/").unwrap(),
52 Regex::new(r"\bchmod\s+-R\s+[ugoa+\-=rwxXst,]+\s+\/").unwrap(),
53 Regex::new(r"(?i)\bchown\s+-R\s+\S+\s+\/").unwrap(),
54 // Fork bomb.
55 Regex::new(r"(?i):\(\)\s*\{\s*:\s*\|\s*:").unwrap(),
56 // Disk / filesystem destruction.
57 Regex::new(r"(?i)>\s*\/dev\/sd[a-z]").unwrap(),
58 Regex::new(r"(?i)\bmkfs(\.|\b)").unwrap(),
59 Regex::new(r"(?i)\bdd\s+if=.+of=\/dev\/").unwrap(),
60 Regex::new(r"(?i)\bshred\s+\/dev\/").unwrap(),
61 Regex::new(r"(?i)\bcryptsetup\b").unwrap(),
62 // System-config destruction.
63 Regex::new(r"(?i)>\s*\/etc\/(?:passwd|shadow|sudoers)\b").unwrap(),
64 Regex::new(r"(?i)\btee\s+(?:-a\s+)?\/etc\/(?:passwd|shadow|sudoers)\b").unwrap(),
65 // Remote-fetch-then-execute.
66 Regex::new(r"(?i)\b(?:curl|wget|fetch)\b[^|]*\|\s*(?:bash|sh|zsh|fish)\b").unwrap(),
67 Regex::new(r"(?i)(?:^|[\s;&|(])(?:bash|sh|zsh|source|\.)\s+<\(\s*(?:curl|wget|fetch)\b").unwrap(),
68 Regex::new(r#"(?i)\beval\s+["'`]?\$\(\s*(?:curl|wget|fetch)\b|\beval\s+`\s*(?:curl|wget|fetch)\b"#).unwrap(),
69 // Process/host control.
70 Regex::new(r"\bkill\s+-9\s+1\b").unwrap(),
71 Regex::new(r"(?i)(?:^|[\s;&|(])(?:shutdown|poweroff|reboot|halt)(?:\s|$|[;|&])").unwrap(),
72 Regex::new(r"(?i)(?:^|[\s;&|(])init\s+0\b").unwrap(),
73 // Network-shell exfil.
74 Regex::new(r"(?i)\bnc\b[^|;]*\s-[a-zA-Z]*[ec][a-zA-Z]*\s").unwrap(),
75 ]
76 });
77 &PATTERNS
78}
79// Constants & lazy statics
80
81const BUFFER_SIZE: usize = 4096;
82
83/// How often the Unix PTY read loop wakes up to run the input-wait watchdog.
84#[cfg(unix)]
85const WATCHDOG_TICK_MS: libc::c_int = 250;
86
87/// How long a child may sit in raw (non-canonical) terminal mode — the
88/// signature of an interactive pager/TUI waiting for keystrokes the harness
89/// can never send — before the input-wait watchdog kills it.
90#[cfg(unix)]
91const WATCHDOG_RAW_MODE_GRACE: Duration = Duration::from_secs(2);
92
93/// Diagnostic emitted on stderr when the input-wait watchdog kills a child.
94///
95/// Deliberately DESCRIPTIVE ONLY — no suggested commands. This message only
96/// appears when the auto-fix layer could not rewrite the command with
97/// 100% certainty (otherwise there is nothing to kill), so any specific
98/// advice could be wrong for the unknown case; the model infers the fix
99/// itself for the retry.
100#[cfg(unix)]
101const WATCHDOG_MESSAGE: &str = "input-wait watchdog: killed a process that entered raw (interactive) \
102 terminal mode and waited for keyboard input the harness cannot provide \
103 (interactive pager/editor/TUI).";
104
105static ENV_VAR_PATTERN: LazyLock<Regex> =
106 LazyLock::new(|| Regex::new(r"^[A-Za-z_][A-Za-z0-9_]*$").unwrap());
107
108#[allow(clippy::unwrap_used)]
109fn env_var_pattern() -> &'static Regex {
110 &ENV_VAR_PATTERN
111}
112
113/// Resolve the bash executable to spawn, per platform.
114///
115/// - Unix: plain `"bash"` (resolved through `PATH`, as before).
116/// - Windows: an absolute path to a Git Bash (`bash.exe`) — never the
117/// `C:\Windows\System32\bash.exe` WSL launcher. `CreateProcessW` searches
118/// System32 **before** `PATH`, so a bare `"bash"` would silently route
119/// commands into WSL (ignoring the Windows `cwd` and filesystem). The
120/// resolution order is:
121/// 1. `where.exe bash.exe` — first hit that is not the WSL launcher
122/// (respects the user's PATH, e.g. a scoop/custom install);
123/// 2. the standard Git-for-Windows install locations
124/// (`...\Git\bin\bash.exe`, then `...\Git\usr\bin\bash.exe`).
125///
126/// `bin\bash.exe` is preferred over `usr\bash.exe`: the former is a login
127/// wrapper that adds Git's `usr\bin` to `PATH` for the child, so commands
128/// like `ls` and `seq` resolve; the latter relies on the inherited `PATH`.
129///
130/// # Panics
131///
132/// Panics if no Git Bash can be located — this is a programming/deployment
133/// error (the harness requires Git for Windows), not a per-command failure.
134#[cfg(windows)]
135fn resolve_bash() -> &'static str {
136 static RESOLVED: LazyLock<String> = LazyLock::new(|| {
137 let is_real_bash = |p: &Path| {
138 p.file_name()
139 .and_then(|n| n.to_str())
140 .is_some_and(|n| n.eq_ignore_ascii_case("bash.exe"))
141 && !p.components().any(|c| {
142 c.as_os_str().to_str().is_some_and(|s| {
143 s.eq_ignore_ascii_case("WindowsApps") || s.eq_ignore_ascii_case("System32")
144 })
145 })
146 };
147
148 // 1. Respect the user's PATH via `where.exe`.
149 if let Ok(output) = std::process::Command::new("where.exe")
150 .arg("bash.exe")
151 .stdin(Stdio::null())
152 .output()
153 && output.status.success()
154 {
155 for line in String::from_utf8_lossy(&output.stdout).lines() {
156 let path = Path::new(line.trim());
157 if is_real_bash(path) {
158 return path.to_string_lossy().into_owned();
159 }
160 }
161 }
162
163 // 2. Standard Git-for-Windows locations. `bin\bash.exe` is a login
164 // wrapper that sets up Git's POSIX PATH; `usr\bin` is the raw
165 // Cygwin binary.
166 for candidate in [
167 r"C:\Program Files\Git\bin\bash.exe",
168 r"C:\Program Files (x86)\Git\bin\bash.exe",
169 r"C:\Git\bin\bash.exe",
170 r"C:\Program Files\Git\usr\bin\bash.exe",
171 r"C:\Program Files (x86)\Git\usr\bin\bash.exe",
172 r"C:\Git\usr\bin\bash.exe",
173 ] {
174 let path = Path::new(candidate);
175 if path.is_file() {
176 return candidate.to_string();
177 }
178 }
179
180 panic!(
181 "bash.exe not found: install Git for Windows (https://git-scm.com/download/win) \
182 or ensure a Git Bash `bash.exe` is on PATH"
183 );
184 });
185 RESOLVED.as_str()
186}
187
188#[cfg(not(windows))]
189fn resolve_bash() -> &'static str {
190 "bash"
191}
192
193/// Maps a Unix signal *description* (from `strsignal(3)`) to its numeric value.
194///
195/// `portable_pty::ExitStatus::signal()` returns the description string
196/// produced by the libc `strsignal()` function (e.g. `"Killed"` for
197/// SIGKILL), not the `"SIGKILL"` constant name. We map back to the
198/// numeric value to match `SpawnOutput::signal` in the non-PTY path
199/// (which uses `std::os::unix::process::ExitStatusExt`).
200#[cfg(unix)]
201fn signal_name_to_number(name: &str) -> Option<i32> {
202 match name {
203 "Hangup" => Some(1),
204 "Interrupt" => Some(2),
205 "Quit" => Some(3),
206 "Illegal instruction" => Some(4),
207 "Trace/breakpoint trap" => Some(5),
208 "Aborted" => Some(6),
209 "Bus error" => Some(7),
210 "Arithmetic exception" | "Floating point exception" => Some(8),
211 "Killed" => Some(9),
212 "User defined signal 1" => Some(10),
213 "Segmentation fault" => Some(11),
214 "User defined signal 2" => Some(12),
215 "Broken pipe" => Some(13),
216 "Alarm clock" => Some(14),
217 "Terminated" => Some(15),
218 "Stack fault" => Some(16),
219 "Child exited" => Some(17),
220 "Continued" => Some(18),
221 "Stopped (signal)" => Some(19),
222 "Stopped" => Some(20),
223 "Stopped (tty input)" => Some(21),
224 "Stopped (tty output)" => Some(22),
225 "Urgent I/O condition" => Some(23),
226 "CPU time limit exceeded" => Some(24),
227 "File size limit exceeded" => Some(25),
228 "Virtual timer expired" => Some(26),
229 "Profiling timer expired" => Some(27),
230 "Window changed" => Some(28),
231 "I/O possible" => Some(29),
232 "Power failure" => Some(30),
233 "Bad system call" => Some(31),
234 _ => None,
235 }
236}
237
238// Types
239
240pub struct BashOutput {
241 pub stdout: Option<String>,
242 pub exit_code: Option<i32>,
243 pub signal: Option<i32>,
244}
245
246pub struct SpawnOutput {
247 pub stdout: Vec<u8>,
248 pub stderr: Vec<u8>,
249 pub exit_code: Option<i32>,
250 pub signal: Option<i32>,
251 pub truncated: bool,
252}
253
254pub struct ExecError {
255 pub stderr: Option<String>,
256 pub signal: Option<i32>,
257}
258
259pub struct BashError {
260 pub text_err: Option<String>,
261 pub exec_err: Option<ExecError>,
262}
263// Public API
264
265pub(super) fn validate_bash_patterns(command: &str) -> Result<(), String> {
266 let patterns = critical_bash_patterns();
267 for pattern in patterns {
268 if pattern.is_match(command) {
269 return Err(format!("Pattern match found: {}", pattern.as_str()));
270 }
271 }
272 Ok(())
273}
274
275/// Execute a bash command and return its output as an async stream.
276///
277/// Validates the command against dangerous patterns and spawns
278/// `bash -c <command>` as a child process in the given `cwd`.
279/// Returns a stream of [`SpawnOutput`] items.
280///
281/// Before spawning, the command goes through [`auto_fix_command`]: when the
282/// command can be rewritten into a provably non-blocking form with 100%
283/// certainty, the rewritten form runs and a `auto-fix: …` note is yielded as
284/// the first stderr item. Ambiguous commands run unchanged.
285///
286/// When `timeout_ms` is `Some(ms)`, the child is killed after `ms`
287/// milliseconds and the stream yields a final item with
288/// `signal: Some(-1)` and `exit_code: None`.
289///
290/// # Errors
291///
292/// Returns [`BashError`] if the command is an absolute path or matches a
293/// dangerous security pattern. Both the original and the auto-fixed command
294/// are validated, so a rewrite can never smuggle a pattern past the guards.
295pub fn run(
296 timeout_ms: Option<u64>,
297 env: &Option<Vec<(String, String)>>,
298 pty: bool,
299 command: &str,
300 cwd: &str,
301) -> Result<Pin<Box<dyn Stream<Item = SpawnOutput> + Send>>, BashError> {
302 // Guards
303 // On Windows, POSIX-style absolute commands (`/bin/echo hi`) are NOT
304 // `Path::is_absolute`, but Git Bash resolves them against its MSYS root
305 // — they are absolute from the shell's perspective and are blocked too.
306 let is_absolute = Path::new(command).is_absolute() || command.starts_with('/');
307 if is_absolute {
308 return Err(BashError {
309 text_err: Some("absolute command not allowed, use relative path".to_string()),
310 exec_err: None,
311 });
312 }
313
314 if let Err(err) = validate_bash_patterns(command) {
315 return Err(BashError {
316 text_err: Some(err),
317 exec_err: None,
318 });
319 }
320
321 // Auto-fix: rewrite paged commands when 100% certain; ambiguous
322 // commands pass through unchanged (the input-wait watchdog remains the
323 // backstop, so the worst case is the pre-existing behavior). The
324 // rewritten command is owned by this frame and moved into the 'static
325 // stream below, so the stream never borrows a local.
326 #[cfg(any(unix, windows))]
327 let (fixed_command, autofix_note) = auto_fix_command(command);
328 #[cfg(any(unix, windows))]
329 if fixed_command != command
330 && let Err(err) = validate_bash_patterns(&fixed_command)
331 {
332 return Err(BashError {
333 text_err: Some(err),
334 exec_err: None,
335 });
336 }
337 #[cfg(any(unix, windows))]
338 let command: String = fixed_command;
339 #[cfg(not(any(unix, windows)))]
340 let command: String = command.to_string();
341 #[cfg(not(any(unix, windows)))]
342 let autofix_note: Option<String> = None;
343
344 let env = (*env).clone();
345 let cwd = cwd.to_string();
346 let use_pty = pty;
347
348 Ok(Box::pin(stream! {
349 // Surface the rewrite first, so consumers see exactly which command
350 // is about to run (the model must know it did not run verbatim).
351 if let Some(note) = autofix_note {
352 yield SpawnOutput {
353 stdout: vec![],
354 stderr: note.into_bytes(),
355 exit_code: None,
356 signal: None,
357 truncated: false,
358 };
359 }
360
361 let mut stream: Pin<Box<dyn Stream<Item = Result<SpawnOutput, Error>> + Send>> = if use_pty {
362 #[cfg(any(unix, windows))]
363 {
364 spawn_bash_pty(env, &cwd, &command, timeout_ms)
365 }
366 #[cfg(not(any(unix, windows)))]
367 {
368 spawn_bash(env, cwd.as_str(), command.as_str(), timeout_ms)
369 }
370 } else {
371 spawn_bash(env, cwd.as_str(), command.as_str(), timeout_ms)
372 };
373 while let Some(item) = stream.next().await {
374 match item {
375 Ok(output) => yield output,
376 Err(e) => {
377 // Surface stream/spawn errors instead of dropping them:
378 // a failed spawn (e.g. an invalid cwd) or a mid-stream
379 // read error used to end the stream SILENTLY with zero
380 // items, leaving callers with no output and no
381 // explanation. Convert the error into a stderr-bearing
382 // item so consumers always see what went wrong.
383 yield SpawnOutput {
384 stdout: vec![],
385 stderr: format!("bash stream error: {e}").into_bytes(),
386 exit_code: None,
387 signal: None,
388 truncated: false,
389 };
390 }
391 }
392 }
393 }))
394}
395
396/// Default environment variables injected into every spawned bash to keep the
397/// child process non-interactive.
398///
399/// Interactive pagers (`less`, `more`) and editors are the most common cause
400/// of a command "hanging": when stdout is a PTY (or the pager is inherited
401/// from the harness's own environment), tools like `git` launch `less` on
402/// `git log` / `git diff`, which then blocks forever waiting for keystrokes
403/// that never arrive — the run only ends when the timeout fires.
404const NON_INTERACTIVE_DEFAULTS: &[(&str, &str)] = &[
405 // Plain-text pagers: never launch an interactive pager.
406 ("PAGER", "cat"),
407 ("GIT_PAGER", "cat"),
408 ("MANPAGER", "cat"),
409 // If a pager still runs (e.g. an explicit `| less`), environment flags
410 // CANNOT make it non-interactive: less reads keystrokes directly from
411 // /dev/tty, so it blocks forever inside a PTY no matter what LESS/PAGER
412 // vars say. That case is handled by the input-wait watchdog in
413 // spawn_bash_pty, not by environment variables.
414 // Editors must never open an interactive UI inside a spawned command.
415 ("GIT_EDITOR", "true"),
416 ("EDITOR", "true"),
417 ("VISUAL", "true"),
418];
419
420/// Compute which non-interactive defaults still need to be injected.
421///
422/// Any variable **explicitly** provided by the caller wins over the default:
423/// only the defaults whose keys are absent from `caller_env` are returned.
424/// Note that an inherited value from the parent process's environment does
425/// NOT win — the parent may itself run under `PAGER=less`, which is exactly
426/// the hang this mechanism prevents.
427fn non_interactive_defaults(
428 caller_env: &Option<Vec<(String, String)>>,
429) -> Vec<(&'static str, &'static str)> {
430 let caller_keys: Vec<&str> = caller_env
431 .as_ref()
432 .map(|pairs| pairs.iter().map(|(k, _)| k.as_str()).collect())
433 .unwrap_or_default();
434 NON_INTERACTIVE_DEFAULTS
435 .iter()
436 .copied()
437 .filter(|(key, _)| !caller_keys.contains(key))
438 .collect()
439}
440
441// --- Auto-fix: deterministic rewrites of commands that would block on an ---
442// --- interactive pager, applied before spawn. Fallback = current behavior. ---
443//
444// The rewrite layer ONLY fires when it can be 100% certain the rewrite is
445// safe (no quoting, no shell operators it doesn't understand). When it
446// cannot be certain, it returns the command untouched and the input-wait
447// watchdog remains the backstop — so the worst case is exactly the
448// pre-existing behavior.
449
450/// Regex: `git [env-prefixes] [global-flags...] <paged-subcommand>` at the
451/// start of the command line. The `git` token is captured so the rewrite can
452/// insert `--no-pager` at its exact position even with env prefixes present.
453/// Global flags that consume a separate value are enumerated explicitly
454/// (`-c`, `-C`, `--git-dir`, `--work-tree`, `--namespace`,
455/// `--super-prefix` — git's fixed global-option list); anything else stays
456/// unmatched → no rewrite.
457#[cfg(any(unix, windows))]
458static AUTO_FIX_GIT_RE: LazyLock<Regex> = LazyLock::new(|| {
459 Regex::new(concat!(
460 r"^(?:\s*[A-Za-z_][A-Za-z0-9_]*=\S*\s+)*(git)",
461 r"(?:\s+(?:",
462 r"-[cC]\s+\S+", // -c <name>=<value> / -C <path>
463 r"|--(?:git-dir|work-tree|namespace|super-prefix)\s+\S+", // --opt <value>
464 r"|-{1,2}[\w-]+\S*", // valueless global flags
465 r"))*",
466 r"(?:\s+(log|diff|show|blame|shortlog|reflog|whatchanged)\b)",
467 ))
468 .unwrap()
469});
470
471/// Regex: a final pipeline segment that is ONLY a pager plus flag tokens:
472/// `... | less`, `| less -X`. Anything else in the tail makes the match
473/// fail → no rewrite (fallback stays the watchdog): a filename argument
474/// would make the pager ignore stdin, and a `>`/`<` redirect must never be
475/// swallowed by the replacement.
476#[cfg(any(unix, windows))]
477static AUTO_FIX_PIPE_PAGER_RE: LazyLock<Regex> =
478 LazyLock::new(|| Regex::new(r"\|\s*(less|more|most)\b(?:\s+-{1,2}[\w][\w-]*)*\s*$").unwrap());
479
480/// Regex: `less`/`more` invoked as the command itself with plain filename
481/// arguments (no flags — flags would be invalid for `cat`).
482#[cfg(any(unix, windows))]
483static AUTO_FIX_PAGER_CMD_RE: LazyLock<Regex> =
484 LazyLock::new(|| Regex::new(r"^(less|more)\s+([^;|&]+)$").unwrap());
485
486/// Rewrite `command` so it cannot block on an interactive pager, when doing
487/// so is 100% certain; otherwise return the command unchanged.
488///
489/// Returns `(possibly-rewritten-command, note)` where `note` — when present —
490/// describes the rewrite and MUST be surfaced to the caller (the model) as
491/// stderr, so it knows exactly which command actually ran.
492///
493/// Certainty rules shared by every pattern below:
494/// - no quotes of any kind in the command (a quoted string containing e.g.
495/// `git log` must not be rewritten);
496/// - anything ambiguous → no rewrite (the watchdog stays the backstop).
497///
498/// Rules are applied repeatedly until the command stabilizes (bounded), so
499/// composite forms like `git log | less` are resolved by a single applicable
500/// rule — the pipeline rule — instead of stacking rewrites.
501#[cfg(any(unix, windows))]
502pub(crate) fn auto_fix_command(command: &str) -> (String, Option<String>) {
503 use regex::Captures;
504
505 // Applies at most one rule. Returns `Some((rewritten, note))` or `None`.
506 fn apply_one(command: &str) -> Option<(String, String)> {
507 // Already explicitly non-paged via git's own flag → nothing to do.
508 // Also stops the fix loop right after a git rewrite.
509 if command.contains("--no-pager") {
510 return None;
511 }
512
513 // Rule A: final pipeline segment is a pager → replace it with `cat`
514 // (same bytes, no pagination). The regex guarantees the segment is
515 // the tail of the command and contains no `|`, `;` or `&`.
516 if AUTO_FIX_PIPE_PAGER_RE.is_match(command) {
517 let fixed = AUTO_FIX_PIPE_PAGER_RE.replace(command, |_: &Captures| "| cat");
518 return Some((
519 fixed.into_owned(),
520 "auto-fix: replaced the trailing pager in the pipeline with `cat` to keep the output complete".to_string(),
521 ));
522 }
523
524 // Rule B: `git [env-prefixes] [global-flags] <paged-subcommand>` →
525 // insert `--no-pager` right after `git`. Only when NOT piped: with
526 // stdout attached to a pipe git never launches its pager anyway,
527 // and the flag would be pure noise. Compound commands are also
528 // refused: in `git log && less file`, fixing only the git segment
529 // would leave a blocking pager behind — a partial rewrite is not
530 // 100% certain, so anything with `;` or `&` passes through.
531 if !command.contains('|')
532 && !command.contains(';')
533 && !command.contains('&')
534 && let Some(cap) = AUTO_FIX_GIT_RE.captures(command)
535 {
536 let git_token = cap.get(1)?;
537 let sub = cap.get(2)?.as_str();
538 // Insert right after the captured `git` token — not at the
539 // match start, which can be preceded by env assignments.
540 let insert_at = git_token.end();
541 let mut fixed = String::with_capacity(command.len() + "--no-pager ".len());
542 fixed.push_str(&command[..insert_at]);
543 fixed.push_str(" --no-pager");
544 fixed.push_str(&command[insert_at..]);
545 return Some((
546 fixed,
547 format!(
548 "auto-fix: inserted `--no-pager` into `git {sub}` to prevent an interactive pager"
549 ),
550 ));
551 }
552
553 // Rule C: `less/more <plain filenames>` → `cat <filenames>`. The
554 // args must not start a token with `-` (less flags are not valid
555 // cat flags) and the regex excludes `;`, `|`, `&`.
556 if let Some(cap) = AUTO_FIX_PAGER_CMD_RE.captures(command) {
557 let args = cap.get(2)?.as_str();
558 if args.split_whitespace().all(|tok| !tok.starts_with('-')) {
559 let pager = cap.get(1)?.as_str();
560 return Some((
561 format!("cat {args}"),
562 format!(
563 "auto-fix: replaced `{pager}` with `cat` to print the file(s) directly"
564 ),
565 ));
566 }
567 // `less -X file` etc. → not certain → no rewrite.
568 }
569
570 None
571 }
572
573 // Global certainty gate: any quoting anywhere in the line disqualifies
574 // every rule. A quoted literal could contain the patterns below without
575 // being a command (`echo "git log"`); regex cannot prove intent, so we
576 // refuse to touch it.
577 if command.contains('"') || command.contains('\'') || command.contains('`') {
578 return (command.to_string(), None);
579 }
580
581 // Multi-line commands are out of scope for every rule: a `\n` can hide
582 // a second command or a heredoc body that a single-line-oriented
583 // rewrite would corrupt or silently delete.
584 if command.contains('\n') {
585 return (command.to_string(), None);
586 }
587
588 // An explicit git pager opt-in (`--paginate`) must win over our
589 // insertion: git's --paginate/--no-pager are last-wins, so inserting
590 // --no-pager BEFORE it would re-enable the pager and the note would
591 // claim a fix that does not happen. Refuse → watchdog fallback.
592 if command.contains("--paginate") {
593 return (command.to_string(), None);
594 }
595
596 // Fix loop: apply one rule at a time until nothing applies. The bound
597 // is defensive — every rule is convergent, but never risk unbounded
598 // rewriting on a pathological input.
599 let mut current = command.to_string();
600 let mut notes: Vec<String> = Vec::new();
601 for _ in 0..4 {
602 match apply_one(¤t) {
603 Some((fixed, note)) => {
604 notes.push(note);
605 current = fixed;
606 }
607 None => break,
608 }
609 }
610
611 let note = (!notes.is_empty()).then(|| notes.join("; "));
612 (current, note)
613}
614
615/// Spawn a bash process and return its output as an async stream.
616///
617/// When `timeout_ms` is `Some(ms)`, the child is killed after `ms`
618/// milliseconds and the stream yields a final item with
619/// `signal: Some(-1)` and `exit_code: None`.
620///
621/// # Panics
622///
623/// Panics if the child process stdout or stderr pipe cannot be taken (this
624/// only happens if [`std::process::Stdio::piped`] was not set).
625pub(crate) fn spawn_bash(
626 env: Option<Vec<(String, String)>>,
627 cwd: impl Into<String>,
628 command: impl Into<String>,
629 timeout_ms: Option<u64>,
630) -> Pin<Box<dyn Stream<Item = Result<SpawnOutput, Error>> + Send>> {
631 // Own `cwd`/`command` up front so the returned stream is 'static — it
632 // must never borrow a caller's string, because `run()` passes an
633 // auto-fixed command that is a local of its own frame.
634 let cwd = cwd.into();
635 let command = command.into();
636 let mut buffer_stdout = [0u8; BUFFER_SIZE];
637 let mut buffer_stderr = [0u8; BUFFER_SIZE];
638 let mut cmd = tokio::process::Command::new(resolve_bash());
639
640 Box::pin(stream! {
641 if let Some(ref env) = env {
642 let valid_pattern = env_var_pattern();
643 for (key, value) in env {
644 if !valid_pattern.is_match(key) {
645 yield Err(Error::new(
646 ErrorKind::InvalidInput,
647 format!("invalid env variable name: {key}"),
648 ));
649 return;
650 }
651 cmd.env(key, value);
652 }
653 }
654
655 // Non-interactive defaults: caller-provided env wins, inherited env
656 // does not. Prevents pagers/editors from hanging on the piped path
657 // when the parent process itself runs under e.g. PAGER=less.
658 for (key, value) in non_interactive_defaults(&env) {
659 cmd.env(key, value);
660 }
661
662 // Empty cwd means "inherit the parent process's working directory".
663 // `current_dir("")` would fail the spawn and (before the stream-error
664 // surfacing fix) produced a silently EMPTY stream.
665 if !cwd.is_empty() {
666 cmd.current_dir(cwd);
667 }
668 cmd.arg("-c").arg(command);
669 cmd.stdin(Stdio::null());
670 cmd.stdout(Stdio::piped()).stderr(Stdio::piped());
671
672 let mut child = match cmd.spawn() {
673 Ok(c) => c,
674 Err(e) => {
675 yield Err(Error::other(format!("failed to spawn bash process: {e}")));
676 return;
677 }
678 };
679
680 #[allow(clippy::expect_used)]
681 let mut stream_stdout = child.stdout.take().expect("stdout pipe should be configured");
682 #[allow(clippy::expect_used)]
683 let mut stream_stderr = child.stderr.take().expect("stderr pipe should be configured");
684 let mut stdout_done = false;
685 let mut stderr_done = false;
686
687 // Zero timeout: the command must not run at all. Kill it immediately
688 // without reading any output, so the result is deterministic — the
689 // child never gets a chance to write before the 0ms timer fires.
690 if timeout_ms == Some(0) {
691 let _ = child.kill().await;
692 let _ = child.wait().await;
693 yield Ok(SpawnOutput {
694 stdout: vec![],
695 stderr: vec![],
696 exit_code: None,
697 signal: Some(-1_i32),
698 truncated: false,
699 });
700 return;
701 }
702
703 let deadline = timeout_ms.map(|ms| tokio::time::Instant::now() + Duration::from_millis(ms));
704
705 loop {
706 tokio::select! {
707 biased;
708
709 () = async {
710 match deadline {
711 Some(dl) => tokio::time::sleep_until(dl).await,
712 None => std::future::pending::<()>().await,
713 }
714 } => {
715 let _ = child.kill().await;
716 yield Ok(SpawnOutput {
717 stdout: vec![],
718 stderr: vec![],
719 exit_code: None,
720 signal: Some(-1_i32),
721 truncated: false,
722 });
723 return;
724 }
725
726 result_stdout = stream_stdout.read(&mut buffer_stdout), if !stdout_done => {
727 match result_stdout {
728 Ok(0) => stdout_done = true,
729 Ok(n) => {
730 let stdout = buffer_stdout[..n].to_vec();
731 yield Ok(SpawnOutput {
732 stdout,
733 stderr: vec![],
734 exit_code: None,
735 signal: None,
736 truncated: n == BUFFER_SIZE,
737 });
738 }
739 Err(e) => {
740 yield Err(Error::other(format!("stdout read error: {e}")));
741 stdout_done = true;
742 }
743 }
744 }
745
746 result_stderr = stream_stderr.read(&mut buffer_stderr), if !stderr_done => {
747 match result_stderr {
748 Ok(0) => stderr_done = true,
749 Ok(n) => {
750 let stderr = buffer_stderr[..n].to_vec();
751 yield Ok(SpawnOutput {
752 stdout: vec![],
753 stderr,
754 exit_code: None,
755 signal: None,
756 truncated: n == BUFFER_SIZE,
757 });
758 }
759 Err(e) => {
760 yield Err(Error::other(format!("stderr read error: {e}")));
761 stderr_done = true;
762 }
763 }
764 }
765 }
766 if stdout_done && stderr_done { break; }
767 }
768
769 let status = match child.wait().await {
770 Ok(s) => s,
771 Err(e) => {
772 yield Err(Error::other(format!("failed to wait for child: {e}")));
773 return;
774 }
775 };
776
777 let signal = {
778 #[cfg(unix)]
779 {
780 use std::os::unix::process::ExitStatusExt;
781 status.signal()
782 }
783 #[cfg(not(unix))]
784 {
785 None
786 }
787 };
788
789 yield Ok(SpawnOutput {
790 stdout: vec![],
791 stderr: vec![],
792 exit_code: status.code(),
793 signal,
794 truncated: false,
795 });
796 })
797}
798
799/// Spawn a bash process into a PTY and return its output as an async stream.
800///
801/// Unlike [`spawn_bash`], this uses a pseudo-terminal so the child process
802/// behaves as if connected to a terminal (colored output, prompts, etc.).
803/// stdout and stderr are multiplexed into the PTY; [`SpawnOutput::stderr`]
804/// is always empty in this path.
805///
806/// I/O is bridged from the synchronous `portable_pty` reader to the async
807/// stream via [`tokio::task::spawn_blocking`] and an `mpsc` channel, adding
808/// one buffer copy per chunk.
809///
810/// # Timeout mechanism
811///
812/// A **self-pipe** trick is used to unblock the blocking [`poll(2)`] call
813/// from the async timeout handler without requiring a separate monitor
814/// thread:
815///
816/// 1. A pipe is created before [`tokio::task::spawn_blocking`]
817/// is called, giving one fd for each side of the pipe (via
818/// [`libc::pipe2`]; on macOS, which has no `pipe2`, via
819/// [`libc::pipe`] + `fcntl(F_SETFD, FD_CLOEXEC)`).
820/// 2. The read end (`pipe_rx`) is moved into the blocking task; the write
821/// end (`pipe_tx`) stays on the async side.
822/// 3. In the blocking task, [`libc::poll`] watches **both** the PTY master
823/// file descriptor and `pipe_rx` for readability.
824/// 4. When the async timeout fires, it writes a single byte to `pipe_tx`.
825/// `poll` wakes immediately, the blocking task detects the byte on the
826/// pipe fd, breaks out of the read loop, kills and reaps the child.
827///
828/// # Safety
829///
830/// The `unsafe` blocks in this function are:
831///
832/// | Location | Call | Invariant |
833/// |---|---|---|
834/// | Stream setup | `pipe2` (macOS: `pipe` + `fcntl`) | `pipe_fds` is a valid pointer to 2 `i32`s |
835/// | Timeout handler | `write` + `close` on `pipe_tx` | `pipe_tx` is a valid fd, not used after |
836/// | Blocking task | `poll`, `read`, `close` on `pipe_rx` and `pty_fd` | Both fds are valid and open for the lifetime of `poll_fds`; `pipe_rx` is closed once after use |
837/// | Watchdog check | `mem::zeroed` + `tcgetattr` on `pty_fd` | `termios` is a fully-owned stack out-buffer, zeroed before the call; `pty_fd` is valid and open for the loop's lifetime (POSIX: `tcgetattr` on a PTY master reads the attached slave's line discipline) |
838///
839/// These are trivially verified by inspection — the pipe fds are created
840/// together, one is consumed per side, and each is closed exactly once.
841///
842/// [`poll(2)`]: https://man7.org/linux/man-pages/man2/poll.2.html
843#[cfg(unix)]
844pub(crate) fn spawn_bash_pty(
845 env: Option<Vec<(String, String)>>,
846 cwd: &str,
847 command: &str,
848 timeout_ms: Option<u64>,
849) -> Pin<Box<dyn Stream<Item = Result<SpawnOutput, Error>> + Send>> {
850 let cwd = cwd.to_string();
851 let command = command.to_string();
852 let kill_flag = Arc::new(AtomicBool::new(false));
853
854 Box::pin(stream! {
855 // Zero timeout: the command must not run at all. Yield the timeout
856 // item immediately without spawning the process (deterministic).
857 if timeout_ms == Some(0) {
858 yield Ok(SpawnOutput {
859 stdout: vec![],
860 stderr: vec![],
861 exit_code: None,
862 signal: Some(-1_i32),
863 truncated: false,
864 });
865 return;
866 }
867
868 let (tx, mut rx) = mpsc::unbounded_channel();
869 let kill = kill_flag.clone();
870
871 // Create a self-pipe before spawn_blocking so the write end
872 // (pipe_tx) is accessible from the async timeout handler.
873 let mut pipe_fds: [libc::c_int; 2] = [0; 2];
874 // SAFETY: pipe/pipe2 are safe per POSIX; fds provides a valid array
875 // pointer. macOS has no pipe2, so there the CLOEXEC flag is set
876 // after the fact via fcntl(F_SETFD) — before any child exists, so
877 // nothing can inherit the fds in practice (the only theoretical
878 // window is a concurrent exec elsewhere in the process, which this
879 // single-threaded setup point doesn't have).
880 let pipe_result = unsafe {
881 #[cfg(not(target_os = "macos"))]
882 {
883 libc::pipe2(pipe_fds.as_mut_ptr(), libc::O_CLOEXEC)
884 }
885 #[cfg(target_os = "macos")]
886 {
887 let rc = libc::pipe(pipe_fds.as_mut_ptr());
888 if rc == 0 {
889 for fd in pipe_fds {
890 libc::fcntl(fd, libc::F_SETFD, libc::FD_CLOEXEC);
891 }
892 }
893 rc
894 }
895 };
896 if pipe_result != 0_i32 {
897 yield Err(Error::other("failed to create self-pipe for timeout"));
898 return;
899 }
900 let [pipe_rx, pipe_tx] = pipe_fds;
901
902 tokio::task::spawn_blocking(move || {
903 // Validate environment variables (same rules as spawn_bash).
904 if let Some(ref env) = env {
905 let valid_pattern = env_var_pattern();
906 for (key, _) in env {
907 if !valid_pattern.is_match(key) {
908 let _ = tx.send(Err(Error::new(
909 ErrorKind::InvalidInput,
910 format!("invalid env variable name: {key}"),
911 )));
912 return;
913 }
914 }
915 }
916
917 let pty_system = native_pty_system();
918 let pair = match pty_system.openpty(PtySize {
919 rows: 24,
920 cols: 80,
921 pixel_width: 0,
922 pixel_height: 0,
923 }) {
924 Ok(p) => p,
925 Err(e) => {
926 let _ = tx.send(Err(Error::other(format!("failed to open pty: {e}"))));
927 return;
928 }
929 };
930
931 let mut cmd_builder = CommandBuilder::new("bash");
932 cmd_builder.arg("-c");
933 cmd_builder.arg(&command);
934 // Empty cwd → inherit the parent's working directory (same
935 // semantics as the piped path; an empty cwd fails the spawn).
936 if !cwd.is_empty() {
937 cmd_builder.cwd(&cwd);
938 }
939
940 // Non-interactive defaults: caller-provided env wins, inherited
941 // env does not. The PTY makes stdout look like an interactive
942 // terminal, so git etc. would otherwise launch a full-screen
943 // pager (less) that blocks forever waiting for keystrokes.
944 let non_interactive = non_interactive_defaults(&env);
945
946 if let Some(env) = env {
947 for (key, value) in env {
948 cmd_builder.env(key, value);
949 }
950 }
951 for (key, value) in non_interactive {
952 cmd_builder.env(key, value);
953 }
954
955 let mut child = match pair.slave.spawn_command(cmd_builder) {
956 Ok(c) => c,
957 Err(e) => {
958 let _ = tx.send(Err(Error::other(
959 format!("failed to spawn command in pty: {e}"),
960 )));
961 return;
962 }
963 };
964
965 let mut reader = match pair.master.try_clone_reader() {
966 Ok(r) => r,
967 Err(e) => {
968 let _ = tx.send(Err(Error::other(
969 format!("failed to clone pty reader: {e}"),
970 )));
971 return;
972 }
973 };
974
975 // Dropping the slave closes its end of the PTY, signalling
976 // EOF to the master reader once the child exits.
977 drop(pair.slave);
978
979 // Get the PTY master fd for poll(). Both pair.master and the
980 // cloned reader share the same underlying PTY, so polling the
981 // master fd for POLLIN and reading from the reader is coherent.
982 let Some(pty_fd) = pair.master.as_raw_fd() else {
983 let _ = tx.send(Err(Error::other("failed to get PTY fd")));
984 return;
985 };
986
987 // poll_fds[0] = PTY master, poll_fds[1] = self-pipe read end.
988 let mut poll_fds = [
989 libc::pollfd { fd: pty_fd, events: libc::POLLIN, revents: 0 },
990 libc::pollfd { fd: pipe_rx, events: libc::POLLIN, revents: 0 },
991 ];
992
993 let mut buf = [0u8; BUFFER_SIZE];
994 // Input-wait watchdog state: timestamp of the first consecutive
995 // observation of raw (non-canonical) terminal mode, if any.
996 let mut watchdog_raw_since: Option<Instant> = None;
997 loop {
998 // Watchdog tick or I/O readiness: poll's return value is
999 // not needed — revents below drive every event handler.
1000 let _n_ready = loop {
1001
1002 // Bounded timeout so the loop wakes up periodically to
1003 // run the input-wait watchdog even with no I/O.
1004 let res = unsafe { libc::poll(poll_fds.as_mut_ptr(), 2, WATCHDOG_TICK_MS) };
1005 if res < 0_i32 {
1006 let err = std::io::Error::last_os_error();
1007 if err.raw_os_error() == Some(libc::EINTR) {
1008 continue;
1009 }
1010 let _ = tx.send(Err(Error::other(format!("pty poll error: {err}"))));
1011 return;
1012 }
1013 break res;
1014 };
1015
1016 // n_ready == 0 (watchdog tick, nothing readable): fall
1017 // through — all revents are zero, so every event handler
1018 // below no-ops and control reaches the input-wait check at
1019 // the bottom of the loop.
1020
1021 // Self-pipe has data → timeout requested by the async handler.
1022 if poll_fds[1].revents & (libc::POLLIN | libc::POLLHUP | libc::POLLERR) != 0 {
1023 // Consume the notification byte (ignore errors — we
1024 // only care that poll woke up).
1025 let _ = unsafe {
1026 libc::read(
1027 pipe_rx,
1028 buf.as_mut_ptr().cast::<libc::c_void>(),
1029 buf.len(),
1030 )
1031 };
1032 break;
1033 }
1034
1035 // PTY has data available.
1036 if poll_fds[0].revents & libc::POLLIN != 0 {
1037 match reader.read(&mut buf) {
1038 Ok(0) => break,
1039 Ok(n) => {
1040 if tx
1041 .send(Ok(SpawnOutput {
1042 stdout: buf[..n].to_vec(),
1043 stderr: vec![],
1044 exit_code: None,
1045 signal: None,
1046 truncated: n == BUFFER_SIZE,
1047 }))
1048 .is_err()
1049 {
1050 // Receiver dropped (stream cancelled).
1051 break;
1052 }
1053 }
1054 Err(e) => {
1055 let _ = tx.send(Err(Error::other(format!("pty read error: {e}"))));
1056 return;
1057 }
1058 }
1059 }
1060
1061 // PTY hung up (child exited).
1062 if poll_fds[0].revents & (libc::POLLHUP | libc::POLLERR) != 0 {
1063 // Drain any remaining data before breaking.
1064 match reader.read(&mut buf) {
1065 Ok(0) | Err(_) => {}
1066 Ok(n) => {
1067 let _ = tx.send(Ok(SpawnOutput {
1068 stdout: buf[..n].to_vec(),
1069 stderr: vec![],
1070 exit_code: None,
1071 signal: None,
1072 truncated: n == BUFFER_SIZE,
1073 }));
1074 }
1075 }
1076 break;
1077 }
1078
1079 // Input-wait watchdog: a live child that keeps the terminal
1080 // in non-canonical (raw) mode is an interactive pager/editor/
1081 // TUI waiting for keystrokes. Those keystrokes would come
1082 // from /dev/tty, which the harness never writes to, so the
1083 // child can only be freed by killing it. Environment vars
1084 // cannot prevent this hang (the pager reads the TTY
1085 // directly). A brief raw-mode flip is normal (line editing,
1086 // `read -n`), hence the sustained-raw-mode grace window.
1087 let child_alive = matches!(child.try_wait(), Ok(None));
1088 if !child_alive {
1089 watchdog_raw_since = None;
1090 } else {
1091 let mut termios: libc::termios = unsafe { std::mem::zeroed() };
1092 let is_raw = unsafe { libc::tcgetattr(pty_fd, &mut termios) } == 0
1093 && termios.c_lflag & libc::ICANON == 0;
1094 if !is_raw {
1095 watchdog_raw_since = None;
1096 } else {
1097 let since =
1098 *watchdog_raw_since.get_or_insert_with(Instant::now);
1099 if since.elapsed() >= WATCHDOG_RAW_MODE_GRACE {
1100 let _ = child.kill();
1101 let _ = tx.send(Ok(SpawnOutput {
1102 stdout: vec![],
1103 stderr: WATCHDOG_MESSAGE.as_bytes().to_vec(),
1104 exit_code: None,
1105 signal: None,
1106 truncated: false,
1107 }));
1108 break;
1109 }
1110 }
1111 }
1112 }
1113
1114 // Close pipe read end (write end is closed by the async handler).
1115 // SAFETY: pipe_rx is a valid fd not used after this point.
1116 unsafe { libc::close(pipe_rx); }
1117
1118 if kill.load(Ordering::Relaxed) {
1119 let _ = child.kill();
1120 }
1121
1122 match child.wait() {
1123 Ok(status) => {
1124 let signal = status.signal().and_then(signal_name_to_number);
1125 // Match non-PTY semantics: exit_code is None when
1126 // killed by a signal; only Some when the process
1127 // exited normally.
1128 let exit_code = if signal.is_some() {
1129 None
1130 } else {
1131 #[allow(clippy::cast_possible_wrap)]
1132 Some(status.exit_code() as i32)
1133 };
1134 let _ = tx.send(Ok(SpawnOutput {
1135 stdout: vec![],
1136 stderr: vec![],
1137 exit_code,
1138 signal,
1139 truncated: false,
1140 }));
1141 }
1142 Err(e) => {
1143 let _ = tx.send(Err(Error::other(format!("pty wait error: {e}"))));
1144 }
1145 }
1146 });
1147
1148 let deadline = timeout_ms.map(|ms| tokio::time::Instant::now() + Duration::from_millis(ms));
1149 loop {
1150 tokio::select! {
1151 biased;
1152
1153 () = async {
1154 match deadline {
1155 Some(dl) => tokio::time::sleep_until(dl).await,
1156 None => std::future::pending::<()>().await,
1157 }
1158 } => {
1159 kill_flag.store(true, Ordering::Relaxed);
1160
1161 // Write a byte to the self-pipe so the blocking task's
1162 // poll() wakes up and detects the timeout.
1163 let byte: u8 = 0;
1164 // SAFETY: pipe_tx is a valid fd owned by this scope.
1165 unsafe {
1166 libc::write(
1167 pipe_tx,
1168 (&raw const byte).cast::<libc::c_void>(),
1169 1,
1170 );
1171 }
1172 // SAFETY: pipe_tx is never used after this point.
1173 unsafe { libc::close(pipe_tx); }
1174
1175 yield Ok(SpawnOutput {
1176 stdout: vec![],
1177 stderr: vec![],
1178 exit_code: None,
1179 signal: Some(-1_i32),
1180 truncated: false,
1181 });
1182 return;
1183 }
1184
1185 item = rx.recv() => {
1186 match item {
1187 Some(result) => yield result,
1188 None => break,
1189 }
1190 }
1191 }
1192 }
1193
1194 // Stream ended normally — close the pipe write end.
1195 // SAFETY: pipe_tx was not closed by the timeout handler (the
1196 // handler returned above).
1197 unsafe { libc::close(pipe_tx); }
1198 })
1199}
1200
1201/// Spawn a bash process into a Windows ConPTY and return its output as an
1202/// async stream.
1203///
1204/// Windows counterpart of the Unix [`spawn_bash_pty`]: uses `portable_pty`'s
1205/// ConPTY backend (Windows 10 1809+) so the child runs attached to a pseudo
1206/// terminal; stdout and stderr are multiplexed and [`SpawnOutput::stderr`] is
1207/// always empty.
1208///
1209/// # Windows ConPTY specifics
1210///
1211/// Unlike a Unix PTY, the ConPTY **output pipe never reaches EOF when the
1212/// child exits** — the conhost side keeps the pipe handle open, so a blocking
1213/// reader would hang forever. Two protocol quirks are handled here (both
1214/// verified experimentally against Git Bash under ConPTY):
1215///
1216/// 1. **DSR cursor query.** bash emits `ESC[6n` (cursor position request) at
1217/// startup and blocks until the terminal answers. The stream watches for
1218/// the sequence in the output and replies `ESC[1;1R` through the PTY
1219/// writer.
1220/// 2. **No EOF.** Stream termination is driven by the child's exit status:
1221/// once the watcher thread reports the reaped child, the reader is given a
1222/// short drain window for in-flight chunks, then the final status item is
1223/// emitted and the (still blocked) reader thread is detached.
1224///
1225/// The blocking reader runs in `tokio::task::spawn_blocking` and forwards
1226/// chunks through an unbounded `mpsc` channel; `child.wait()` runs in its own
1227/// thread because it may block longer than the stream lives.
1228///
1229/// # Timeout mechanism
1230///
1231/// The async deadline arms `tokio::select!` against the channel receiver. On
1232/// timeout the child is killed via `portable_pty::ChildKiller::kill`
1233/// (`TerminateProcess`) — callable from the async side through a pre-cloned
1234/// killer — in-flight chunks are drained, and the final `signal: Some(-1)`
1235/// item is yielded. The watcher's exit status is discarded in that case,
1236/// matching the Unix path (timeout item wins).
1237#[cfg(windows)]
1238use portable_pty::ChildKiller;
1239
1240#[cfg(windows)]
1241pub(crate) fn spawn_bash_pty(
1242 env: Option<Vec<(String, String)>>,
1243 cwd: &str,
1244 command: &str,
1245 timeout_ms: Option<u64>,
1246) -> Pin<Box<dyn Stream<Item = Result<SpawnOutput, Error>> + Send>> {
1247 let cwd = cwd.to_string();
1248 let command = command.to_string();
1249
1250 Box::pin(stream! {
1251 // Zero timeout: the command must not run at all. Yield the timeout
1252 // item immediately without spawning the process (deterministic).
1253 if timeout_ms == Some(0) {
1254 yield Ok(SpawnOutput {
1255 stdout: vec![],
1256 stderr: vec![],
1257 exit_code: None,
1258 signal: Some(-1_i32),
1259 truncated: false,
1260 });
1261 return;
1262 }
1263
1264 // Validate environment variables (same rules as spawn_bash) before
1265 // spawning anything.
1266 if let Some(env) = &env {
1267 let valid_pattern = env_var_pattern();
1268 for (key, _) in env {
1269 if !valid_pattern.is_match(key) {
1270 yield Err(Error::new(
1271 ErrorKind::InvalidInput,
1272 format!("invalid env variable name: {key}"),
1273 ));
1274 return;
1275 }
1276 }
1277 }
1278
1279 let pty_system = native_pty_system();
1280 let pair = match pty_system.openpty(PtySize {
1281 rows: 24,
1282 cols: 80,
1283 pixel_width: 0,
1284 pixel_height: 0,
1285 }) {
1286 Ok(p) => p,
1287 Err(e) => {
1288 yield Err(Error::other(format!("failed to open pty: {e}")));
1289 return;
1290 }
1291 };
1292
1293 let mut cmd_builder = CommandBuilder::new(resolve_bash());
1294 cmd_builder.arg("-c");
1295 cmd_builder.arg(&command);
1296 // Empty cwd → inherit the parent's working directory (same
1297 // semantics as the piped path; an empty cwd fails the spawn).
1298 if !cwd.is_empty() {
1299 cmd_builder.cwd(&cwd);
1300 }
1301
1302 if let Some(env) = env {
1303 for (key, value) in env {
1304 cmd_builder.env(key, value);
1305 }
1306 }
1307
1308 let mut child = match pair.slave.spawn_command(cmd_builder) {
1309 Ok(c) => c,
1310 Err(e) => {
1311 yield Err(Error::other(format!("failed to spawn command in pty: {e}")));
1312 return;
1313 }
1314 };
1315
1316 let mut reader = match pair.master.try_clone_reader() {
1317 Ok(r) => r,
1318 Err(e) => {
1319 yield Err(Error::other(format!("failed to clone pty reader: {e}")));
1320 return;
1321 }
1322 };
1323 // Take the writer so we can answer DSR cursor queries.
1324 let mut writer = match pair.master.take_writer() {
1325 Ok(w) => w,
1326 Err(e) => {
1327 yield Err(Error::other(format!("failed to take pty writer: {e}")));
1328 return;
1329 }
1330 };
1331 // Pre-clone the killer so the async side can terminate the child on
1332 // timeout while the watcher thread holds the blocking `wait()`.
1333 let mut killer = ChildKiller::clone_killer(&*child);
1334
1335 // Reader thread: forwards chunks through an unbounded channel. An
1336 // empty Vec signals EOF/read-error. The thread may stay blocked in
1337 // read() forever (ConPTY never EOFs) — it is detached when the
1338 // stream ends; the OS reclaims it when the process exits.
1339 let (tx, mut rx) = mpsc::unbounded_channel::<Vec<u8>>();
1340 tokio::task::spawn_blocking(move || {
1341 let mut buf = [0u8; BUFFER_SIZE];
1342 loop {
1343 match reader.read(&mut buf) {
1344 Ok(0) | Err(_) => break,
1345 Ok(n) => {
1346 if tx.send(buf[..n].to_vec()).is_err() {
1347 break;
1348 }
1349 }
1350 }
1351 }
1352 let _ = tx.send(Vec::new());
1353 });
1354
1355 // Watcher thread: reaps the child and reports its status as
1356 // (exit_code, signal). ConPTY has no signal concept — a killed child
1357 // surfaces as a non-zero exit code.
1358 let (stx, mut srx) = mpsc::unbounded_channel::<(Option<i32>, Option<i32>)>();
1359 std::thread::spawn(move || match child.wait() {
1360 Ok(status) => {
1361 #[allow(clippy::cast_possible_wrap)]
1362 let _ = stx.send((Some(status.exit_code() as i32), None));
1363 }
1364 Err(e) => {
1365 let _ = stx.send((None, None));
1366 let _ = e;
1367 }
1368 });
1369
1370 const DSR_QUERY: &[u8] = b"\x1b[6n";
1371 const DSR_REPLY: &[u8] = b"\x1b[1;1R";
1372 const DRAIN_WINDOW: Duration = Duration::from_millis(300);
1373
1374 let deadline =
1375 timeout_ms.map(|ms| tokio::time::Instant::now() + Duration::from_millis(ms));
1376 let mut timed_out = false;
1377 let mut answered_dsr = false;
1378 let mut status: Option<(Option<i32>, Option<i32>)> = None;
1379 // Chunks received during the post-exit drain window.
1380 let mut drained: Vec<Vec<u8>> = Vec::new();
1381
1382 loop {
1383 tokio::select! {
1384 biased;
1385
1386 () = async {
1387 match deadline {
1388 Some(dl) => tokio::time::sleep_until(dl).await,
1389 None => std::future::pending::<()>().await,
1390 }
1391 }, if deadline.is_some() && !timed_out => {
1392 timed_out = true;
1393 killer.kill().ok();
1394 break;
1395 }
1396
1397 status_msg = srx.recv() => {
1398 match status_msg {
1399 Some(s) => {
1400 status = Some(s);
1401 // Child reaped: give in-flight chunks a short
1402 // drain window, then end the stream (no EOF).
1403 let _ = tokio::time::timeout(DRAIN_WINDOW, async {
1404 while let Some(chunk) = rx.recv().await {
1405 if chunk.is_empty() {
1406 break;
1407 }
1408 drained.push(chunk);
1409 }
1410 })
1411 .await;
1412 break;
1413 }
1414 None => break,
1415 }
1416 }
1417
1418 chunk = rx.recv() => {
1419 match chunk {
1420 Some(data) => {
1421 if data.is_empty() {
1422 // Reader EOF (unusual on ConPTY): keep
1423 // waiting for the watcher's status.
1424 continue;
1425 }
1426 if !answered_dsr
1427 && data.windows(DSR_QUERY.len()).any(|w| w == DSR_QUERY)
1428 {
1429 answered_dsr = true;
1430 let _ = writer.write_all(DSR_REPLY);
1431 let _ = writer.flush();
1432 }
1433 yield Ok(SpawnOutput {
1434 stdout: data,
1435 stderr: vec![],
1436 exit_code: None,
1437 signal: None,
1438 truncated: false,
1439 });
1440 }
1441 None => break,
1442 }
1443 }
1444 }
1445 }
1446
1447 // Emit chunks drained during the post-exit window before the final
1448 // status item (chunks always precede the status item).
1449 for chunk in drained {
1450 yield Ok(SpawnOutput {
1451 stdout: chunk,
1452 stderr: vec![],
1453 exit_code: None,
1454 signal: None,
1455 truncated: false,
1456 });
1457 }
1458
1459 if timed_out {
1460 yield Ok(SpawnOutput {
1461 stdout: vec![],
1462 stderr: vec![],
1463 exit_code: None,
1464 signal: Some(-1_i32),
1465 truncated: false,
1466 });
1467 return;
1468 }
1469
1470 if let Some((exit_code, signal)) = status {
1471 yield Ok(SpawnOutput {
1472 stdout: vec![],
1473 stderr: vec![],
1474 exit_code,
1475 signal,
1476 truncated: false,
1477 });
1478 }
1479 // Watcher died without a status (should not happen) — end the
1480 // stream silently, matching the Unix path's behavior on error.
1481 })
1482}