Skip to main content

subc_daemon/
stderr_tail.rs

1//! Bounded per-module stderr capture.
2//!
3//! # Why this exists
4//!
5//! `last_exit_code` survives a respawn because the supervisor holds it in memory.
6//! Stderr had no such path: it went to the daemon's inherited fd, and from there
7//! to whatever rotates or evicts it. On the box this was written for, that window
8//! was about three hours; on another it was bounded by a log file reaching 908 MB
9//! with one module accounting for 98% of it. Two hosts, two mechanisms, the same
10//! outcome -- the text explaining a crash is gone by the time anyone asks.
11//!
12//! So this keeps the last few lines where `last_exit` already lives: in supervisor
13//! memory, immune to whatever happens to the log.
14//!
15//! # What it is not
16//!
17//! Not a log. The ring is deliberately small and lossy, and callers are expected
18//! to know they are reading a tail rather than a history. The daemon log keeps
19//! doing its job; this exists because that job has a time limit.
20
21use std::collections::VecDeque;
22use std::io::{self, Write};
23use std::path::{Path, PathBuf};
24use std::sync::{
25    atomic::{AtomicBool, Ordering},
26    Arc, Mutex,
27};
28
29use tokio::io::AsyncReadExt;
30
31/// Longest single line admitted to the ring before truncation.
32///
33/// A module emitting a 40 MB backtrace on one line satisfies any line-count cap
34/// while evicting everything else -- the pathological emitter wins twice, once by
35/// filling the ring and once by being unreadable itself. Truncating on the way in
36/// costs that emitter one line instead of the whole tail.
37pub const DEFAULT_MAX_LINE_BYTES: usize = 2048;
38
39/// Lines retained per module.
40pub const DEFAULT_MAX_LINES: usize = 200;
41
42/// Total bytes retained per module, across all lines.
43///
44/// Both this and [`DEFAULT_MAX_LINES`] apply; whichever binds first wins. A line
45/// cap alone is satisfied by 200 lines of 2 KB, which is not a budget worth
46/// holding for fifteen modules.
47pub const DEFAULT_MAX_BYTES: usize = 64 * 1024;
48
49/// Whether a module's stderr is being captured, and if not, why not.
50///
51/// Typed rather than nullable so `NotCaptured` has to be handled rather than
52/// defaulted past. An empty tail and an uncaptured one send an operator in
53/// opposite directions -- one says the module printed nothing before dying, the
54/// other says nobody was listening -- and rendering them alike is the defect this
55/// module exists to fix, reproduced one layer up.
56#[derive(Debug, Clone, PartialEq, Eq)]
57pub enum CaptureState {
58    /// A reader is attached, or was attached and reached clean EOF.
59    Captured,
60    /// Retained entries are valid, but capture ended before clean EOF.
61    Incomplete { reason: String },
62    /// No reader was attached. The tail says nothing about what the module wrote.
63    NotCaptured { reason: String },
64}
65
66/// One retained entry.
67///
68/// Boundaries are in-band rather than a separate field because their position
69/// relative to the lines is the whole point: "these three lines came from the
70/// process that died, those came from its replacement" is unanswerable from a
71/// count.
72#[derive(Debug, Clone, PartialEq, Eq)]
73pub enum TailEntry {
74    Line {
75        text: String,
76        /// This line was cut at the per-line cap.
77        truncated: bool,
78    },
79    /// The supervisor spawned a new process for this module. Lines after this
80    /// entry come from the new one.
81    ProcessStart,
82}
83
84impl TailEntry {
85    fn cost(&self) -> usize {
86        match self {
87            Self::Line { text, .. } => text.len(),
88            Self::ProcessStart => 0,
89        }
90    }
91}
92
93#[derive(Debug, Clone, Copy, PartialEq, Eq)]
94pub struct StderrTailConfig {
95    max_lines: usize,
96    max_bytes: usize,
97    max_line_bytes: usize,
98}
99
100impl StderrTailConfig {
101    /// Keeps every retained line within the ring's total byte budget.
102    /// Clamp rather than reject so diagnostics degrade without blocking supervisor startup.
103    pub const fn new(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> Self {
104        Self {
105            max_lines,
106            max_bytes,
107            max_line_bytes: if max_line_bytes > max_bytes {
108                max_bytes
109            } else {
110                max_line_bytes
111            },
112        }
113    }
114}
115
116impl Default for StderrTailConfig {
117    fn default() -> Self {
118        Self::new(DEFAULT_MAX_LINES, DEFAULT_MAX_BYTES, DEFAULT_MAX_LINE_BYTES)
119    }
120}
121
122/// A module's retained stderr, oldest first.
123#[derive(Debug, Clone, PartialEq, Eq)]
124pub struct StderrTailSnapshot {
125    pub capture: CaptureState,
126    pub entries: Vec<TailEntry>,
127    /// Lines evicted since the module was first supervised.
128    ///
129    /// Non-zero means the tail starts mid-stream. That is the ring working as
130    /// intended, but a reader diagnosing a crash needs to know the first retained
131    /// line is not the first line the module wrote -- otherwise an absent cause
132    /// reads as a module that never explained itself.
133    pub dropped_lines: u64,
134}
135
136impl StderrTailSnapshot {
137    /// The uncaptured case, for a module whose stderr was never piped.
138    pub fn not_captured(reason: impl Into<String>) -> Self {
139        Self {
140            capture: CaptureState::NotCaptured {
141                reason: reason.into(),
142            },
143            entries: Vec::new(),
144            dropped_lines: 0,
145        }
146    }
147}
148
149/// Bounded ring of a single module's stderr lines.
150///
151/// Survives respawn deliberately. The stderr explaining an exit is written
152/// *before* that exit, so clearing on restart would discard the lines exactly
153/// when they become the thing being asked for. [`TailEntry::ProcessStart`] keeps
154/// the generations distinguishable instead.
155#[derive(Debug)]
156pub struct StderrRing {
157    config: StderrTailConfig,
158    entries: VecDeque<TailEntry>,
159    // Count of `TailEntry::Line` entries, kept running because eviction checks
160    // it on every push and recounting would walk the whole ring under the
161    // mutex each time.
162    lines: usize,
163    bytes: usize,
164    dropped_lines: u64,
165    capture: CaptureState,
166}
167
168impl StderrRing {
169    pub fn new(config: StderrTailConfig) -> Self {
170        Self {
171            config,
172            entries: VecDeque::new(),
173            lines: 0,
174            bytes: 0,
175            dropped_lines: 0,
176            // Until a reader attaches, the honest answer is that nothing is
177            // listening -- not that the module has been quiet.
178            capture: CaptureState::NotCaptured {
179                reason: "stderr reader has not started".to_string(),
180            },
181        }
182    }
183
184    pub fn mark_captured(&mut self) {
185        if matches!(self.capture, CaptureState::NotCaptured { .. }) {
186            self.capture = CaptureState::Captured;
187        }
188    }
189
190    pub fn mark_incomplete(&mut self, reason: impl Into<String>) {
191        self.capture = CaptureState::Incomplete {
192            reason: reason.into(),
193        };
194    }
195
196    pub fn mark_not_captured(&mut self, reason: impl Into<String>) {
197        self.capture = CaptureState::NotCaptured {
198            reason: reason.into(),
199        };
200    }
201
202    /// Record that a new process was spawned for this module.
203    ///
204    /// A boundary separates output on either side of it, so one with nothing
205    /// before it separates nothing: on the FIRST spawn it would make a module
206    /// that printed nothing render as a marker rather than as empty, and the
207    /// caller then has to decide whether a one-marker tail counts as silence.
208    /// Recording it only once there is something to divide keeps "captured and
209    /// empty" literally empty.
210    pub fn push_process_start(&mut self) {
211        if self.entries.is_empty() && self.dropped_lines == 0 {
212            return;
213        }
214        if matches!(self.entries.back(), Some(TailEntry::ProcessStart)) {
215            return;
216        }
217        self.push_entry(TailEntry::ProcessStart);
218    }
219
220    /// Admit one complete line, truncating it if it exceeds the per-line cap.
221    ///
222    /// `line` must not contain a trailing newline; the reader strips it so the
223    /// stored text and the byte accounting agree.
224    pub fn push_line(&mut self, line: &str) {
225        let (text, truncated) = truncate_line(line, self.config.max_line_bytes);
226        self.push_entry(TailEntry::Line { text, truncated });
227    }
228
229    fn push_entry(&mut self, entry: TailEntry) {
230        self.bytes += entry.cost();
231        if matches!(entry, TailEntry::Line { .. }) {
232            self.lines += 1;
233        }
234        self.entries.push_back(entry);
235        self.evict_to_fit();
236    }
237
238    fn evict_to_fit(&mut self) {
239        while self.lines > self.config.max_lines
240            || (self.bytes > self.config.max_bytes && self.entries.len() > 1)
241        {
242            let Some(evicted) = self.entries.pop_front() else {
243                break;
244            };
245            self.bytes -= evicted.cost();
246            if matches!(evicted, TailEntry::Line { .. }) {
247                self.lines -= 1;
248                self.dropped_lines += 1;
249            }
250        }
251    }
252
253    /// The most recent entries, oldest first, bounded by the caller's limits.
254    ///
255    /// `max_lines`/`max_bytes` narrow the ring's own caps; they cannot widen them.
256    pub fn snapshot(
257        &self,
258        max_lines: Option<usize>,
259        max_bytes: Option<usize>,
260    ) -> StderrTailSnapshot {
261        let line_limit = max_lines.unwrap_or(self.config.max_lines);
262        let byte_limit = max_bytes.unwrap_or(self.config.max_bytes);
263
264        let mut taken: Vec<TailEntry> = Vec::new();
265        let mut bytes = 0usize;
266        let mut lines = 0usize;
267        // Walk backwards: a tail is anchored at the newest end, so a caller
268        // asking for 20 lines wants the last 20, not the first 20.
269        for entry in self.entries.iter().rev() {
270            match entry {
271                TailEntry::Line { .. } => {
272                    if lines >= line_limit {
273                        break;
274                    }
275                    let cost = entry.cost();
276                    if lines > 0 && bytes + cost > byte_limit {
277                        break;
278                    }
279                    bytes += cost;
280                    lines += 1;
281                    taken.push(entry.clone());
282                }
283                TailEntry::ProcessStart if lines > 0 => taken.push(entry.clone()),
284                TailEntry::ProcessStart => {}
285            }
286        }
287        taken.reverse();
288
289        let withheld = self
290            .entries
291            .iter()
292            .filter(|entry| matches!(entry, TailEntry::Line { .. }))
293            .count()
294            .saturating_sub(
295                taken
296                    .iter()
297                    .filter(|entry| matches!(entry, TailEntry::Line { .. }))
298                    .count(),
299            );
300
301        StderrTailSnapshot {
302            capture: self.capture.clone(),
303            entries: taken,
304            // Lines the ring evicted plus lines this request's own limits held
305            // back. Both mean the same thing to the reader -- the text above is
306            // not the beginning -- and separating them would invite treating a
307            // narrow request as evidence of a quiet module.
308            dropped_lines: self.dropped_lines + withheld as u64,
309        }
310    }
311}
312
313/// Reassembly buffer ceiling for a line with no newline in sight.
314///
315/// The ring truncates what it stores, but the READER has to hold the bytes until
316/// it finds a delimiter. A module emitting a gigabyte with no newline would grow
317/// this buffer without bound and take the daemon down with it -- a module fault
318/// escalating into a fleet fault, which is exactly what supervision exists to
319/// prevent. At this ceiling the pending bytes are flushed as a line and
320/// reassembly restarts.
321const MAX_PENDING_LINE_BYTES: usize = 1024 * 1024;
322
323/// Read a child's stderr to EOF, retaining the bounded crash tail and forwarding
324/// every complete line to the selected capture sink.
325pub async fn pump_stderr<R>(source: R, ring: Arc<Mutex<StderrRing>>)
326where
327    R: AsyncReadExt + Unpin,
328{
329    pump_stderr_into(source, ring, &mut StderrSink).await
330}
331
332/// Shared destination for a child's stdout and stderr pumps.
333///
334/// The file keeps the historical `.stderr.log` name even though it carries both
335/// streams; the stable name is part of the operator contract. Both pumps share
336/// one mutex, and cortexkit-log writes each framed line in one call, so partial
337/// lines from the two pipes cannot interleave.
338#[derive(Clone)]
339pub(crate) enum ChildOutputSink {
340    File {
341        sink: Arc<Mutex<cortexkit_log::LineSink>>,
342        path: Arc<PathBuf>,
343        failure_reported: Arc<AtomicBool>,
344    },
345    Stderr,
346}
347
348impl ChildOutputSink {
349    pub(crate) fn open(path: &Path, retention: cortexkit_log::Retention) -> io::Result<Self> {
350        Ok(Self::File {
351            sink: Arc::new(Mutex::new(cortexkit_log::LineSink::open(path, retention)?)),
352            path: Arc::new(path.to_path_buf()),
353            failure_reported: Arc::new(AtomicBool::new(false)),
354        })
355    }
356}
357
358/// Where forwarded complete lines go. It exists so tests can observe framing
359/// and so production can serialize the two child pipes through one file sink.
360pub trait OutputSink {
361    fn write_line(&mut self, line: &[u8]);
362}
363
364struct StderrSink;
365
366impl OutputSink for StderrSink {
367    fn write_line(&mut self, line: &[u8]) {
368        let stderr = std::io::stderr();
369        let mut handle = stderr.lock();
370        let _ = handle.write_all(line);
371    }
372}
373
374impl OutputSink for ChildOutputSink {
375    fn write_line(&mut self, line: &[u8]) {
376        match self {
377            Self::File {
378                sink,
379                path,
380                failure_reported,
381            } => {
382                let result = sink
383                    .lock()
384                    .unwrap_or_else(|poisoned| poisoned.into_inner())
385                    .write_line(line);
386                if let Err(error) = result {
387                    if !failure_reported.swap(true, Ordering::Relaxed) {
388                        tracing::warn!(
389                            path = %path.display(),
390                            error = %error,
391                            "child output capture write failed; later failures are suppressed"
392                        );
393                    }
394                }
395            }
396            Self::Stderr => StderrSink.write_line(line),
397        }
398    }
399}
400
401pub(crate) async fn pump_stderr_to<R>(
402    source: R,
403    ring: Arc<Mutex<StderrRing>>,
404    mut sink: ChildOutputSink,
405) where
406    R: AsyncReadExt + Unpin,
407{
408    pump_stderr_into(source, ring, &mut sink).await;
409}
410
411pub(crate) async fn pump_stdout_to<R>(source: R, mut sink: ChildOutputSink)
412where
413    R: AsyncReadExt + Unpin,
414{
415    pump_lines_into(source, None, &mut sink, "stdout").await;
416}
417
418async fn pump_stderr_into<R, S>(source: R, ring: Arc<Mutex<StderrRing>>, sink: &mut S)
419where
420    R: AsyncReadExt + Unpin,
421    S: OutputSink,
422{
423    pump_lines_into(source, Some(&ring), sink, "stderr").await;
424}
425
426async fn pump_lines_into<R, S>(
427    mut source: R,
428    ring: Option<&Arc<Mutex<StderrRing>>>,
429    sink: &mut S,
430    stream_name: &str,
431) where
432    R: AsyncReadExt + Unpin,
433    S: OutputSink,
434{
435    if let Some(ring) = ring {
436        lock_ring(ring).mark_captured();
437    }
438
439    let mut pending: Vec<u8> = Vec::new();
440    // Bytes before `scanned_upto` are already known to hold no newline; searching
441    // them again would rescan the whole buffer on every chunk -- for a line with
442    // no newline that is about 64 MiB examined per MiB of module output.
443    let mut scanned_upto = 0usize;
444    // Bytes before `cursor` were emitted as complete lines. They are removed in
445    // one compaction per chunk rather than shifting the buffer once per line.
446    let mut cursor = 0usize;
447    let mut chunk = [0u8; 8192];
448    loop {
449        let read = match source.read(&mut chunk).await {
450            Ok(0) => break,
451            Ok(n) => n,
452            Err(error) => {
453                if let Some(ring) = ring {
454                    lock_ring(ring).mark_incomplete(format!("{stream_name} read failed: {error}"));
455                } else {
456                    tracing::warn!(stream = stream_name, error = %error, "child output capture read failed");
457                }
458                return;
459            }
460        };
461        pending.extend_from_slice(&chunk[..read]);
462
463        while let Some(relative) = find_newline(&pending[scanned_upto..]) {
464            let newline = scanned_upto + relative;
465            emit_line(ring, sink, &pending[cursor..newline], true);
466            cursor = newline + 1;
467            scanned_upto = cursor;
468        }
469        scanned_upto = pending.len();
470
471        if cursor > 0 {
472            pending.drain(..cursor);
473            scanned_upto -= cursor;
474            cursor = 0;
475        }
476
477        if pending.len() >= MAX_PENDING_LINE_BYTES {
478            let line = std::mem::take(&mut pending);
479            emit_line(ring, sink, &line, false);
480            scanned_upto = 0;
481        }
482    }
483
484    if !pending.is_empty() {
485        emit_line(ring, sink, &pending, false);
486    }
487}
488
489// Bytes examined by newline searches, summed across a pump. Tests use this to
490// assert the reader does not rescan bytes it already knows contain no newline.
491// Per thread, not process-wide: other tests pump concurrently on their own
492// threads, and a shared counter picked up their searches too, failing the
493// bound by a few bytes under a parallel run. `#[tokio::test]` runs its body
494// and every future it awaits on one thread, so a pump's searches land on the
495// test's own counter.
496#[cfg(test)]
497thread_local! {
498    static SCANNED_BYTES: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
499}
500
501#[cfg(test)]
502fn take_scanned_bytes() -> usize {
503    SCANNED_BYTES.with(|scanned| scanned.replace(0))
504}
505
506/// Locate the next newline in `haystack`, counting the bytes examined so a
507/// test can observe how much of the pending buffer each search walks.
508fn find_newline(haystack: &[u8]) -> Option<usize> {
509    let found = memchr::memchr(b'\n', haystack);
510    #[cfg(test)]
511    SCANNED_BYTES.with(|scanned| {
512        scanned.set(scanned.get() + found.map(|index| index + 1).unwrap_or(haystack.len()));
513    });
514    found
515}
516
517fn emit_line<S: OutputSink>(
518    ring: Option<&Arc<Mutex<StderrRing>>>,
519    sink: &mut S,
520    raw: &[u8],
521    terminated: bool,
522) {
523    if let Some(ring) = ring {
524        lock_ring(ring).push_line(&String::from_utf8_lossy(raw));
525    }
526
527    // Framed and written in ONE call. Two writes would let the other pipe land
528    // between the body and newline.
529    if terminated {
530        let mut framed = Vec::with_capacity(raw.len() + 1);
531        framed.extend_from_slice(raw);
532        framed.push(b'\n');
533        sink.write_line(&framed);
534    } else {
535        sink.write_line(raw);
536    }
537}
538
539fn lock_ring(ring: &Arc<Mutex<StderrRing>>) -> std::sync::MutexGuard<'_, StderrRing> {
540    ring.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
541}
542
543/// Cut `line` to at most `max_bytes`, reporting whether it was shortened.
544///
545/// Cuts on a char boundary: slicing a multi-byte sequence would produce invalid
546/// UTF-8, and a panic while capturing a crash message is the worst possible time
547/// to discover that.
548fn truncate_line(line: &str, max_bytes: usize) -> (String, bool) {
549    if line.len() <= max_bytes {
550        return (line.to_string(), false);
551    }
552    let mut end = max_bytes;
553    while end > 0 && !line.is_char_boundary(end) {
554        end -= 1;
555    }
556    (line[..end].to_string(), true)
557}
558
559#[cfg(test)]
560mod tests {
561    use std::{
562        io,
563        pin::Pin,
564        task::{Context, Poll},
565    };
566
567    use super::*;
568    use tokio::io::{AsyncRead, ReadBuf};
569
570    fn ring(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> StderrRing {
571        StderrRing::new(StderrTailConfig::new(max_lines, max_bytes, max_line_bytes))
572    }
573
574    fn lines(snapshot: &StderrTailSnapshot) -> Vec<String> {
575        snapshot
576            .entries
577            .iter()
578            .filter_map(|entry| match entry {
579                TailEntry::Line { text, .. } => Some(text.clone()),
580                TailEntry::ProcessStart => None,
581            })
582            .collect()
583    }
584
585    #[test]
586    fn a_fresh_ring_reports_not_captured_rather_than_empty() {
587        // The distinction this whole module exists for: "nobody was listening"
588        // must not render as "the module said nothing".
589        let ring = ring(10, 1024, 128);
590        let snapshot = ring.snapshot(None, None);
591        assert!(matches!(snapshot.capture, CaptureState::NotCaptured { .. }));
592        assert!(snapshot.entries.is_empty());
593    }
594
595    #[test]
596    fn a_captured_module_that_printed_nothing_is_distinguishable_from_an_uncaptured_one() {
597        let mut captured = ring(10, 1024, 128);
598        captured.mark_captured();
599        let uncaptured = ring(10, 1024, 128);
600
601        let captured = captured.snapshot(None, None);
602        let uncaptured = uncaptured.snapshot(None, None);
603
604        // Both are empty. Only the capture state separates them, which is the
605        // point -- an assertion on emptiness alone would pass either way.
606        assert!(captured.entries.is_empty());
607        assert!(uncaptured.entries.is_empty());
608        assert_eq!(captured.capture, CaptureState::Captured);
609        assert!(matches!(
610            uncaptured.capture,
611            CaptureState::NotCaptured { .. }
612        ));
613    }
614
615    #[test]
616    fn the_line_cap_evicts_oldest_first_and_counts_what_it_dropped() {
617        let mut ring = ring(3, 10_000, 128);
618        ring.mark_captured();
619        for i in 0..6 {
620            ring.push_line(&format!("line{i}"));
621        }
622        let snapshot = ring.snapshot(None, None);
623        assert_eq!(lines(&snapshot), vec!["line3", "line4", "line5"]);
624        // Without this the tail silently becomes "the last lines that happened
625        // to survive" and reads as complete.
626        assert_eq!(snapshot.dropped_lines, 3);
627    }
628
629    #[test]
630    fn the_byte_cap_binds_before_the_line_cap_when_lines_are_large() {
631        // 100 lines allowed, but only ~30 bytes of them.
632        let mut ring = ring(100, 30, 128);
633        ring.mark_captured();
634        for i in 0..10 {
635            ring.push_line(&format!("{i}--------")); // 9 bytes each
636        }
637        let snapshot = ring.snapshot(None, None);
638        assert!(
639            snapshot.entries.len() < 10,
640            "byte cap did not bind: {} entries retained",
641            snapshot.entries.len()
642        );
643        let retained: usize = lines(&snapshot).iter().map(String::len).sum();
644        assert!(
645            retained <= 30,
646            "retained {retained} bytes over a 30 byte cap"
647        );
648        assert!(snapshot.dropped_lines > 0);
649    }
650
651    #[test]
652    fn one_enormous_line_is_truncated_rather_than_evicting_the_tail() {
653        // The pathological-emitter case: without per-line truncation this single
654        // line would evict every other line AND be unreadable itself.
655        let mut ring = ring(10, 10_000, 64);
656        ring.mark_captured();
657        ring.push_line("context line that must survive");
658        ring.push_line(&"x".repeat(40_000));
659
660        let snapshot = ring.snapshot(None, None);
661        let kept = &snapshot.entries;
662        assert!(matches!(
663            &kept[0],
664            TailEntry::Line { text, truncated: false }
665                if text == "context line that must survive"
666        ));
667        let TailEntry::Line { text, truncated } = &kept[1] else {
668            panic!("expected a truncated line");
669        };
670        assert_eq!(text, &"x".repeat(64));
671        assert!(*truncated);
672    }
673
674    #[test]
675    fn truncation_is_visible_so_a_cut_line_is_not_mistaken_for_a_short_one() {
676        let mut ring = ring(10, 10_000, 16);
677        ring.mark_captured();
678        ring.push_line("0123456789abcdefghij");
679        ring.push_line("short");
680
681        let snapshot = ring.snapshot(None, None);
682        let TailEntry::Line { truncated, .. } = &snapshot.entries[0] else {
683            panic!("expected a line");
684        };
685        assert!(truncated);
686        let TailEntry::Line { truncated, .. } = &snapshot.entries[1] else {
687            panic!("expected a line");
688        };
689        assert!(!truncated, "a short line must not be reported as truncated");
690    }
691
692    #[test]
693    fn truncation_cuts_on_a_char_boundary_rather_than_splitting_utf8() {
694        // A panic message with non-ASCII in it is not exotic, and slicing mid
695        // sequence would panic while capturing a crash.
696        let mut ring = ring(10, 10_000, 5);
697        ring.mark_captured();
698        ring.push_line("aa€€€€");
699        let snapshot = ring.snapshot(None, None);
700        let TailEntry::Line { text, truncated } = &snapshot.entries[0] else {
701            panic!("expected a line");
702        };
703        assert!(truncated);
704        assert!(text.starts_with("aa"));
705    }
706
707    #[test]
708    fn a_restart_boundary_keeps_generations_distinguishable() {
709        let mut ring = ring(10, 10_000, 128);
710        ring.mark_captured();
711        ring.push_line("before the crash");
712        ring.push_process_start();
713        ring.push_line("after the respawn");
714
715        let snapshot = ring.snapshot(None, None);
716        assert_eq!(
717            snapshot.entries,
718            vec![
719                TailEntry::Line {
720                    text: "before the crash".to_string(),
721                    truncated: false
722                },
723                TailEntry::ProcessStart,
724                TailEntry::Line {
725                    text: "after the respawn".to_string(),
726                    truncated: false
727                },
728            ]
729        );
730    }
731
732    #[test]
733    fn the_ring_survives_respawn_because_the_cause_is_written_before_the_exit() {
734        // Clearing on restart would discard the lines at the exact moment they
735        // become the thing being asked for.
736        let mut ring = ring(10, 10_000, 128);
737        ring.mark_captured();
738        ring.push_line("Error: storage section missing");
739        ring.push_process_start();
740
741        let snapshot = ring.snapshot(None, None);
742        assert!(lines(&snapshot).contains(&"Error: storage section missing".to_string()));
743    }
744
745    #[test]
746    fn a_caller_limit_returns_the_newest_lines_not_the_oldest() {
747        let mut ring = ring(100, 100_000, 128);
748        ring.mark_captured();
749        for i in 0..10 {
750            ring.push_line(&format!("line{i}"));
751        }
752        let snapshot = ring.snapshot(Some(3), None);
753        assert_eq!(lines(&snapshot), vec!["line7", "line8", "line9"]);
754    }
755
756    #[test]
757    fn a_caller_line_limit_keeps_the_boundary_before_the_selected_line() {
758        let mut ring = ring(100, 100_000, 128);
759        ring.mark_captured();
760        ring.push_line("before restart");
761        ring.push_process_start();
762        ring.push_line("after restart");
763
764        let snapshot = ring.snapshot(Some(1), None);
765        assert_eq!(
766            snapshot.entries,
767            vec![
768                TailEntry::ProcessStart,
769                TailEntry::Line {
770                    text: "after restart".to_string(),
771                    truncated: false,
772                },
773            ]
774        );
775    }
776
777    #[test]
778    fn a_caller_line_limit_omits_a_trailing_boundary_after_the_selected_line() {
779        let mut ring = ring(100, 100_000, 128);
780        ring.mark_captured();
781        ring.push_line("before restart");
782        ring.push_process_start();
783
784        let snapshot = ring.snapshot(Some(1), None);
785        assert_eq!(
786            snapshot.entries,
787            vec![TailEntry::Line {
788                text: "before restart".to_string(),
789                truncated: false,
790            }]
791        );
792    }
793
794    #[test]
795    fn a_caller_limit_reports_what_it_withheld_rather_than_looking_complete() {
796        let mut ring = ring(100, 100_000, 128);
797        ring.mark_captured();
798        for i in 0..10 {
799            ring.push_line(&format!("line{i}"));
800        }
801        // Nothing was evicted; the narrowing is the caller's own. It still has to
802        // be reported, or a 3-line request reads as a module that wrote 3 lines.
803        assert_eq!(ring.snapshot(Some(3), None).dropped_lines, 7);
804        assert_eq!(ring.snapshot(None, None).dropped_lines, 0);
805    }
806
807    #[test]
808    fn a_caller_limit_cannot_widen_the_rings_own_caps() {
809        let mut ring = ring(2, 10_000, 128);
810        ring.mark_captured();
811        for i in 0..5 {
812            ring.push_line(&format!("line{i}"));
813        }
814        let snapshot = ring.snapshot(Some(1000), Some(1_000_000));
815        assert_eq!(lines(&snapshot).len(), 2);
816    }
817
818    fn shared(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> Arc<Mutex<StderrRing>> {
819        Arc::new(Mutex::new(ring(max_lines, max_bytes, max_line_bytes)))
820    }
821
822    /// Records each forwarded write separately, so a test can tell one write of
823    /// `b"abc\n"` from two writes of `b"abc"` and `b"\n"`.
824    #[derive(Default)]
825    struct RecordingSink {
826        writes: Vec<Vec<u8>>,
827    }
828
829    impl OutputSink for RecordingSink {
830        fn write_line(&mut self, line: &[u8]) {
831            self.writes.push(line.to_vec());
832        }
833    }
834
835    /// Yields predetermined chunks, one per read, so a test controls exactly
836    /// where the byte stream is split.
837    struct ChunkedReader {
838        chunks: VecDeque<Vec<u8>>,
839    }
840
841    impl AsyncRead for ChunkedReader {
842        fn poll_read(
843            mut self: Pin<&mut Self>,
844            _cx: &mut Context<'_>,
845            buf: &mut ReadBuf<'_>,
846        ) -> Poll<io::Result<()>> {
847            match self.chunks.pop_front() {
848                None => Poll::Ready(Ok(())),
849                Some(chunk) => {
850                    buf.put_slice(&chunk);
851                    Poll::Ready(Ok(()))
852                }
853            }
854        }
855    }
856
857    struct FailingReader {
858        bytes: Vec<u8>,
859        emitted: bool,
860    }
861
862    impl AsyncRead for FailingReader {
863        fn poll_read(
864            mut self: Pin<&mut Self>,
865            _cx: &mut Context<'_>,
866            buf: &mut ReadBuf<'_>,
867        ) -> Poll<io::Result<()>> {
868            if self.emitted {
869                return Poll::Ready(Err(io::Error::other("reader failed")));
870            }
871            self.emitted = true;
872            buf.put_slice(&self.bytes);
873            Poll::Ready(Ok(()))
874        }
875    }
876
877    #[tokio::test]
878    async fn the_pump_splits_on_newlines_and_keeps_a_trailing_fragment() {
879        let ring = shared(10, 10_000, 128);
880        // No trailing newline on the last line: a crashing process routinely dies
881        // mid-line, and that fragment is often the message worth reading.
882        let source = std::io::Cursor::new(b"one\ntwo\nthree".to_vec());
883        let mut sink = RecordingSink::default();
884        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
885
886        let snapshot = lock_ring(&ring).snapshot(None, None);
887        assert_eq!(lines(&snapshot), vec!["one", "two", "three"]);
888        assert_eq!(snapshot.capture, CaptureState::Captured);
889        assert_eq!(
890            sink.writes,
891            vec![b"one\n".to_vec(), b"two\n".to_vec(), b"three".to_vec()]
892        );
893    }
894
895    #[tokio::test]
896    async fn a_read_failure_keeps_prior_lines_and_marks_the_capture_incomplete() {
897        let ring = shared(10, 10_000, 128);
898        let source = FailingReader {
899            bytes: b"crash cause\n".to_vec(),
900            emitted: false,
901        };
902        let mut sink = RecordingSink::default();
903        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
904
905        let snapshot = lock_ring(&ring).snapshot(None, None);
906        assert_eq!(lines(&snapshot), vec!["crash cause"]);
907        assert!(matches!(
908            snapshot.capture,
909            CaptureState::Incomplete { ref reason } if reason.contains("reader failed")
910        ));
911        assert_eq!(sink.writes, vec![b"crash cause\n".to_vec()]);
912    }
913
914    #[tokio::test]
915    async fn every_captured_line_is_also_forwarded() {
916        // Forwarding is not optional. The daemon log is overwhelmingly module
917        // output; a tap that captured without forwarding would leave it nearly
918        // empty and every existing reader would report clean on nothing.
919        let ring = shared(10, 10_000, 128);
920        let source = std::io::Cursor::new(b"alpha\nbeta\n".to_vec());
921        let mut sink = RecordingSink::default();
922        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
923
924        assert_eq!(sink.writes, vec![b"alpha\n".to_vec(), b"beta\n".to_vec()]);
925    }
926
927    #[tokio::test]
928    async fn each_forwarded_line_is_exactly_one_write() {
929        // Inheriting the fd gave line atomicity for free. Reading a pipe and
930        // re-emitting can split a line that used to be atomic, so the framing
931        // must be one syscall per complete line -- asserted as one write per
932        // line, not merely as correct bytes.
933        let ring = shared(10, 10_000, 128);
934        let source = std::io::Cursor::new(b"first\nsecond\nthird\n".to_vec());
935        let mut sink = RecordingSink::default();
936        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
937
938        assert_eq!(sink.writes.len(), 3);
939        for write in &sink.writes {
940            assert_eq!(
941                write.iter().filter(|byte| **byte == b'\n').count(),
942                1,
943                "a write carried something other than exactly one complete line"
944            );
945            assert_eq!(*write.last().unwrap(), b'\n');
946        }
947    }
948
949    #[test]
950    fn the_first_process_start_is_not_recorded_because_it_divides_nothing() {
951        // Otherwise a module that printed nothing renders as a lone boundary
952        // marker, and every caller has to decide whether that counts as silence.
953        let mut ring = ring(10, 10_000, 128);
954        ring.push_process_start();
955        assert!(ring.snapshot(None, None).entries.is_empty());
956
957        ring.push_line("first process said this");
958        ring.push_process_start();
959        assert!(
960            matches!(ring.entries.back(), Some(TailEntry::ProcessStart)),
961            "a boundary with output before it must be recorded"
962        );
963    }
964
965    #[test]
966    fn a_process_start_is_recorded_when_only_dropped_lines_precede_it() {
967        // The ring can be non-empty in the sense that matters -- lines were
968        // written and evicted -- while `entries` is empty. Suppressing the
969        // boundary there would attribute surviving output to the wrong process.
970        let mut ring = ring(1, 10_000, 128);
971        ring.push_line("evicted");
972        ring.push_line("also evicted");
973        // Emptying `entries` by hand must also zero the running totals kept
974        // beside it, or the ring holds counts for lines it no longer has and
975        // any later eviction decision is made against the stale numbers.
976        ring.entries.clear();
977        ring.lines = 0;
978        ring.bytes = 0;
979        ring.push_process_start();
980        assert!(matches!(
981            ring.entries.front(),
982            Some(TailEntry::ProcessStart)
983        ));
984    }
985
986    #[tokio::test]
987    async fn the_pump_marks_captured_even_when_the_module_writes_nothing() {
988        // Clean EOF with no output is a module that was quiet, not one nobody
989        // listened to -- and the two must not render alike.
990        let ring = shared(10, 10_000, 128);
991        let source = std::io::Cursor::new(Vec::new());
992        let mut sink = RecordingSink::default();
993        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
994
995        let snapshot = lock_ring(&ring).snapshot(None, None);
996        assert!(snapshot.entries.is_empty());
997        assert_eq!(snapshot.capture, CaptureState::Captured);
998        assert!(sink.writes.is_empty());
999    }
1000
1001    #[tokio::test]
1002    async fn a_line_with_no_newline_cannot_grow_the_reader_without_bound() {
1003        // A module fault must not become a daemon fault: without the pending
1004        // ceiling this buffer grows to whatever the module writes.
1005        let ring = shared(10, 10_000_000, 4 * 1024 * 1024);
1006        let source = std::io::Cursor::new(vec![b'x'; MAX_PENDING_LINE_BYTES + 4096]);
1007        let mut sink = RecordingSink::default();
1008        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1009
1010        let snapshot = lock_ring(&ring).snapshot(None, None);
1011        assert_eq!(
1012            lines(&snapshot).len(),
1013            2,
1014            "expected a forced flush at the ceiling plus the remainder"
1015        );
1016        assert_eq!(
1017            sink.writes,
1018            vec![vec![b'x'; MAX_PENDING_LINE_BYTES], vec![b'x'; 4096],],
1019            "forced flushes and EOF fragments must not invent delimiters"
1020        );
1021    }
1022
1023    #[tokio::test]
1024    async fn boundaries_truncation_and_framing_do_not_depend_on_chunk_splits() {
1025        // The same stream split at hostile boundaries -- mid-line, between a CR
1026        // and its LF, and a line sitting exactly on the per-line cap -- must
1027        // produce the same ring entries and forwarded bytes as any other split.
1028        let ring = shared(100, 100_000, 8);
1029        let source = ChunkedReader {
1030            chunks: vec![
1031                b"fir".to_vec(),
1032                b"st\nsec".to_vec(),
1033                b"ond\ncarry\r".to_vec(),
1034                b"\nover\n".to_vec(),
1035                b"12345678\n".to_vec(),
1036                b"1234567".to_vec(),
1037                b"89\n".to_vec(),
1038                b"tail".to_vec(),
1039            ]
1040            .into_iter()
1041            .collect(),
1042        };
1043        let mut sink = RecordingSink::default();
1044        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1045
1046        let snapshot = lock_ring(&ring).snapshot(None, None);
1047        assert_eq!(snapshot.capture, CaptureState::Captured);
1048        assert_eq!(
1049            snapshot.entries,
1050            vec![
1051                TailEntry::Line {
1052                    text: "first".to_string(),
1053                    truncated: false
1054                },
1055                TailEntry::Line {
1056                    text: "second".to_string(),
1057                    truncated: false
1058                },
1059                // The pump delimits on '\n' alone; a CR belongs to the line body.
1060                TailEntry::Line {
1061                    text: "carry\r".to_string(),
1062                    truncated: false
1063                },
1064                TailEntry::Line {
1065                    text: "over".to_string(),
1066                    truncated: false
1067                },
1068                // Exactly at the per-line cap: kept whole.
1069                TailEntry::Line {
1070                    text: "12345678".to_string(),
1071                    truncated: false
1072                },
1073                // One byte past the cap: cut, and marked as cut.
1074                TailEntry::Line {
1075                    text: "12345678".to_string(),
1076                    truncated: true
1077                },
1078                TailEntry::Line {
1079                    text: "tail".to_string(),
1080                    truncated: false
1081                },
1082            ]
1083        );
1084        assert_eq!(
1085            sink.writes,
1086            vec![
1087                b"first\n".to_vec(),
1088                b"second\n".to_vec(),
1089                b"carry\r\n".to_vec(),
1090                b"over\n".to_vec(),
1091                b"12345678\n".to_vec(),
1092                b"123456789\n".to_vec(),
1093                b"tail".to_vec(),
1094            ]
1095        );
1096    }
1097
1098    #[tokio::test]
1099    async fn a_line_with_no_newline_is_not_rescanned_from_byte_zero_on_every_chunk() {
1100        // One 1 MiB line arrives in 8192-byte reads. Searching the whole
1101        // pending buffer for a newline on every chunk scans each byte once per
1102        // chunk that arrived after it -- about 64 MiB examined per MiB of
1103        // output. Searching only the bytes that arrived since the last search
1104        // scans each byte once.
1105        let input = vec![b'x'; MAX_PENDING_LINE_BYTES + 4096];
1106        let ring = shared(10, 10_000_000, 4 * 1024 * 1024);
1107        let source = std::io::Cursor::new(input.clone());
1108        let mut sink = RecordingSink::default();
1109
1110        take_scanned_bytes();
1111        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1112        let scanned = take_scanned_bytes();
1113
1114        assert!(
1115            scanned <= 2 * input.len(),
1116            "newline searches examined {scanned} bytes for {} bytes of input; \
1117             each chunk must search only newly arrived bytes",
1118            input.len()
1119        );
1120    }
1121
1122    #[test]
1123    fn a_byte_limit_smaller_than_one_line_still_returns_that_line() {
1124        // Returning nothing would be indistinguishable from a quiet module, which
1125        // is the failure this module exists to prevent.
1126        let mut ring = ring(10, 10_000, 128);
1127        ring.mark_captured();
1128        ring.push_line("a line considerably longer than the request limit");
1129        let snapshot = ring.snapshot(None, Some(4));
1130        assert_eq!(snapshot.entries.len(), 1);
1131    }
1132
1133    #[test]
1134    fn an_incoherent_config_clamps_the_line_cap_and_keeps_its_restart_boundary() {
1135        let config = StderrTailConfig::new(2, 10, 100);
1136        assert_eq!(config.max_line_bytes, config.max_bytes);
1137        let mut ring = StderrRing::new(config);
1138        ring.mark_captured();
1139        ring.push_line("old");
1140        ring.push_process_start();
1141        ring.push_line("new process line longer than the ring byte cap");
1142
1143        assert_eq!(
1144            ring.snapshot(None, None).entries,
1145            vec![
1146                TailEntry::ProcessStart,
1147                TailEntry::Line {
1148                    text: "new proces".to_string(),
1149                    truncated: true,
1150                },
1151            ]
1152        );
1153    }
1154}