Skip to main content

rightkit_process/
stderr_tail.rs

1//! Bounded capture of a child's stderr so a crash is diagnosable without a
2//! debugger (ScrapeRight keeps a 16 KiB tail; HeardRight redacts each line
3//! before it reaches telemetry). Draining stderr is also what keeps a chatty
4//! child from blocking on a full pipe.
5
6use std::collections::VecDeque;
7use std::io::{BufRead, BufReader, Read};
8use std::sync::{Arc, Mutex};
9use std::thread::{self, JoinHandle};
10
11/// A background drain of one stderr stream. Keeps the most recent lines whose
12/// total size is at most `max_bytes`; older lines fall off the front.
13pub struct StderrTail {
14    lines: Arc<Mutex<(VecDeque<String>, usize)>>,
15    handle: Option<JoinHandle<()>>,
16}
17
18impl StderrTail {
19    /// Drain `stderr` on a thread. `map` sees every line first and returns the
20    /// text to keep (redaction), or `None` to drop it. `on_line` receives the
21    /// kept text as it arrives (forward to logs/telemetry).
22    pub fn capture<R: Read + Send + 'static>(
23        stderr: R,
24        max_bytes: usize,
25        map: impl Fn(&str) -> Option<String> + Send + 'static,
26        on_line: impl Fn(&str) + Send + 'static,
27    ) -> Self {
28        let lines: Arc<Mutex<(VecDeque<String>, usize)>> = Arc::default();
29        let shared = Arc::clone(&lines);
30        let handle = thread::Builder::new()
31            .name("rightkit-stderr-tail".into())
32            .spawn(move || {
33                let mut reader = BufReader::new(stderr);
34                let mut raw = Vec::new();
35                loop {
36                    raw.clear();
37                    // A pathological unterminated line must not grow without bound.
38                    match reader
39                        .by_ref()
40                        .take(max_bytes.max(1) as u64 * 4)
41                        .read_until(b'\n', &mut raw)
42                    {
43                        Ok(0) | Err(_) => return,
44                        Ok(_) => {}
45                    }
46                    let text = String::from_utf8_lossy(&raw);
47                    let Some(kept) = map(text.trim_end_matches(['\r', '\n'])) else {
48                        continue;
49                    };
50                    on_line(&kept);
51                    let mut g = shared.lock().unwrap_or_else(|e| e.into_inner());
52                    g.1 += kept.len() + 1;
53                    g.0.push_back(kept);
54                    while g.1 > max_bytes && g.0.len() > 1 {
55                        if let Some(old) = g.0.pop_front() {
56                            g.1 -= old.len() + 1;
57                        }
58                    }
59                }
60            });
61        Self {
62            lines,
63            handle: handle.ok(),
64        }
65    }
66
67    /// Capture with no redaction and no forwarding.
68    pub fn plain<R: Read + Send + 'static>(stderr: R, max_bytes: usize) -> Self {
69        Self::capture(stderr, max_bytes, |l| Some(l.to_string()), |_| {})
70    }
71
72    /// The retained lines, oldest first, joined by newlines.
73    pub fn snapshot(&self) -> String {
74        let g = self.lines.lock().unwrap_or_else(|e| e.into_inner());
75        g.0.iter()
76            .map(String::as_str)
77            .collect::<Vec<_>>()
78            .join("\n")
79    }
80
81    /// Like [`Self::finish`] but gives up after `timeout` (a grandchild that
82    /// inherited stderr can keep the stream open after the child is gone) and
83    /// returns what was captured so far.
84    pub fn finish_within(&mut self, timeout: std::time::Duration) -> String {
85        let deadline = std::time::Instant::now() + timeout;
86        while self.handle.as_ref().is_some_and(|h| !h.is_finished())
87            && std::time::Instant::now() < deadline
88        {
89            thread::sleep(std::time::Duration::from_millis(5));
90        }
91        if self.handle.as_ref().is_some_and(|h| h.is_finished()) {
92            if let Some(h) = self.handle.take() {
93                let _ = h.join();
94            }
95        }
96        self.snapshot()
97    }
98
99    /// Wait for the stream to end (the child exited) so the tail is complete.
100    pub fn finish(&mut self) -> String {
101        if let Some(h) = self.handle.take() {
102            let _ = h.join();
103        }
104        self.snapshot()
105    }
106}