Skip to main content

qframe/runtime/
process.rs

1//! Running a child process and reading its output line by line, for showing in a log view.
2//!
3//! Two modes, and the difference matters:
4//!
5//! - **Pipes** (the default) keep standard output and standard error apart, so a failure stays
6//!   recognisable as a failure. Programs that check for a terminal drop their progress bar and
7//!   their colour when they write to a pipe.
8//! - **A pseudo-terminal** ([`Process::pty`]) gives the child a terminal of the size we choose,
9//!   so it draws its progress. Both of its streams land on that one terminal, so every line
10//!   arrives as [`Line::Out`].
11//!
12//! A line that a `\r` overwrites, such as each frame of a progress bar, is dropped as a screen
13//! would drop it, unless the frames are asked for with [`Process::run_with_overwritten`].
14//!
15//! Output nobody reads line by line is the other kind: a build's whole log, a program's answer,
16//! a flood of progress. [`Process::collect`] runs the child the same way but keeps only the end
17//! of what it wrote, within [`Keep`], and gives it back as one piece of text.
18
19use std::collections::VecDeque;
20use std::ffi::OsString;
21use std::io::{self, Read};
22use std::path::PathBuf;
23use std::process::{Child, Command, Stdio};
24use std::sync::mpsc::{self, RecvTimeoutError, SyncSender};
25use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
26use std::time::{Duration, Instant};
27
28/// How long the loop waits for the next line before it looks at `cancel` again.
29const POLL: Duration = Duration::from_millis(10);
30
31/// How much is read from a stream at a time.
32pub(super) const CHUNK: usize = 4096;
33
34/// The longest line held at once. A program that writes more without a newline has its line
35/// delivered in pieces of at most this many bytes, so the reader never holds all of it.
36const MAX_LINE: usize = 64 * 1024;
37
38/// How many lines wait for `on_line` at most. Beyond that the readers stop reading, the pipe
39/// fills and the child waits, so a child that writes faster than the application reads does not
40/// pile its output up in memory.
41const QUEUE: usize = 1024;
42
43/// A child process whose output is read line by line, for showing in a log view.
44///
45/// ```no_run
46/// use qframe::runtime::{Line, Process};
47///
48/// let mut lines = Vec::new();
49/// let outcome = Process::new("sh")
50///     .args(["-c", "echo ready"])
51///     .env("LC_ALL", "C")
52///     .run(&|| false, &mut |line| lines.push(line))?;
53/// assert_eq!(lines, vec![Line::Out("ready".to_owned())]);
54/// # Ok::<(), std::io::Error>(())
55/// ```
56///
57/// For output that is not read line by line, [`Process::collect`] keeps only the end of it.
58#[derive(Debug, Clone, PartialEq, Eq)]
59pub struct Process {
60    program: OsString,
61    args: Vec<OsString>,
62    dir: Option<PathBuf>,
63    env: Vec<(OsString, OsString)>,
64    pty: Option<(u16, u16)>,
65    no_stdin: bool,
66    /// The child gets nothing of the environment but what was named with `env`.
67    cleared: bool,
68}
69
70/// Where a line came from. With a pseudo-terminal both streams share one line, so only
71/// [`Line::Out`] appears.
72#[derive(Debug, Clone, PartialEq, Eq)]
73pub enum Line {
74    /// A line the child wrote to its standard output.
75    Out(String),
76    /// A line the child wrote to its standard error.
77    Err(String),
78}
79
80/// How the child ended.
81#[derive(Debug, Clone, PartialEq, Eq)]
82pub enum ProcessOutcome {
83    /// The child ran to its end; `code` is `None` when a signal ended it.
84    Finished {
85        /// The exit code, or `None` when a signal ended the child.
86        code: Option<i32>,
87    },
88    /// `cancel` turned true, so the child was killed and its pending output dropped.
89    Cancelled,
90}
91
92/// How much of a child's output to keep when it is [`collected`](Process::collect), and how long
93/// it may run.
94///
95/// It is built with [`Keep::bytes`] and narrowed with [`Keep::lines`] and [`Keep::limit`]. What
96/// is kept is the end of the output, never the beginning: a program that writes more than the
97/// bytes asked for has its oldest output dropped as the newest arrives, and
98/// [`Collected::trimmed`] says that it did.
99#[derive(Debug, Clone, Copy, PartialEq, Eq)]
100#[non_exhaustive]
101pub struct Keep {
102    bytes: usize,
103    lines: Option<usize>,
104    limit: Option<Duration>,
105}
106
107impl Keep {
108    /// Keeps at most `bytes` bytes of the output, the end of it. A program that writes more has
109    /// its oldest output dropped while it runs and never costs more memory than this plus one
110    /// read, and the cut falls between characters, so no character is half in the text. Zero
111    /// bytes is a way of asking for the outcome alone.
112    #[must_use]
113    pub fn bytes(bytes: usize) -> Self {
114        Self { bytes, lines: None, limit: None }
115    }
116
117    /// Keeps at most `lines` lines as well, the last of them. The line a program is still writing
118    /// counts as one, so output that ends without a newline keeps what was written.
119    #[must_use]
120    pub fn lines(mut self, lines: usize) -> Self {
121        self.lines = Some(lines);
122        self
123    }
124
125    /// Ends the program when `limit` is over, the way cancelling does: with
126    /// [`Process::no_stdin`] on Unix its whole process group goes with it, and
127    /// [`Collected::timed_out`] says that the time is what ended it.
128    #[must_use]
129    pub fn limit(mut self, limit: Duration) -> Self {
130        self.limit = Some(limit);
131        self
132    }
133}
134
135/// What a child wrote, kept to what [`Keep`] asked for, and how it ended.
136#[derive(Debug, Clone, PartialEq, Eq)]
137#[non_exhaustive]
138pub struct Collected {
139    /// The end of the child's output, both of its streams in the order it wrote them. Bytes that
140    /// are not UTF-8 became the replacement character, and the text begins on a character: the
141    /// cut the size asked for made never falls inside one.
142    pub text: String,
143    /// How the child ended: [`ProcessOutcome::Cancelled`] when `cancel` turned true, and
144    /// [`ProcessOutcome::Finished`] with no code when [`Keep::limit`] ended it.
145    pub outcome: ProcessOutcome,
146    /// Whether older output fell out to hold what [`Keep`] asked for.
147    pub trimmed: bool,
148    /// Whether [`Keep::limit`] ended the child.
149    pub timed_out: bool,
150    /// Whether `cancel` ended the child.
151    pub cancelled: bool,
152}
153
154impl Process {
155    /// A child process that runs `program`, with pipes and the application's own environment.
156    #[must_use]
157    pub fn new(program: impl Into<OsString>) -> Self {
158        Self {
159            program: program.into(),
160            args: Vec::new(),
161            dir: None,
162            env: Vec::new(),
163            pty: None,
164            no_stdin: false,
165            cleared: false,
166        }
167    }
168
169    /// Adds one argument.
170    #[must_use]
171    pub fn arg(mut self, arg: impl Into<OsString>) -> Self {
172        self.args.push(arg.into());
173        self
174    }
175
176    /// Adds several arguments, in order.
177    #[must_use]
178    pub fn args(mut self, args: impl IntoIterator<Item = impl Into<OsString>>) -> Self {
179        self.args.extend(args.into_iter().map(Into::into));
180        self
181    }
182
183    /// Runs the child in `dir` instead of the application's working directory.
184    #[must_use]
185    pub fn dir(mut self, dir: impl Into<PathBuf>) -> Self {
186        self.dir = Some(dir.into());
187        self
188    }
189
190    /// Sets one environment variable for the child. The rest of the environment is inherited,
191    /// and setting a variable the application already has replaces it for the child only.
192    #[must_use]
193    pub fn env(mut self, key: impl Into<OsString>, value: impl Into<OsString>) -> Self {
194        self.env.push((key.into(), value.into()));
195        self
196    }
197
198    /// Gives the child nothing of the application's environment: what it gets is what was set
199    /// with [`Process::env`] and nothing else, so a program sees what the application said it
200    /// would see rather than what the shell the application was started from happened to have.
201    ///
202    /// `PATH` is one of the variables that has to be named, since it is what the child looks
203    /// other programs up with, as in the example. `HOME`, `LANG` and the rest are just as gone.
204    ///
205    /// ```no_run
206    /// use qframe::runtime::{Line, Process, ProcessOutcome};
207    ///
208    /// let path = std::env::var("PATH").expect("a path to look other programs up with");
209    /// let mut lines = Vec::new();
210    /// let outcome = Process::new("sh")
211    ///     .args(["-c", "printf 'sadece bu'"])
212    ///     .clear_env()
213    ///     .env("PATH", path)
214    ///     .run(&|| false, &mut |line| lines.push(line))?;
215    /// assert_eq!(lines, vec![Line::Out("sadece bu".to_owned())]);
216    /// assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
217    /// # Ok::<(), std::io::Error>(())
218    /// ```
219    ///
220    /// [`Process::env`] still sets a variable here, but that is now all the child gets.
221    #[must_use]
222    pub fn clear_env(mut self) -> Self {
223        self.cleared = true;
224        self
225    }
226
227    /// Runs the child on a pseudo-terminal `cols` wide and `rows` tall, so programs that check
228    /// for a terminal draw their progress and colour. Standard input stays the application's
229    /// own unless [`Process::no_stdin`] is asked for, and the child keeps the controlling
230    /// terminal, which is what keeps a warm `sudo` ticket shared. Without this the child gets
231    /// pipes and sees no terminal.
232    ///
233    /// The child reads the size given here, not the real terminal's, so its progress bar fits
234    /// the space the application is going to draw it in.
235    #[must_use]
236    pub fn pty(mut self, cols: u16, rows: u16) -> Self {
237        self.pty = Some((cols, rows));
238        self
239    }
240
241    /// Gives the child no standard input: it reads an empty stream (`/dev/null`) instead of the
242    /// application's terminal. A program that asks a question then gets no answer rather than
243    /// the keys meant for the application, which it would otherwise take from under it.
244    ///
245    /// On Unix the child also starts in a process group of its own, so cancelling ends the
246    /// programs it started as well; see [`Process::run`]. It keeps the application's session
247    /// and controlling terminal, so a warm `sudo` ticket still applies. A program that reads
248    /// the terminal itself anyway, as `sudo` does to ask for a password, is stopped by the
249    /// system until it is cancelled, because its group is not the one the terminal belongs to:
250    /// warm the ticket first with a [`Handoff`](crate::runtime::Handoff) of `sudo -v`, or pass
251    /// `sudo -n` so it fails at once instead of asking.
252    #[must_use]
253    pub fn no_stdin(mut self) -> Self {
254        self.no_stdin = true;
255        self
256    }
257
258    /// Runs the child, handing every line to `on_line`, and returns how it ended.
259    ///
260    /// Lines arrive one by one, without their newline. A `\r` overwrites the line being built
261    /// rather than starting a new one, which is how progress bars are written, and the last
262    /// line is delivered even when the output does not end with a newline. Bytes that are not
263    /// UTF-8 become the replacement character instead of being dropped. A line longer than
264    /// 64 KiB is delivered in pieces of at most that size, cut between characters, so a program
265    /// that never writes a newline cannot make the reader hold all of its output.
266    ///
267    /// When the child writes faster than `on_line` takes its lines, the reading waits and the
268    /// child waits with it, rather than its output piling up in memory.
269    ///
270    /// `cancel` is asked between lines, and every few milliseconds while there is none, also
271    /// after the child has closed its output but keeps running; when it turns true the child is
272    /// killed, its pending output is dropped and the outcome is [`ProcessOutcome::Cancelled`].
273    ///
274    /// What cancelling kills depends on standard input. With [`Process::no_stdin`] on Unix, the
275    /// child runs in a process group of its own and the whole group is killed, so the programs
276    /// it started go with it (`podman` with `buildah` and the build's steps), except those that
277    /// moved to a group or session of their own. Without it the child shares the application's
278    /// standard input, which is the terminal: in a group of its own it would be stopped by the
279    /// system the first time it read from it, so it stays in the application's group and only
280    /// the child itself is killed; a program that started children of its own can leave them
281    /// running.
282    ///
283    /// Meant to be called inside a [`Task`](crate::runtime::Task), with `cancel` reading
284    /// [`TaskCx::is_cancelled`](crate::runtime::TaskCx::is_cancelled).
285    ///
286    /// # Errors
287    ///
288    /// Returns an I/O error when the child cannot be started, when a pseudo-terminal was asked
289    /// for and cannot be opened, or when a reading thread cannot be started.
290    pub fn run(self, cancel: &dyn Fn() -> bool, on_line: &mut dyn FnMut(Line)) -> io::Result<ProcessOutcome> {
291        self.run_inner(cancel, on_line, None)
292    }
293
294    /// Runs the child like [`Process::run`], and also hands every line a `\r` overwrites to
295    /// `on_overwritten` instead of dropping it: the frames of a progress bar, as `cargo`,
296    /// `pacman`, `curl` and `git` write them.
297    ///
298    /// A frame is the text built since the last line end or `\r`, delivered when the byte
299    /// after the `\r` shows that the line really is overwritten; `\r\n` and `\r\r\n` stay
300    /// plain line ends and give no frame, and an empty frame is not delivered. When the output
301    /// ends right after a `\r`, its last frame is delivered too. Colour codes and erase codes
302    /// such as `ESC [K` are passed on untouched. A frame comes tagged like a line: [`Line::Out`]
303    /// or [`Line::Err`] for the stream it was written to, and always [`Line::Out`] on a
304    /// pseudo-terminal. Lines and frames arrive in the order the child wrote them; `on_line`
305    /// receives exactly what [`Process::run`] would hand it.
306    ///
307    /// ```no_run
308    /// use qframe::runtime::{Line, Process};
309    ///
310    /// let (mut lines, mut frames) = (Vec::new(), Vec::new());
311    /// Process::new("sh").args(["-c", r"printf '10%\r50%\rdone\n'"]).run_with_overwritten(
312    ///     &|| false,
313    ///     &mut |line| lines.push(line),
314    ///     &mut |frame| frames.push(frame),
315    /// )?;
316    /// assert_eq!(lines, vec![Line::Out("done".to_owned())]);
317    /// assert_eq!(frames, vec![Line::Out("10%".to_owned()), Line::Out("50%".to_owned())]);
318    /// # Ok::<(), std::io::Error>(())
319    /// ```
320    ///
321    /// # Errors
322    ///
323    /// The same as [`Process::run`].
324    pub fn run_with_overwritten(
325        self,
326        cancel: &dyn Fn() -> bool,
327        on_line: &mut dyn FnMut(Line),
328        on_overwritten: &mut dyn FnMut(Line),
329    ) -> io::Result<ProcessOutcome> {
330        self.run_inner(cancel, on_line, Some(on_overwritten))
331    }
332
333    /// Runs the child; frames are read at all only when someone takes them, so a plain run
334    /// never queues them.
335    fn run_inner(
336        self,
337        cancel: &dyn Fn() -> bool,
338        on_line: &mut dyn FnMut(Line),
339        mut on_overwritten: Option<&mut dyn FnMut(Line)>,
340    ) -> io::Result<ProcessOutcome> {
341        let frames = on_overwritten.is_some();
342        // The group is what lets cancelling reach the child's own children; see `run`'s notes.
343        let group = self.no_stdin && cfg!(unix);
344        let command = self.command(group);
345        let (sender, receiver) = mpsc::sync_channel(QUEUE);
346        let mut child = match self.pty {
347            Some(size) => spawn_on_pty(command, size, &sender, frames, group)?,
348            None => spawn_on_pipes(command, &sender, frames, group)?,
349        };
350        // The readers hold the only remaining senders, so the channel ends when they do.
351        drop(sender);
352        loop {
353            if cancel() {
354                kill(&mut child, group);
355                return Ok(ProcessOutcome::Cancelled);
356            }
357            match receiver.recv_timeout(POLL) {
358                Ok(Sent::Line(line)) => on_line(line),
359                Ok(Sent::Overwritten(frame)) => {
360                    if let Some(on_overwritten) = on_overwritten.as_deref_mut() {
361                        on_overwritten(frame);
362                    }
363                }
364                Err(RecvTimeoutError::Timeout) => {}
365                Err(RecvTimeoutError::Disconnected) => break,
366            }
367        }
368        // Its streams are closed, but the child may still be running: it closed them itself, or
369        // they were handed to a program of its own. Waiting keeps asking `cancel`.
370        loop {
371            if let Some(status) = child.try_wait()? {
372                return Ok(ProcessOutcome::Finished { code: status.code() });
373            }
374            if cancel() {
375                kill(&mut child, group);
376                return Ok(ProcessOutcome::Cancelled);
377            }
378            std::thread::sleep(POLL);
379        }
380    }
381
382    /// Runs the child, keeping only the end of what it writes, and gives that back as one piece
383    /// of text.
384    ///
385    /// Both of the child's streams land on one pipe, so the text is in the order the child wrote
386    /// it and a failure stays among the rest of the output; with [`Process::pty`] both land on
387    /// the terminal itself, as they always do there. Nothing is delivered line by line: the point
388    /// is to hold the end of the output, so a program that prints megabytes, a build's whole log,
389    /// a flood of progress, costs no more memory than [`Keep`] asks for however long it writes.
390    ///
391    /// [`Keep::limit`] ends the child the way `cancel` does, its process group too where
392    /// [`Process::no_stdin`] asked for one, and both give back the output kept so far in bounded
393    /// time.
394    ///
395    /// Meant to be called inside a [`Task`](crate::runtime::Task), with `cancel` reading
396    /// [`TaskCx::is_cancelled`](crate::runtime::TaskCx::is_cancelled).
397    ///
398    /// ```no_run
399    /// use std::time::Duration;
400    /// use qframe::runtime::{Keep, Process, ProcessOutcome};
401    ///
402    /// let collected = Process::new("sh")
403    ///     .args(["-c", "echo ready"])
404    ///     .env("LC_ALL", "C")
405    ///     .collect(Keep::bytes(4096).lines(20).limit(Duration::from_secs(30)), &|| false)?;
406    /// assert_eq!(collected.text, "ready\n");
407    /// assert_eq!(collected.outcome, ProcessOutcome::Finished { code: Some(0) });
408    /// # Ok::<(), std::io::Error>(())
409    /// ```
410    ///
411    /// # Errors
412    ///
413    /// The same as [`Process::run`].
414    pub fn collect(self, keep: Keep, cancel: &dyn Fn() -> bool) -> io::Result<Collected> {
415        let Keep { bytes, lines, limit } = keep;
416        let group = self.no_stdin && cfg!(unix);
417        let command = self.command(group);
418        let merged = Arc::new(Mutex::new(Merged { tail: Tail::new(bytes, lines), open: true }));
419        let mut child = match self.pty {
420            Some(size) => spawn_merged_on_pty(command, size, group, &merged)?,
421            None => spawn_merged(command, group, &merged)?,
422        };
423        let deadline = limit.map(|limit| Instant::now() + limit);
424        // The stream ending is not the child ending: it may have closed its output and kept
425        // running. The loop waits in short steps and asks both again in each of them, so nothing
426        // waits for a child that will not answer.
427        let (outcome, timed_out, cancelled) = loop {
428            if cancel() {
429                kill(&mut child, group);
430                break (ProcessOutcome::Cancelled, false, true);
431            }
432            if deadline.is_some_and(|deadline| Instant::now() >= deadline) {
433                kill(&mut child, group);
434                break (ProcessOutcome::Finished { code: None }, true, false);
435            }
436            if !lock(&merged).open
437                && let Some(status) = child.try_wait()?
438            {
439                break (ProcessOutcome::Finished { code: status.code() }, false, false);
440            }
441            std::thread::sleep(POLL);
442        };
443        let mut merged = lock(&merged);
444        Ok(Collected { text: merged.tail.text(), trimmed: merged.tail.trimmed, outcome, timed_out, cancelled })
445    }
446
447    /// The command line with everything but the streams: the group the child runs in, where it
448    /// runs, and the environment it is given.
449    fn command(&self, group: bool) -> Command {
450        let mut command = Command::new(&self.program);
451        command.args(&self.args);
452        if self.no_stdin {
453            command.stdin(Stdio::null());
454        } else {
455            command.stdin(Stdio::inherit());
456        }
457        #[cfg(unix)]
458        if group {
459            use std::os::unix::process::CommandExt;
460            command.process_group(0);
461        }
462        #[cfg(not(unix))]
463        let _ = group;
464        if let Some(dir) = &self.dir {
465            command.current_dir(dir);
466        }
467        if self.cleared {
468            command.env_clear();
469        }
470        for (key, value) in &self.env {
471            command.env(key, value);
472        }
473        command
474    }
475}
476
477/// What a reading thread hands to the loop in [`Process::run_inner`]. One channel carries both
478/// so lines and frames keep the order the child wrote them in.
479enum Sent {
480    Line(Line),
481    Overwritten(Line),
482}
483
484/// Kills the child, and with `group` every process still in its process group, and waits for the
485/// child so it leaves nothing behind.
486fn kill(child: &mut Child, group: bool) {
487    #[cfg(unix)]
488    if group {
489        let leader = rustix::process::Pid::from_child(child);
490        // Fails only when nothing is left in the group; the child itself is killed below in any
491        // case.
492        let _ = rustix::process::kill_process_group(leader, rustix::process::Signal::KILL);
493    }
494    #[cfg(not(unix))]
495    let _ = group;
496    // The kill fails only on a child that ended by itself, which is what was wanted, and the
497    // wait collects it rather than being asked anything. This function has no answer to give:
498    // its callers have already decided the child is to go.
499    let _ = child.kill();
500    let _ = child.wait();
501}
502
503/// Starts the child with a pipe per stream and a reading thread for each, so a failure stays
504/// recognisable as one.
505fn spawn_on_pipes(mut command: Command, sender: &SyncSender<Sent>, frames: bool, group: bool) -> io::Result<Child> {
506    command.stdout(Stdio::piped()).stderr(Stdio::piped());
507    let mut child = command.spawn()?;
508    drop(command);
509    let taken = child.stdout.take().zip(child.stderr.take());
510    let started = match taken {
511        Some((out, err)) => spawn_reader("out", out, Line::Out, frames, sender.clone())
512            .and_then(|()| spawn_reader("err", err, Line::Err, frames, sender.clone())),
513        None => Err(io::Error::other("the child was started without its pipes")),
514    };
515    match started {
516        Ok(()) => Ok(child),
517        Err(error) => {
518            kill(&mut child, group);
519            Err(error)
520        }
521    }
522}
523
524/// Starts the child on a pseudo-terminal of the given size, reading the one stream both of its
525/// streams land on.
526#[cfg(unix)]
527fn spawn_on_pty(
528    mut command: Command,
529    size: (u16, u16),
530    sender: &SyncSender<Sent>,
531    frames: bool,
532    group: bool,
533) -> io::Result<Child> {
534    let terminal = open_pty(&mut command, size)?;
535    let mut child = command.spawn()?;
536    // The command holds the child's side of the terminal until it is dropped, and while it is
537    // open the reading side never reaches its end of file.
538    drop(command);
539    match spawn_reader("pty", terminal, Line::Out, frames, sender.clone()) {
540        Ok(()) => Ok(child),
541        Err(error) => {
542            kill(&mut child, group);
543            Err(error)
544        }
545    }
546}
547
548/// Without Unix there is no pseudo-terminal to open, so the caller is told instead of being
549/// given a child that quietly sees no terminal.
550#[cfg(not(unix))]
551fn spawn_on_pty(
552    mut command: Command,
553    size: (u16, u16),
554    _sender: &SyncSender<Sent>,
555    _frames: bool,
556    _group: bool,
557) -> io::Result<Child> {
558    open_pty(&mut command, size)
559}
560
561/// Starts the child with one pipe carrying both of its streams, so what it writes is read in the
562/// order it wrote it, and a thread that keeps only the end of it.
563fn spawn_merged(mut command: Command, group: bool, merged: &Shared) -> io::Result<Child> {
564    let (reader, writer) = io::pipe()?;
565    command.stdout(writer.try_clone()?).stderr(writer);
566    let mut child = command.spawn()?;
567    // The command holds the write end of the pipe until it is dropped, and while it is open the
568    // reading end never reaches its end of file.
569    drop(command);
570    match spawn_collector(reader, Arc::clone(merged)) {
571        Ok(()) => Ok(child),
572        Err(error) => {
573            kill(&mut child, group);
574            Err(error)
575        }
576    }
577}
578
579/// Starts the child on a pseudo-terminal, where both of its streams land on the one stream
580/// already, and keeps the end of that.
581#[cfg(unix)]
582fn spawn_merged_on_pty(mut command: Command, size: (u16, u16), group: bool, merged: &Shared) -> io::Result<Child> {
583    let terminal = open_pty(&mut command, size)?;
584    let mut child = command.spawn()?;
585    // The command holds the child's side of the terminal until it is dropped, and while it is
586    // open the reading side never reaches its end of file.
587    drop(command);
588    match spawn_collector(terminal, Arc::clone(merged)) {
589        Ok(()) => Ok(child),
590        Err(error) => {
591            kill(&mut child, group);
592            Err(error)
593        }
594    }
595}
596
597/// Without Unix there is no pseudo-terminal, as in [`spawn_on_pty`].
598#[cfg(not(unix))]
599fn spawn_merged_on_pty(mut command: Command, size: (u16, u16), _group: bool, _merged: &Shared) -> io::Result<Child> {
600    open_pty(&mut command, size)
601}
602
603/// Opens a pseudo-terminal `cols` wide and `rows` tall, gives the child both of its streams on
604/// it and hands back the one stream they land on.
605#[cfg(unix)]
606fn open_pty(command: &mut Command, (cols, rows): (u16, u16)) -> io::Result<std::fs::File> {
607    use std::fs::File;
608    use std::os::fd::OwnedFd;
609
610    use rustix::fs::{Mode, OFlags};
611    use rustix::io::{FdFlags, fcntl_setfd};
612    use rustix::pty::{OpenptFlags, grantpt, openpt, ptsname, unlockpt};
613    use rustix::termios::{Winsize, tcsetwinsize};
614
615    // Our own side must not reach the child: it would then hold the terminal open itself and
616    // reading would never end. Where the flag can be given at once, no program another thread
617    // starts in between can inherit it either; elsewhere it is set right after.
618    #[cfg(any(target_os = "linux", target_os = "android", target_os = "freebsd", target_os = "netbsd"))]
619    let flags = OpenptFlags::RDWR | OpenptFlags::NOCTTY | OpenptFlags::CLOEXEC;
620    #[cfg(not(any(target_os = "linux", target_os = "android", target_os = "freebsd", target_os = "netbsd")))]
621    let flags = OpenptFlags::RDWR | OpenptFlags::NOCTTY;
622    let controller = openpt(flags)?;
623    fcntl_setfd(&controller, FdFlags::CLOEXEC)?;
624    grantpt(&controller)?;
625    unlockpt(&controller)?;
626    tcsetwinsize(&controller, Winsize { ws_row: rows, ws_col: cols, ws_xpixel: 0, ws_ypixel: 0 })?;
627    let name = ptsname(&controller, Vec::new())?;
628    // `NOCTTY` leaves the child on the application's controlling terminal instead of making this
629    // pseudo-terminal the controlling one; a `sudo` ticket is held per controlling terminal, so
630    // taking it away would ask for the password again.
631    // `CLOEXEC` keeps this descriptor itself out of the child, which gets the terminal only as
632    // its standard output and error, and out of any program another thread starts meanwhile:
633    // a stray copy would keep the terminal open after the child closed its streams.
634    let device: OwnedFd = rustix::fs::open(name, OFlags::RDWR | OFlags::NOCTTY | OFlags::CLOEXEC, Mode::empty())?;
635    command.stdout(Stdio::from(device.try_clone()?)).stderr(Stdio::from(device));
636    Ok(File::from(controller))
637}
638
639/// Without Unix there is no pseudo-terminal to open, so the caller is told instead of being
640/// given a child that quietly sees no terminal.
641#[cfg(not(unix))]
642fn open_pty(_command: &mut Command, _size: (u16, u16)) -> io::Result<std::fs::File> {
643    Err(io::Error::new(io::ErrorKind::Unsupported, "a pseudo-terminal needs a Unix system"))
644}
645
646/// Reads `source` on its own thread, sending one message per line, and with `frames` one per
647/// overwritten frame, until the stream ends or the receiver is gone.
648fn spawn_reader(
649    name: &str,
650    source: impl Read + Send + 'static,
651    tag: fn(String) -> Line,
652    frames: bool,
653    sender: SyncSender<Sent>,
654) -> io::Result<()> {
655    std::thread::Builder::new()
656        .name(format!("quvyta-process-{name}"))
657        .spawn(move || read_lines(source, tag, frames, &sender))
658        .map(|_| ())
659}
660
661/// Sends every line of `source` as a message, stopping as soon as the receiver is gone.
662fn read_lines(mut source: impl Read, tag: fn(String) -> Line, frames: bool, sender: &SyncSender<Sent>) {
663    let mut chunk = [0_u8; CHUNK];
664    let mut lines = Lines::default();
665    // Both closures send; a failed send from either means nobody listens any more.
666    let listening = std::cell::Cell::new(true);
667    let mut on_line = |line| listening.set(listening.get() && sender.send(Sent::Line(tag(line))).is_ok());
668    let mut on_frame = |frame| listening.set(listening.get() && sender.send(Sent::Overwritten(tag(frame))).is_ok());
669    loop {
670        match source.read(&mut chunk) {
671            Ok(0) => break,
672            Ok(count) => {
673                lines.feed_keeping(&chunk[..count], &mut on_line, frames.then_some(&mut on_frame));
674                if !listening.get() {
675                    return;
676                }
677            }
678            Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
679            // A pseudo-terminal answers with an I/O error once the child's side is gone, and a
680            // broken pipe says the same thing; both are the end of the stream.
681            Err(_) => break,
682        }
683    }
684    lines.finish_keeping(&mut on_line, frames.then_some(&mut on_frame));
685}
686
687/// Starts a thread that keeps the end of `source` in `merged`, which the caller reads as the
688/// child writes.
689fn spawn_collector(source: impl Read + Send + 'static, merged: Shared) -> io::Result<()> {
690    std::thread::Builder::new()
691        .name("quvyta-process-merged".to_owned())
692        .spawn(move || read_tail(source, &merged))
693        .map(|_| ())
694}
695
696/// The one stream both of a collected child's streams land on: what is worth keeping of it, and
697/// whether the thread reading it still holds it. The thread fills this in and the caller reads
698/// it, so the caller can take what has been kept the moment it stops the child.
699#[derive(Debug)]
700struct Merged {
701    tail: Tail,
702    open: bool,
703}
704
705/// The one stream, shared by the thread that reads it and the caller that waits for it.
706type Shared = Arc<Mutex<Merged>>;
707
708fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
709    mutex.lock().unwrap_or_else(PoisonError::into_inner)
710}
711
712/// Keeps what one read brought in `merged` until the stream ends, and lets go of it only then,
713/// so the caller never hears the end before the last of the output is in.
714fn read_tail(mut source: impl Read, merged: &Shared) {
715    let mut chunk = [0_u8; CHUNK];
716    loop {
717        match source.read(&mut chunk) {
718            Ok(0) => break,
719            Ok(count) => lock(merged).tail.feed(&chunk[..count]),
720            Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
721            // A pseudo-terminal answers with an I/O error once the child's side is gone, and a
722            // broken pipe says the same thing; both are the end of the stream.
723            Err(_) => break,
724        }
725    }
726    lock(merged).open = false;
727}
728
729/// The end of a byte stream: at most the bytes asked for, and at most the lines asked for when
730/// there are any, the oldest dropped as the newest arrives. Nothing is held beyond this and one
731/// read, so a program that writes for an hour costs what it is allowed to.
732#[derive(Debug)]
733struct Tail {
734    bytes: VecDeque<u8>,
735    /// How many line ends the kept bytes hold.
736    lines: usize,
737    /// Whether the last byte read was a line end, so nothing is being written after it.
738    ends_line: bool,
739    /// How many bytes are kept.
740    kept: usize,
741    /// How many lines are kept, when the caller named a count.
742    lines_kept: Option<usize>,
743    /// Whether anything was dropped for being older than what is kept.
744    trimmed: bool,
745}
746
747impl Tail {
748    /// A tail keeping `bytes` bytes, and `lines` lines as well when the caller named a count.
749    fn new(bytes: usize, lines: Option<usize>) -> Self {
750        Self { bytes: VecDeque::new(), lines: 0, ends_line: true, kept: bytes, lines_kept: lines, trimmed: false }
751    }
752
753    /// Adds what one read brought, dropping the oldest until what is kept is within what was
754    /// asked for.
755    fn feed(&mut self, bytes: &[u8]) {
756        let Some(&last) = bytes.last() else {
757            return;
758        };
759        self.ends_line = last == b'\n';
760        self.lines += bytes.iter().filter(|&&byte| byte == b'\n').count();
761        self.bytes.extend(bytes);
762        if let Some(kept) = self.lines_kept {
763            for _ in 0..self.counting().saturating_sub(kept) {
764                self.drop_oldest_line();
765            }
766        }
767        for _ in 0..self.bytes.len().saturating_sub(self.kept) {
768            if self.bytes.pop_front() == Some(b'\n') {
769                self.lines -= 1;
770            }
771            self.trimmed = true;
772        }
773    }
774
775    /// How many lines the kept bytes hold: the line ends in them, and the line being written,
776    /// which has no end yet.
777    fn counting(&self) -> usize {
778        self.lines + usize::from(!self.ends_line && !self.bytes.is_empty())
779    }
780
781    /// Drops the oldest line, whatever is left of it and its end with it.
782    fn drop_oldest_line(&mut self) {
783        self.trimmed = true;
784        while let Some(byte) = self.bytes.pop_front() {
785            if byte == b'\n' {
786                self.lines -= 1;
787                break;
788            }
789        }
790    }
791
792    /// The kept bytes as text, from a character: a cut made for the size asked for can leave the
793    /// first half of a character behind, and half a character is no text.
794    fn text(&mut self) -> String {
795        let bytes = self.bytes.make_contiguous();
796        let start = bytes.iter().copied().take(3).take_while(|byte| byte & 0b1100_0000 == 0b1000_0000).count();
797        String::from_utf8_lossy(&bytes[start..]).into_owned()
798    }
799}
800
801/// Splits a byte stream into lines, letting `\r` overwrite the line being built. What it
802/// overwrites is dropped, or handed to a second callback by [`Lines::feed_keeping`].
803#[derive(Debug, Default)]
804pub(super) struct Lines {
805    buffer: Vec<u8>,
806    /// A `\r` was read and it is not yet known whether a `\n` follows it.
807    pending_return: bool,
808}
809
810impl Lines {
811    /// Feeds `bytes`, calling `emit` once per finished line.
812    pub(super) fn feed(&mut self, bytes: &[u8], emit: &mut impl FnMut(String)) {
813        self.feed_keeping(bytes, emit, None);
814    }
815
816    /// Feeds `bytes` like [`Lines::feed`], also handing each non-empty frame a `\r`
817    /// overwrites to `overwritten` when it is given.
818    pub(super) fn feed_keeping(
819        &mut self,
820        bytes: &[u8],
821        emit: &mut impl FnMut(String),
822        mut overwritten: Option<&mut dyn FnMut(String)>,
823    ) {
824        for &byte in bytes {
825            if self.pending_return {
826                // A terminal ends its lines with `\r\n`, so a `\r` right before a newline ends
827                // the line rather than overwriting it. A second `\r` changes nothing, as on a
828                // screen: a program's own `\r\n` arrives as `\r\r\n` from a pseudo-terminal.
829                match byte {
830                    b'\r' => continue,
831                    b'\n' => {
832                        self.pending_return = false;
833                        emit(self.take());
834                        continue;
835                    }
836                    _ => {
837                        self.pending_return = false;
838                        self.overwrite(&mut overwritten);
839                    }
840                }
841            }
842            match byte {
843                b'\r' => self.pending_return = true,
844                b'\n' => emit(self.take()),
845                _ => {
846                    self.buffer.push(byte);
847                    if self.buffer.len() >= MAX_LINE {
848                        self.emit_piece(emit);
849                    }
850                }
851            }
852        }
853    }
854
855    /// Delivers the full buffer as a line of its own, keeping back the start of a character
856    /// that is not complete yet so no character is cut in two.
857    fn emit_piece(&mut self, emit: &mut impl FnMut(String)) {
858        // A character is at most four bytes, so only the last three can start one that is not
859        // complete yet. Anything else that is not UTF-8 is replaced as usual.
860        let len = self.buffer.len();
861        let mut cut = len;
862        for back in 1..=len.min(3) {
863            let byte = self.buffer[len - back];
864            if byte & 0b1100_0000 != 0b1000_0000 {
865                let width = match byte {
866                    0xc0..=0xdf => 2,
867                    0xe0..=0xef => 3,
868                    0xf0..=0xf7 => 4,
869                    _ => 1,
870                };
871                if width > back {
872                    cut = len - back;
873                }
874                break;
875            }
876        }
877        let rest = self.buffer.split_off(cut);
878        emit(self.take());
879        self.buffer = rest;
880    }
881
882    /// Drops the line a `\r` overwrites, or hands it to `overwritten` when there is one.
883    fn overwrite(&mut self, overwritten: &mut Option<&mut dyn FnMut(String)>) {
884        match overwritten {
885            Some(overwritten) if !self.buffer.is_empty() => overwritten(self.take()),
886            _ => self.buffer.clear(),
887        }
888    }
889
890    /// Delivers the last line when the stream ended without a newline.
891    pub(super) fn finish(&mut self, emit: &mut impl FnMut(String)) {
892        self.finish_keeping(emit, None);
893    }
894
895    /// Ends the stream like [`Lines::finish`]; a last line followed by a `\r` goes to
896    /// `overwritten` when it is given.
897    pub(super) fn finish_keeping(
898        &mut self,
899        emit: &mut impl FnMut(String),
900        mut overwritten: Option<&mut dyn FnMut(String)>,
901    ) {
902        if self.pending_return {
903            // The line was overwritten and nothing was written in its place.
904            self.overwrite(&mut overwritten);
905            self.pending_return = false;
906        }
907        if !self.buffer.is_empty() {
908            emit(self.take());
909        }
910    }
911
912    /// The line built so far, with anything that is not UTF-8 replaced rather than dropped.
913    fn take(&mut self) -> String {
914        let line = String::from_utf8_lossy(&self.buffer).into_owned();
915        self.buffer.clear();
916        line
917    }
918}
919
920#[cfg(test)]
921mod tests {
922    use std::sync::atomic::{AtomicUsize, Ordering};
923    use std::time::Duration;
924
925    use super::{Collected, Keep, Line, Lines, MAX_LINE, Process, ProcessOutcome, Tail};
926
927    /// Runs a shell command to its end and returns its lines and outcome.
928    fn shell(script: &str) -> (Vec<Line>, ProcessOutcome) {
929        run(Process::new("sh").args(["-c", script]))
930    }
931
932    /// Runs `process` to its end, never cancelling.
933    fn run(process: Process) -> (Vec<Line>, ProcessOutcome) {
934        let mut lines = Vec::new();
935        let outcome = process.run(&|| false, &mut |line| lines.push(line)).expect("the shell starts");
936        (lines, outcome)
937    }
938
939    /// Runs `process` to its end and keeps only the end of what it wrote, never cancelling.
940    fn collect(process: Process, keep: Keep) -> Collected {
941        process.collect(keep, &|| false).expect("the shell starts")
942    }
943
944    #[test]
945    fn keeps_the_two_streams_apart_and_reports_the_exit_code() {
946        let (lines, outcome) = shell("echo bir; echo iki >&2; exit 3");
947        assert_eq!(lines.len(), 2, "{lines:?}");
948        assert!(lines.contains(&Line::Out("bir".to_owned())), "{lines:?}");
949        assert!(lines.contains(&Line::Err("iki".to_owned())), "{lines:?}");
950        assert_eq!(outcome, ProcessOutcome::Finished { code: Some(3) });
951    }
952
953    #[test]
954    fn delivers_the_last_line_without_a_newline() {
955        let (lines, outcome) = shell("printf 'son satir'");
956        assert_eq!(lines, vec![Line::Out("son satir".to_owned())]);
957        assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
958    }
959
960    #[test]
961    fn carriage_returns_collapse_into_one_line() {
962        let (lines, _) = shell(r"printf 'a\rbb\rccc\n'");
963        assert_eq!(lines, vec![Line::Out("ccc".to_owned())]);
964    }
965
966    #[test]
967    fn invalid_utf8_becomes_the_replacement_character() {
968        let (lines, _) = shell(r"printf 'a\377b\n'");
969        assert_eq!(lines, vec![Line::Out("a\u{fffd}b".to_owned())]);
970    }
971
972    #[test]
973    fn the_environment_is_inherited_and_one_variable_can_be_replaced() {
974        let (lines, _) = shell("echo ${PATH:+inherited}");
975        assert_eq!(lines, vec![Line::Out("inherited".to_owned())]);
976        let (lines, _) = run(Process::new("sh").args(["-c", "echo $LC_ALL"]).env("LC_ALL", "C"));
977        assert_eq!(lines, vec![Line::Out("C".to_owned())]);
978    }
979
980    #[test]
981    fn a_cleared_environment_gives_the_child_only_what_it_was_told() {
982        // A test process cannot set a variable for itself, since `set_var` is unsafe, so the
983        // home it really runs with is what the first run reads.
984        let script = r#"printf '%s' "${HOME-unset}""#;
985        let (lines, _) = shell(script);
986        assert_ne!(lines, vec![Line::Out("unset".to_owned())], "the test process really has a home");
987        let path = std::env::var("PATH").expect("the test process was started with a path");
988        let (lines, _) = run(Process::new("sh").args(["-c", script]).clear_env().env("PATH", path.clone()));
989        assert_eq!(lines, vec![Line::Out("unset".to_owned())], "nothing is inherited");
990        // What was named is what the child gets, the path included: it looks other programs up
991        // with that.
992        let (lines, _) = run(Process::new("sh")
993            .args(["-c", r#"printf '%s' "${PATH-unset}""#])
994            .clear_env()
995            .env("PATH", path.clone()));
996        assert_eq!(lines, vec![Line::Out(path)]);
997        let (lines, _) =
998            run(Process::new("sh").args(["-c", r#"printf '%s' "${EV-unset}""#]).clear_env().env("EV", "1"));
999        assert_eq!(lines, vec![Line::Out("1".to_owned())], "a variable that was given arrives");
1000    }
1001
1002    #[test]
1003    fn runs_in_the_directory_it_is_given() {
1004        let (lines, _) = run(Process::new("sh").args(["-c", "pwd"]).dir("/"));
1005        assert_eq!(lines, vec![Line::Out("/".to_owned())]);
1006    }
1007
1008    #[test]
1009    fn cancelling_kills_a_long_running_child() {
1010        let seen = AtomicUsize::new(0);
1011        let outcome = Process::new("sh")
1012            .args(["-c", "while true; do echo tik; sleep 0.05; done"])
1013            .run(&|| seen.load(Ordering::Relaxed) > 0, &mut |line| {
1014                assert_eq!(line, Line::Out("tik".to_owned()));
1015                seen.fetch_add(1, Ordering::Relaxed);
1016            })
1017            .expect("the shell starts");
1018        assert_eq!(outcome, ProcessOutcome::Cancelled);
1019        assert!(seen.load(Ordering::Relaxed) > 0);
1020    }
1021
1022    #[test]
1023    fn a_missing_program_is_an_error_and_not_a_panic() {
1024        let error = Process::new("quvyta-no-such-program")
1025            .run(&|| false, &mut |_| unreachable!("a missing program writes nothing"))
1026            .expect_err("a missing program cannot run");
1027        assert_eq!(error.kind(), std::io::ErrorKind::NotFound);
1028    }
1029
1030    #[cfg(unix)]
1031    #[test]
1032    fn on_a_pseudo_terminal_the_child_sees_a_terminal_of_the_size_we_gave() {
1033        // `stty` reads its standard input, which is the application's own; reading the size the
1034        // child was given means asking about the stream it writes to.
1035        let (lines, outcome) = run(Process::new("sh").args(["-c", "test -t 1 && stty size <&1"]).pty(100, 24));
1036        assert_eq!(lines, vec![Line::Out("24 100".to_owned())]);
1037        assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
1038    }
1039
1040    #[cfg(unix)]
1041    #[test]
1042    fn on_a_pseudo_terminal_both_streams_arrive_as_output() {
1043        let (lines, outcome) = run(Process::new("sh").args(["-c", "echo bir; echo iki >&2"]).pty(80, 24));
1044        assert_eq!(lines, vec![Line::Out("bir".to_owned()), Line::Out("iki".to_owned())]);
1045        assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
1046    }
1047
1048    /// The process group, session and controlling terminal in a line of `/proc/<pid>/stat`.
1049    #[cfg(target_os = "linux")]
1050    fn stat_ids(stat: &str) -> [String; 3] {
1051        // The command name may hold spaces; after its closing parenthesis come state, parent,
1052        // group, session and terminal.
1053        let fields: Vec<&str> = stat[stat.rfind(')').expect("name") + 2..].split(' ').collect();
1054        [fields[2], fields[3], fields[4]].map(str::to_owned)
1055    }
1056
1057    #[cfg(target_os = "linux")]
1058    #[test]
1059    fn without_stdin_the_child_reads_an_empty_stream() {
1060        let script = r#"readlink /proc/$$/fd/0; read answer; echo "read $?""#;
1061        for process in [Process::new("sh").args(["-c", script]), Process::new("sh").args(["-c", script]).pty(80, 24)] {
1062            let (lines, outcome) = run(process.no_stdin());
1063            assert_eq!(lines, vec![Line::Out("/dev/null".to_owned()), Line::Out("read 1".to_owned())]);
1064            assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
1065        }
1066    }
1067
1068    #[cfg(target_os = "linux")]
1069    #[test]
1070    fn only_a_child_without_stdin_gets_a_group_of_its_own_and_it_keeps_the_session() {
1071        let script = "cat /proc/$$/stat";
1072        let ours = stat_ids(&std::fs::read_to_string("/proc/self/stat").expect("stat"));
1073        let ids = |process: Process| {
1074            let (lines, _) = run(process);
1075            let [Line::Out(stat)] = &lines[..] else { panic!("one line: {lines:?}") };
1076            stat_ids(stat)
1077        };
1078        let shared = ids(Process::new("sh").args(["-c", script]));
1079        assert_eq!(shared, ours, "a child reading the terminal stays in the application's group");
1080        for process in [Process::new("sh").args(["-c", script]), Process::new("sh").args(["-c", script]).pty(80, 24)] {
1081            let [group, session, terminal] = ids(process.no_stdin());
1082            assert_ne!(group, ours[0], "a group of its own");
1083            // `sudo` keeps its ticket per controlling terminal and session, so both must stay.
1084            assert_eq!(session, ours[1], "the application's session");
1085            assert_eq!(terminal, ours[2], "the application's controlling terminal");
1086        }
1087    }
1088
1089    /// Whether `pid` has ended: gone, or ended and waiting for its parent to collect it.
1090    #[cfg(target_os = "linux")]
1091    fn ended(pid: &str) -> bool {
1092        std::fs::read_to_string(format!("/proc/{pid}/stat"))
1093            .map_or(true, |stat| stat[stat.rfind(')').expect("name") + 2..].starts_with('Z'))
1094    }
1095
1096    #[cfg(target_os = "linux")]
1097    #[test]
1098    fn cancelling_a_child_without_stdin_ends_the_programs_it_started() {
1099        for pty in [false, true] {
1100            let seen = std::cell::RefCell::new(Vec::new());
1101            let process = Process::new("sh").args(["-c", "sleep 60 & echo $!; sleep 60 & echo $!; wait"]).no_stdin();
1102            let process = if pty { process.pty(80, 24) } else { process };
1103            let outcome = process
1104                .run(&|| seen.borrow().len() == 2, &mut |line| match line {
1105                    Line::Out(pid) => seen.borrow_mut().push(pid),
1106                    Line::Err(text) => panic!("nothing on standard error: {text}"),
1107                })
1108                .expect("the shell starts");
1109            assert_eq!(outcome, ProcessOutcome::Cancelled);
1110            let pids = seen.into_inner();
1111            let started = std::time::Instant::now();
1112            while !pids.iter().all(|pid| ended(pid)) {
1113                assert!(started.elapsed() < std::time::Duration::from_secs(20), "still running: {pids:?} (pty {pty})");
1114                std::thread::sleep(std::time::Duration::from_millis(20));
1115            }
1116        }
1117    }
1118
1119    #[cfg(target_os = "linux")]
1120    #[test]
1121    fn cancelling_a_child_that_shares_stdin_ends_only_the_child() {
1122        // Documented: without `no_stdin` the child stays in the application's group, and its
1123        // own children outlive it. They are ended here by hand so the test leaves nothing behind.
1124        let seen = std::cell::RefCell::new(Vec::new());
1125        let outcome = Process::new("sh")
1126            .args(["-c", "sleep 60 & echo $!; wait"])
1127            .run(&|| seen.borrow().len() == 1, &mut |line| {
1128                if let Line::Out(pid) = line {
1129                    seen.borrow_mut().push(pid);
1130                }
1131            })
1132            .expect("the shell starts");
1133        assert_eq!(outcome, ProcessOutcome::Cancelled);
1134        let pid = seen.into_inner().remove(0);
1135        std::thread::sleep(std::time::Duration::from_millis(200));
1136        let survived = !ended(&pid);
1137        let raw: i32 = pid.parse().expect("a process id");
1138        if let Some(pid) = rustix::process::Pid::from_raw(raw) {
1139            let _ = rustix::process::kill_process(pid, rustix::process::Signal::KILL);
1140        }
1141        assert!(survived, "the grandchild outlives a cancel of the child");
1142    }
1143
1144    #[test]
1145    fn a_line_ended_twice_by_a_return_is_kept() {
1146        // A program that ends its own lines with `\r\n` writes `\r\r\n` on a pseudo-terminal,
1147        // which turns every `\n` into `\r\n`. The line is on the screen, so it is not lost.
1148        let mut lines = Lines::default();
1149        let mut seen = Vec::new();
1150        lines.feed(b"hazir\r\r\nbitti\r\r\r\n", &mut |line| seen.push(line));
1151        assert_eq!(seen, vec!["hazir".to_owned(), "bitti".to_owned()]);
1152    }
1153
1154    #[cfg(unix)]
1155    #[test]
1156    fn a_pseudo_terminal_line_ended_by_the_program_itself_arrives_whole() {
1157        let (lines, _) = run(Process::new("sh").args(["-c", r"printf 'bir\r\niki\r\n'"]).pty(80, 24));
1158        assert_eq!(lines, vec![Line::Out("bir".to_owned()), Line::Out("iki".to_owned())]);
1159    }
1160
1161    #[test]
1162    fn a_line_without_an_end_is_delivered_in_pieces_of_bounded_size() {
1163        // A program that never writes a newline must not make the reader hold all of it.
1164        let (lines, _) = shell("head -c 300000 /dev/zero | tr '\\0' a");
1165        let total: usize = lines
1166            .iter()
1167            .map(|line| match line {
1168                Line::Out(text) => {
1169                    assert!(text.len() <= MAX_LINE, "a piece of {} bytes", text.len());
1170                    assert!(text.bytes().all(|byte| byte == b'a'));
1171                    text.len()
1172                }
1173                Line::Err(text) => panic!("nothing was written to standard error: {text}"),
1174            })
1175            .sum();
1176        assert_eq!(total, 300_000, "nothing is lost between the pieces");
1177    }
1178
1179    #[test]
1180    fn a_long_line_is_never_cut_inside_a_character() {
1181        let mut lines = Lines::default();
1182        let mut seen = Vec::new();
1183        // One byte of padding, so the two-byte `ç` straddles every piece boundary.
1184        let mut text = vec![b'a'];
1185        for _ in 0..MAX_LINE {
1186            text.extend_from_slice("ç".as_bytes());
1187        }
1188        lines.feed(&text, &mut |line| seen.push(line));
1189        lines.finish(&mut |line| seen.push(line));
1190        assert!(seen.len() > 1, "the line was split");
1191        assert!(seen.iter().all(|line| !line.contains('\u{fffd}')), "no character was cut in two");
1192        assert_eq!(seen.concat().as_bytes(), text.as_slice());
1193    }
1194
1195    #[test]
1196    fn a_child_that_closes_its_output_can_still_be_cancelled() {
1197        // Its streams end at once, but it keeps running; cancelling must still stop it.
1198        let started = std::time::Instant::now();
1199        let outcome = Process::new("sh")
1200            .args(["-c", "exec >&- 2>&-; sleep 20"])
1201            .run(&|| started.elapsed() > std::time::Duration::from_millis(200), &mut |_| {})
1202            .expect("the shell starts");
1203        assert_eq!(outcome, ProcessOutcome::Cancelled);
1204        assert!(started.elapsed() < std::time::Duration::from_secs(10), "took {:?}", started.elapsed());
1205    }
1206
1207    #[test]
1208    fn a_flood_of_output_waits_for_the_reader_instead_of_piling_up() {
1209        let dir = std::env::temp_dir().join(format!("quvyta-process-flood-{}", std::process::id()));
1210        let _ = std::fs::remove_dir_all(&dir);
1211        std::fs::create_dir_all(&dir).expect("test directory");
1212        let marker = dir.join("done");
1213        let script = format!("yes | head -n 200000; touch '{}'", marker.display());
1214        let mut first = true;
1215        let mut finished_while_the_reader_slept = false;
1216        let mut count = 0_usize;
1217        let outcome = Process::new("sh")
1218            .args(["-c", &script])
1219            .run(&|| false, &mut |_| {
1220                count += 1;
1221                if first {
1222                    first = false;
1223                    // Far longer than writing 200000 short lines takes when nothing holds it back.
1224                    std::thread::sleep(std::time::Duration::from_millis(700));
1225                    finished_while_the_reader_slept = marker.exists();
1226                }
1227            })
1228            .expect("the shell starts");
1229        assert_eq!(outcome, ProcessOutcome::Finished { code: Some(0) });
1230        assert_eq!(count, 200_000);
1231        assert!(!finished_while_the_reader_slept, "the child wrote everything into memory while nobody read");
1232        std::fs::remove_dir_all(&dir).expect("clean");
1233    }
1234
1235    #[cfg(target_os = "linux")]
1236    #[test]
1237    fn the_child_on_a_pseudo_terminal_holds_it_only_on_its_own_streams() {
1238        // Any other descriptor of the terminal would outlive the streams in the child and in the
1239        // programs it starts, and keep the reader waiting after they are closed.
1240        let script = r#"t=$(readlink /proc/$$/fd/1); n=0; for f in /proc/$$/fd/*; do [ "$(readlink "$f")" = "$t" ] && n=$((n+1)); done; echo $n"#;
1241        let (lines, _) = run(Process::new("sh").args(["-c", script]).pty(80, 24));
1242        assert_eq!(lines, vec![Line::Out("2".to_owned())], "standard output and standard error, nothing else");
1243    }
1244
1245    #[test]
1246    fn a_line_split_across_reads_stays_one_line() {
1247        let mut lines = Lines::default();
1248        let mut seen = Vec::new();
1249        let mut emit = |line: String| seen.push(line);
1250        lines.feed(b"ilk par", &mut emit);
1251        lines.feed(b"\xc3", &mut emit);
1252        lines.feed(b"\xa7a\r\nson", &mut emit);
1253        lines.finish(&mut emit);
1254        assert_eq!(seen, vec!["ilk parça".to_owned(), "son".to_owned()]);
1255    }
1256
1257    /// Feeds `bytes` in one go and ends the stream, returning the lines and the overwritten
1258    /// frames.
1259    fn split_keeping_frames(bytes: &[u8]) -> (Vec<String>, Vec<String>) {
1260        let mut lines = Lines::default();
1261        let (mut seen, mut frames) = (Vec::new(), Vec::new());
1262        lines.feed_keeping(bytes, &mut |line| seen.push(line), Some(&mut |frame| frames.push(frame)));
1263        lines.finish_keeping(&mut |line| seen.push(line), Some(&mut |frame| frames.push(frame)));
1264        (seen, frames)
1265    }
1266
1267    #[test]
1268    fn frames_overwritten_by_a_return_are_kept_only_when_asked_for() {
1269        let (seen, frames) = split_keeping_frames(b"bir\riki\ruc\rbitti\r\n");
1270        assert_eq!(seen, vec!["bitti".to_owned()]);
1271        assert_eq!(frames, vec!["bir".to_owned(), "iki".to_owned(), "uc".to_owned()]);
1272        let mut lines = Lines::default();
1273        let mut seen = Vec::new();
1274        lines.feed(b"bir\riki\ruc\rbitti\r\n", &mut |line| seen.push(line));
1275        lines.finish(&mut |line| seen.push(line));
1276        assert_eq!(seen, vec!["bitti".to_owned()]);
1277    }
1278
1279    #[test]
1280    fn a_line_ended_by_returns_and_a_newline_is_no_frame() {
1281        let (seen, frames) = split_keeping_frames(b"hazir\r\r\nbitti\r\n\rbos\r\r\r\n");
1282        assert_eq!(seen, vec!["hazir".to_owned(), "bitti".to_owned(), "bos".to_owned()]);
1283        assert!(frames.is_empty(), "{frames:?}");
1284    }
1285
1286    #[test]
1287    fn a_stream_ending_in_a_return_delivers_its_last_frame() {
1288        let (seen, frames) = split_keeping_frames(b"once\r10%\r20%\r");
1289        assert!(seen.is_empty(), "{seen:?}");
1290        assert_eq!(frames, vec!["once".to_owned(), "10%".to_owned(), "20%".to_owned()]);
1291        // A return split from what follows it by a read still waits for that byte.
1292        let mut lines = Lines::default();
1293        let (mut seen, mut frames) = (Vec::new(), Vec::new());
1294        lines.feed_keeping(b"30%\r", &mut |line| seen.push(line), Some(&mut |frame| frames.push(frame)));
1295        assert!(frames.is_empty(), "a return before a newline is not yet known to overwrite");
1296        lines.feed_keeping(b"\n", &mut |line| seen.push(line), Some(&mut |frame| frames.push(frame)));
1297        assert_eq!((seen, frames), (vec!["30%".to_owned()], Vec::new()));
1298    }
1299
1300    #[test]
1301    fn frames_keep_their_colour_and_erase_codes() {
1302        let (seen, frames) = split_keeping_frames(b"\x1b[1mFetch\x1b[0m 1\r\x1b[K\x1b[92mDone\x1b[0m\r\n");
1303        assert_eq!(frames, vec!["\x1b[1mFetch\x1b[0m 1".to_owned()]);
1304        assert_eq!(seen, vec!["\x1b[K\x1b[92mDone\x1b[0m".to_owned()]);
1305    }
1306
1307    #[test]
1308    fn every_frame_of_a_recorded_cargo_install_is_kept() {
1309        let recorded = include_bytes!("../../tests/fixtures/cargo-install-pty.txt");
1310        let (seen, frames) = split_keeping_frames(recorded);
1311        assert_eq!(frames.len(), 159, "every overwritten frame");
1312        assert_eq!(seen.len(), 76, "the lines themselves are unchanged");
1313        let building: Vec<&String> = frames.iter().filter(|frame| frame.contains("Building")).collect();
1314        assert_eq!(building.len(), 51);
1315        assert!(building[0].contains("] 0/46: anstyle"), "{:?}", building[0]);
1316        assert!(building.iter().any(|frame| frame.contains("] 45/46: hexyl")), "{building:?}");
1317        // Without asking, the same bytes give the same lines and no frame reaches anyone.
1318        let mut lines = Lines::default();
1319        let mut plain = Vec::new();
1320        lines.feed(recorded, &mut |line| plain.push(line));
1321        lines.finish(&mut |line| plain.push(line));
1322        assert_eq!(plain, seen);
1323    }
1324
1325    /// Runs `process` to its end asking for overwritten frames, returning the lines and frames in
1326    /// the order they arrived.
1327    fn run_keeping_frames(process: Process) -> Vec<(bool, Line)> {
1328        let seen = std::cell::RefCell::new(Vec::new());
1329        process
1330            .run_with_overwritten(&|| false, &mut |line| seen.borrow_mut().push((false, line)), &mut |frame| {
1331                seen.borrow_mut().push((true, frame));
1332            })
1333            .expect("the shell starts");
1334        seen.into_inner()
1335    }
1336
1337    #[test]
1338    fn overwritten_frames_arrive_through_a_pipe_in_order_and_tagged_by_stream() {
1339        let seen = run_keeping_frames(Process::new("sh").args(["-c", r"printf 'a\rb\rc\n'; printf '1%%\r2%%\r' >&2"]));
1340        let out: Vec<_> = seen.iter().filter(|(_, line)| matches!(line, Line::Out(_))).cloned().collect();
1341        let err: Vec<_> = seen.iter().filter(|(_, line)| matches!(line, Line::Err(_))).cloned().collect();
1342        assert_eq!(
1343            out,
1344            vec![
1345                (true, Line::Out("a".to_owned())),
1346                (true, Line::Out("b".to_owned())),
1347                (false, Line::Out("c".to_owned()))
1348            ]
1349        );
1350        assert_eq!(err, vec![(true, Line::Err("1%".to_owned())), (true, Line::Err("2%".to_owned()))]);
1351    }
1352
1353    #[cfg(unix)]
1354    #[test]
1355    fn overwritten_frames_arrive_from_a_pseudo_terminal() {
1356        let seen = run_keeping_frames(Process::new("sh").args(["-c", r"printf 'a\rb\rc\nd\r\n'"]).pty(80, 24));
1357        assert_eq!(
1358            seen,
1359            vec![
1360                (true, Line::Out("a".to_owned())),
1361                (true, Line::Out("b".to_owned())),
1362                (false, Line::Out("c".to_owned())),
1363                (false, Line::Out("d".to_owned())),
1364            ]
1365        );
1366    }
1367
1368    #[test]
1369    fn collected_output_arrives_in_the_order_both_streams_were_written() {
1370        let script = "printf a; printf b >&2; printf c";
1371        let collected = collect(Process::new("sh").args(["-c", script]), Keep::bytes(64));
1372        assert_eq!(collected.text, "abc");
1373        assert_eq!(collected.outcome, ProcessOutcome::Finished { code: Some(0) });
1374        assert!(!collected.trimmed, "nothing was dropped: {}", collected.text);
1375        assert!(!collected.timed_out && !collected.cancelled);
1376        // A pseudo-terminal is the one place the two streams share a line anyway, and it comes
1377        // through the same one pipe.
1378        #[cfg(unix)]
1379        {
1380            let collected = collect(Process::new("sh").args(["-c", script]).pty(80, 24), Keep::bytes(64));
1381            assert_eq!(collected.text, "abc");
1382        }
1383    }
1384
1385    #[test]
1386    fn only_the_last_bytes_of_a_flood_of_output_are_kept() {
1387        // Eight megabytes through, sixty-four kilobytes out: the end of it, and nothing else.
1388        let script = "head -c 8388608 /dev/zero | tr '\\0' x; printf 'SON'";
1389        let collected = collect(Process::new("sh").args(["-c", script]), Keep::bytes(65536));
1390        assert_eq!(collected.text.len(), 65536);
1391        assert!(collected.text.ends_with("SON"), "the newest bytes are the ones kept");
1392        assert!(collected.text[..65533].bytes().all(|byte| byte == b'x'), "only the flood is before them");
1393        assert!(collected.trimmed, "eight megabytes did not fit in sixty-four");
1394    }
1395
1396    #[test]
1397    fn the_lines_asked_for_are_the_last_ones() {
1398        let script = "for name in bir iki uc dort bes; do printf '%s\\n' \"$name\"; done";
1399        let collected = collect(Process::new("sh").args(["-c", script]), Keep::bytes(4096).lines(3));
1400        assert_eq!(collected.text, "uc\ndort\nbes\n");
1401        assert!(collected.trimmed);
1402        // The line a program is still writing counts as one, so nothing is lost while it writes.
1403        let collected = collect(Process::new("sh").args(["-c", r"printf 'bir\niki\nuc'"]), Keep::bytes(4096).lines(3));
1404        assert_eq!(collected.text, "bir\niki\nuc");
1405        assert!(!collected.trimmed);
1406        // And bytes still decide when both are asked for.
1407        let collected = collect(Process::new("sh").args(["-c", script]), Keep::bytes(9).lines(3));
1408        assert_eq!(collected.text, "dort\nbes\n");
1409    }
1410
1411    #[test]
1412    fn a_tail_keeps_the_end_of_what_it_is_fed_and_never_half_a_character() {
1413        let mut tail = Tail::new(16, None);
1414        tail.feed(b"bir iki ");
1415        tail.feed(b"uc dort");
1416        assert_eq!(tail.text(), "bir iki uc dort");
1417        assert!(!tail.trimmed, "nothing fell out of sixteen bytes");
1418        tail.feed(b"!\n");
1419        assert_eq!(tail.text(), "ir iki uc dort!\n", "the oldest byte fell out");
1420        assert!(tail.trimmed);
1421        // One `ç` at a time, so its two bytes straddle every cut.
1422        let mut tail = Tail::new(5, None);
1423        for _ in 0..8 {
1424            tail.feed("ç".as_bytes());
1425        }
1426        assert_eq!(tail.text(), "çç", "the half a character at the front is dropped");
1427    }
1428
1429    #[test]
1430    fn a_tail_keeps_the_last_lines_it_was_given() {
1431        let mut tail = Tail::new(64, Some(2));
1432        for name in ["bir\n", "iki\n", "uc\n", "dort"] {
1433            tail.feed(name.as_bytes());
1434        }
1435        assert_eq!(tail.text(), "uc\ndort", "the line being written counts as one");
1436        assert!(tail.trimmed);
1437    }
1438
1439    #[test]
1440    fn collecting_nothing_but_the_outcome_is_allowed() {
1441        let collected = collect(Process::new("sh").args(["-c", "printf 'gorunmez'"]), Keep::bytes(0));
1442        assert_eq!(collected.text, "");
1443        assert_eq!(collected.outcome, ProcessOutcome::Finished { code: Some(0) });
1444        assert!(collected.trimmed);
1445    }
1446
1447    #[test]
1448    fn cancelling_a_collected_child_stops_it_the_way_a_limit_does() {
1449        // The child prints the program it leaves behind first, so on Linux the test can look for it
1450        // in `/proc` without running `ps`.
1451        let script = "sleep 30 & printf 'hazir\\n%s\\n' \"$!\"; sleep 30";
1452        let started = std::time::Instant::now();
1453        let process = Process::new("sh").args(["-c", script]).no_stdin();
1454        let collected = process
1455            .collect(Keep::bytes(1024).limit(Duration::from_secs(30)), &|| started.elapsed() > Duration::from_secs(1))
1456            .expect("the shell starts");
1457        assert!(collected.cancelled && !collected.timed_out, "{collected:?}");
1458        assert_eq!(collected.outcome, ProcessOutcome::Cancelled);
1459        assert!(started.elapsed() < Duration::from_secs(15), "took {:?}", started.elapsed());
1460        let mut lines = collected.text.lines();
1461        assert_eq!(lines.next(), Some("hazir"), "what was written before the cancel is still there: {collected:?}");
1462        // The child is in a process group of its own, so the cancel reached the `sleep` too.
1463        #[cfg(target_os = "linux")]
1464        {
1465            let pid = lines.next().unwrap_or_default().to_owned();
1466            assert!(pid.bytes().all(|byte| byte.is_ascii_digit()), "a process id: {pid:?}");
1467            let started = std::time::Instant::now();
1468            while !ended(&pid) {
1469                assert!(started.elapsed() < Duration::from_secs(20), "the `sleep` is still running: {pid}");
1470                std::thread::sleep(Duration::from_millis(20));
1471            }
1472        }
1473    }
1474
1475    #[cfg(target_os = "linux")]
1476    #[test]
1477    fn the_limit_ends_the_child_and_the_programs_it_started() {
1478        // The child prints the program it leaves behind first, so the test can look for it in
1479        // `/proc` without running `ps`.
1480        let script = "sleep 30 & printf '%s\\n' \"$!\"; sleep 30";
1481        let started = std::time::Instant::now();
1482        let process = Process::new("sh").args(["-c", script]).no_stdin();
1483        let collected =
1484            process.collect(Keep::bytes(1024).limit(Duration::from_secs(1)), &|| false).expect("the shell starts");
1485        assert!(collected.timed_out && !collected.cancelled, "{collected:?}");
1486        assert_eq!(collected.outcome, ProcessOutcome::Finished { code: None }, "a signal ended it");
1487        assert!(started.elapsed() < Duration::from_secs(15), "took {:?}", started.elapsed());
1488        // The child is in a process group of its own, so the limit reached the `sleep` too.
1489        let pid = collected.text.trim().to_owned();
1490        assert!(pid.bytes().all(|byte| byte.is_ascii_digit()), "a process id: {pid:?}");
1491        let started = std::time::Instant::now();
1492        while !ended(&pid) {
1493            assert!(started.elapsed() < Duration::from_secs(20), "the `sleep` is still running: {pid}");
1494            std::thread::sleep(Duration::from_millis(20));
1495        }
1496    }
1497}