Skip to main content

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(&current) {
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}