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::{BTreeMap, VecDeque};
22use std::io::{self, Write};
23use std::path::{Path, PathBuf};
24use std::sync::{
25    atomic::{AtomicBool, Ordering},
26    Arc, Mutex,
27};
28use std::time::{SystemTime, UNIX_EPOCH};
29
30use tokio::io::AsyncReadExt;
31
32/// Longest single line admitted to the ring before truncation.
33///
34/// A module emitting a 40 MB backtrace on one line satisfies any line-count cap
35/// while evicting everything else -- the pathological emitter wins twice, once by
36/// filling the ring and once by being unreadable itself. Truncating on the way in
37/// costs that emitter one line instead of the whole tail.
38pub const DEFAULT_MAX_LINE_BYTES: usize = 2048;
39
40/// Lines retained per module.
41pub const DEFAULT_MAX_LINES: usize = 200;
42
43/// Total bytes retained per module, across all lines.
44///
45/// Both this and [`DEFAULT_MAX_LINES`] apply; whichever binds first wins. A line
46/// cap alone is satisfied by 200 lines of 2 KB, which is not a budget worth
47/// holding for fifteen modules.
48pub const DEFAULT_MAX_BYTES: usize = 64 * 1024;
49
50/// Whether a module's stderr is being captured, and if not, why not.
51///
52/// Typed rather than nullable so `NotCaptured` has to be handled rather than
53/// defaulted past. An empty tail and an uncaptured one send an operator in
54/// opposite directions -- one says the module printed nothing before dying, the
55/// other says nobody was listening -- and rendering them alike is the defect this
56/// module exists to fix, reproduced one layer up.
57#[derive(Debug, Clone, PartialEq, Eq)]
58pub enum CaptureState {
59    /// A reader is attached, or was attached and reached clean EOF.
60    Captured,
61    /// Retained entries are valid, but capture ended before clean EOF, or the
62    /// pipe of a process the supervisor has already moved on from has not
63    /// reached EOF yet. The second kind clears when that pipe does reach EOF,
64    /// because from then on nothing that process wrote is missing.
65    Incomplete { reason: String },
66    /// No reader was attached. The tail says nothing about what the module wrote.
67    NotCaptured { reason: String },
68}
69
70/// One retained entry.
71///
72/// Boundaries are in-band rather than a separate field because their position
73/// relative to the lines is the whole point: "these three lines came from the
74/// process that died, those came from its replacement" is unanswerable from a
75/// count.
76#[derive(Debug, Clone, PartialEq, Eq)]
77pub enum TailEntry {
78    Line {
79        text: String,
80        /// This line was cut at the per-line cap.
81        truncated: bool,
82        /// Wall-clock Unix milliseconds at which the reader framed this line,
83        /// the same instant stamped on its capture-file line. `None` for a line
84        /// admitted without one (see [`StderrRing::push_line`]).
85        at_ms: Option<u64>,
86    },
87    /// The supervisor spawned a new process for this module. Lines after this
88    /// entry come from the new one.
89    ProcessStart,
90}
91
92impl TailEntry {
93    fn cost(&self) -> usize {
94        match self {
95            Self::Line { text, .. } => text.len(),
96            Self::ProcessStart => 0,
97        }
98    }
99}
100
101/// One stored entry. Unlike [`TailEntry`], a boundary remembers which process
102/// generation it starts, so a line that arrives late from an older process can
103/// be put back in that process's section instead of after its successor's
104/// boundary.
105#[derive(Debug, Clone, PartialEq, Eq)]
106enum Slot {
107    Line {
108        text: String,
109        truncated: bool,
110        at_ms: Option<u64>,
111    },
112    ProcessStart {
113        generation: u64,
114    },
115}
116
117impl Slot {
118    fn cost(&self) -> usize {
119        match self {
120            Self::Line { text, .. } => text.len(),
121            Self::ProcessStart { .. } => 0,
122        }
123    }
124}
125
126/// Where the stderr reader of one process generation stands.
127#[derive(Debug, Clone, PartialEq, Eq)]
128enum PumpPhase {
129    /// The process is the supervisor's current concern; its lines go at the end.
130    Attached,
131    /// The supervisor has moved on from the process, but its pipe is still open.
132    /// Lines it still delivers belong in its own section.
133    Retired,
134    /// Retired, and the pipe was still open when the supervisor stopped waiting
135    /// for it. Reported as `Incomplete` until the pipe reaches EOF.
136    Late { reason: String },
137}
138
139#[derive(Debug, Clone, Copy, PartialEq, Eq)]
140pub struct StderrTailConfig {
141    max_lines: usize,
142    max_bytes: usize,
143    max_line_bytes: usize,
144}
145
146impl StderrTailConfig {
147    /// Keeps every retained line within the ring's total byte budget.
148    /// Clamp rather than reject so diagnostics degrade without blocking supervisor startup.
149    pub const fn new(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> Self {
150        Self {
151            max_lines,
152            max_bytes,
153            max_line_bytes: if max_line_bytes > max_bytes {
154                max_bytes
155            } else {
156                max_line_bytes
157            },
158        }
159    }
160}
161
162impl Default for StderrTailConfig {
163    fn default() -> Self {
164        Self::new(DEFAULT_MAX_LINES, DEFAULT_MAX_BYTES, DEFAULT_MAX_LINE_BYTES)
165    }
166}
167
168/// A module's retained stderr, oldest first.
169#[derive(Debug, Clone, PartialEq, Eq)]
170pub struct StderrTailSnapshot {
171    pub capture: CaptureState,
172    pub entries: Vec<TailEntry>,
173    /// Lines evicted since the module was first supervised.
174    ///
175    /// Non-zero means the tail starts mid-stream. That is the ring working as
176    /// intended, but a reader diagnosing a crash needs to know the first retained
177    /// line is not the first line the module wrote -- otherwise an absent cause
178    /// reads as a module that never explained itself.
179    pub dropped_lines: u64,
180}
181
182impl StderrTailSnapshot {
183    /// The uncaptured case, for a module whose stderr was never piped.
184    pub fn not_captured(reason: impl Into<String>) -> Self {
185        Self {
186            capture: CaptureState::NotCaptured {
187                reason: reason.into(),
188            },
189            entries: Vec::new(),
190            dropped_lines: 0,
191        }
192    }
193}
194
195/// Bounded ring of a single module's stderr lines.
196///
197/// Survives respawn deliberately. The stderr explaining an exit is written
198/// *before* that exit, so clearing on restart would discard the lines exactly
199/// when they become the thing being asked for. [`TailEntry::ProcessStart`] keeps
200/// the generations distinguishable instead.
201///
202/// A process's stderr reader can outlive the supervisor's interest in it: the
203/// reader may not have been scheduled yet when the process exited, or a
204/// descendant may still hold the pipe open. Lines such a reader delivers after
205/// the next process started are placed before that next process's boundary, so
206/// the section a line appears in always names the process that wrote it.
207#[derive(Debug)]
208pub struct StderrRing {
209    config: StderrTailConfig,
210    entries: VecDeque<Slot>,
211    // Count of `TailEntry::Line` entries, kept running because eviction checks
212    // it on every push and recounting would walk the whole ring under the
213    // mutex each time.
214    lines: usize,
215    bytes: usize,
216    dropped_lines: u64,
217    capture: CaptureState,
218    /// Generation of the newest process boundary; 0 before the first.
219    generation: u64,
220    /// Readers that have not reached EOF yet, by the generation they read for.
221    pumps: BTreeMap<u64, PumpPhase>,
222    /// Highest generation whose boundary was evicted from the front. A late line
223    /// from an older generation belongs in front of that boundary, which is
224    /// evicted territory, so it is counted as dropped rather than stored.
225    evicted_through: u64,
226}
227
228impl StderrRing {
229    pub fn new(config: StderrTailConfig) -> Self {
230        Self {
231            config,
232            entries: VecDeque::new(),
233            lines: 0,
234            bytes: 0,
235            dropped_lines: 0,
236            // Until a reader attaches, the honest answer is that nothing is
237            // listening -- not that the module has been quiet.
238            capture: CaptureState::NotCaptured {
239                reason: "stderr reader has not started".to_string(),
240            },
241            generation: 0,
242            pumps: BTreeMap::new(),
243            evicted_through: 0,
244        }
245    }
246
247    /// Generation of the newest process boundary.
248    pub fn generation(&self) -> u64 {
249        self.generation
250    }
251
252    pub fn mark_captured(&mut self) {
253        if matches!(self.capture, CaptureState::NotCaptured { .. }) {
254            self.capture = CaptureState::Captured;
255        }
256    }
257
258    pub fn mark_incomplete(&mut self, reason: impl Into<String>) {
259        self.capture = CaptureState::Incomplete {
260            reason: reason.into(),
261        };
262    }
263
264    pub fn mark_not_captured(&mut self, reason: impl Into<String>) {
265        self.capture = CaptureState::NotCaptured {
266            reason: reason.into(),
267        };
268    }
269
270    /// Record that a new process was spawned for this module, returning the
271    /// generation number its lines are attributed to.
272    ///
273    /// A boundary separates output on either side of it, so one with nothing
274    /// before it separates nothing: on the FIRST spawn it would make a module
275    /// that printed nothing render as a marker rather than as empty, and the
276    /// caller then has to decide whether a one-marker tail counts as silence.
277    /// The boundary is still stored, because a late line from an older process
278    /// needs it to find its section, but [`Self::snapshot`] shows it only once
279    /// there is output before it to divide, which keeps "captured and empty"
280    /// literally empty.
281    pub fn push_process_start(&mut self) -> u64 {
282        self.generation += 1;
283        let generation = self.generation;
284        // Two boundaries in a row mean the process between them has written
285        // nothing retained. The earlier one can go only if no reader from its
286        // generation onward is still open: such a reader may yet deliver a line
287        // that belongs between the two. Without this a module that restarts
288        // silently would grow the ring by one boundary per restart forever.
289        if let Some(Slot::ProcessStart {
290            generation: previous,
291        }) = self.entries.back()
292        {
293            if self.pumps.range(*previous..).next().is_none() {
294                self.entries.pop_back();
295            }
296        }
297        self.push_entry(Slot::ProcessStart { generation });
298        generation
299    }
300
301    /// [`Self::push_process_start`] for a process whose stderr reader is
302    /// attached: the reader is tracked until it calls [`Self::finish_pump`].
303    pub(crate) fn begin_process(&mut self) -> u64 {
304        let generation = self.push_process_start();
305        self.pumps.insert(generation, PumpPhase::Attached);
306        generation
307    }
308
309    /// The supervisor has moved on from this generation's process. Its reader
310    /// keeps running, and any line it still delivers is kept in its own section.
311    pub(crate) fn retire_pump(&mut self, generation: u64) {
312        if let Some(phase @ PumpPhase::Attached) = self.pumps.get_mut(&generation) {
313            *phase = PumpPhase::Retired;
314        }
315    }
316
317    /// The supervisor stopped waiting for this generation's reader before its
318    /// pipe reached EOF. The capture reads as `Incomplete` with `reason` until
319    /// the reader finishes. A reader that already finished is left alone: it
320    /// got everything.
321    pub(crate) fn mark_pump_late(&mut self, generation: u64, reason: impl Into<String>) {
322        if let Some(phase) = self.pumps.get_mut(&generation) {
323            *phase = PumpPhase::Late {
324                reason: reason.into(),
325            };
326        }
327    }
328
329    /// This generation's reader has stopped: at EOF, or on a read error that
330    /// has already been recorded with [`Self::mark_incomplete`].
331    pub(crate) fn finish_pump(&mut self, generation: u64) {
332        self.pumps.remove(&generation);
333    }
334
335    /// Admit one complete line, truncating it if it exceeds the per-line cap.
336    ///
337    /// `line` must not contain a trailing newline; the reader strips it so the
338    /// stored text and the byte accounting agree. The line carries no capture
339    /// time: only the pipe reader knows when a line arrived, and it records
340    /// that through [`Self::push_line_from_at`].
341    pub fn push_line(&mut self, line: &str) {
342        self.push_line_from(self.generation, line);
343    }
344
345    /// [`Self::push_line`] for a line read from `generation`'s pipe.
346    pub(crate) fn push_line_from(&mut self, generation: u64, line: &str) {
347        self.admit_line(generation, line, None);
348    }
349
350    /// [`Self::push_line_from`] for a line the reader framed at `at_ms`
351    /// (wall-clock Unix milliseconds).
352    pub(crate) fn push_line_from_at(&mut self, generation: u64, line: &str, at_ms: u64) {
353        self.admit_line(generation, line, Some(at_ms));
354    }
355
356    /// A line from a retired process that arrives after a newer process started
357    /// goes in front of the first boundary newer than its own generation. Lines
358    /// from a process the supervisor has not retired (the incumbent during a
359    /// swap's overlap) go at the end, as they arrive.
360    fn admit_line(&mut self, generation: u64, line: &str, at_ms: Option<u64>) {
361        let (text, truncated) = truncate_line(line, self.config.max_line_bytes);
362        let slot = Slot::Line {
363            text,
364            truncated,
365            at_ms,
366        };
367        let retired = matches!(
368            self.pumps.get(&generation),
369            Some(PumpPhase::Retired | PumpPhase::Late { .. })
370        );
371        if !retired || generation >= self.generation {
372            self.push_entry(slot);
373            return;
374        }
375        if self.evicted_through > generation {
376            // The section this line belongs to has been evicted, so the line is
377            // older than everything retained.
378            self.dropped_lines += 1;
379            return;
380        }
381        let index = self.entries.iter().position(
382            |slot| matches!(slot, Slot::ProcessStart { generation: start } if *start > generation),
383        );
384        match index {
385            Some(index) => self.insert_entry(index, slot),
386            None => self.push_entry(slot),
387        }
388    }
389
390    fn push_entry(&mut self, entry: Slot) {
391        self.insert_entry(self.entries.len(), entry);
392    }
393
394    fn insert_entry(&mut self, index: usize, entry: Slot) {
395        self.bytes += entry.cost();
396        if matches!(entry, Slot::Line { .. }) {
397            self.lines += 1;
398        }
399        self.entries.insert(index, entry);
400        self.evict_to_fit();
401    }
402
403    fn evict_to_fit(&mut self) {
404        while self.lines > self.config.max_lines
405            || (self.bytes > self.config.max_bytes && self.entries.len() > 1)
406        {
407            let Some(evicted) = self.entries.pop_front() else {
408                break;
409            };
410            self.bytes -= evicted.cost();
411            match evicted {
412                Slot::Line { .. } => {
413                    self.lines -= 1;
414                    self.dropped_lines += 1;
415                }
416                Slot::ProcessStart { generation } => {
417                    self.evicted_through = self.evicted_through.max(generation);
418                }
419            }
420        }
421    }
422
423    /// The most recent entries, oldest first, bounded by the caller's limits.
424    ///
425    /// `max_lines`/`max_bytes` narrow the ring's own caps; they cannot widen them.
426    pub fn snapshot(
427        &self,
428        max_lines: Option<usize>,
429        max_bytes: Option<usize>,
430    ) -> StderrTailSnapshot {
431        let line_limit = max_lines.unwrap_or(self.config.max_lines);
432        let byte_limit = max_bytes.unwrap_or(self.config.max_bytes);
433
434        // The stored boundaries, reduced to the ones worth showing: a boundary
435        // with no output before it (retained or evicted) divides nothing, and
436        // two in a row say no more than one.
437        let mut visible: Vec<TailEntry> = Vec::with_capacity(self.entries.len());
438        let mut output_before = self.dropped_lines > 0;
439        for slot in &self.entries {
440            match slot {
441                Slot::Line {
442                    text,
443                    truncated,
444                    at_ms,
445                } => {
446                    visible.push(TailEntry::Line {
447                        text: text.clone(),
448                        truncated: *truncated,
449                        at_ms: *at_ms,
450                    });
451                    output_before = true;
452                }
453                Slot::ProcessStart { .. } => {
454                    if output_before && !matches!(visible.last(), Some(TailEntry::ProcessStart)) {
455                        visible.push(TailEntry::ProcessStart);
456                    }
457                }
458            }
459        }
460
461        let mut taken: Vec<TailEntry> = Vec::new();
462        let mut bytes = 0usize;
463        let mut lines = 0usize;
464        // Walk backwards: a tail is anchored at the newest end, so a caller
465        // asking for 20 lines wants the last 20, not the first 20.
466        for entry in visible.iter().rev() {
467            match entry {
468                TailEntry::Line { .. } => {
469                    if lines >= line_limit {
470                        break;
471                    }
472                    let cost = entry.cost();
473                    if lines > 0 && bytes + cost > byte_limit {
474                        break;
475                    }
476                    bytes += cost;
477                    lines += 1;
478                    taken.push(entry.clone());
479                }
480                TailEntry::ProcessStart if lines > 0 => taken.push(entry.clone()),
481                TailEntry::ProcessStart => {}
482            }
483        }
484        taken.reverse();
485
486        let withheld = self.lines.saturating_sub(lines);
487
488        // A retired process whose pipe the supervisor stopped waiting for may
489        // still be writing; until its reader reaches EOF the tail cannot claim
490        // to hold everything. A permanent state (a read failure, no pipe at
491        // all) is the more specific fact and is reported as is.
492        let late = self.pumps.values().find_map(|phase| match phase {
493            PumpPhase::Late { reason } => Some(reason),
494            _ => None,
495        });
496        let capture = match (&self.capture, late) {
497            (CaptureState::Captured, Some(reason)) => CaptureState::Incomplete {
498                reason: reason.clone(),
499            },
500            (capture, _) => capture.clone(),
501        };
502
503        StderrTailSnapshot {
504            capture,
505            entries: taken,
506            // Lines the ring evicted plus lines this request's own limits held
507            // back. Both mean the same thing to the reader -- the text above is
508            // not the beginning -- and separating them would invite treating a
509            // narrow request as evidence of a quiet module.
510            dropped_lines: self.dropped_lines + withheld as u64,
511        }
512    }
513}
514
515/// Reassembly buffer ceiling for a line with no newline in sight.
516///
517/// The ring truncates what it stores, but the READER has to hold the bytes until
518/// it finds a delimiter. A module emitting a gigabyte with no newline would grow
519/// this buffer without bound and take the daemon down with it -- a module fault
520/// escalating into a fleet fault, which is exactly what supervision exists to
521/// prevent. At this ceiling the pending bytes are flushed as a line and
522/// reassembly restarts.
523const MAX_PENDING_LINE_BYTES: usize = 1024 * 1024;
524
525/// Read a child's stderr to EOF, retaining the bounded crash tail and forwarding
526/// every complete line to the selected capture sink.
527pub async fn pump_stderr<R>(source: R, ring: Arc<Mutex<StderrRing>>)
528where
529    R: AsyncReadExt + Unpin,
530{
531    pump_stderr_into(source, ring, &mut StderrSink).await
532}
533
534/// Shared destination for a child's stdout and stderr pumps.
535///
536/// The file keeps the historical `.stderr.log` name even though it carries both
537/// streams; the stable name is part of the operator contract. Both pumps share
538/// one mutex, and cortexkit-log writes each framed line in one call, so partial
539/// lines from the two pipes cannot interleave.
540#[derive(Clone)]
541pub(crate) enum ChildOutputSink {
542    File {
543        sink: Arc<Mutex<cortexkit_log::LineSink>>,
544        path: Arc<PathBuf>,
545        failure_reported: Arc<AtomicBool>,
546    },
547    Stderr,
548}
549
550impl ChildOutputSink {
551    pub(crate) fn open(path: &Path, retention: cortexkit_log::Retention) -> io::Result<Self> {
552        Ok(Self::File {
553            sink: Arc::new(Mutex::new(cortexkit_log::LineSink::open(path, retention)?)),
554            path: Arc::new(path.to_path_buf()),
555            failure_reported: Arc::new(AtomicBool::new(false)),
556        })
557    }
558}
559
560/// Where forwarded complete lines go. It exists so tests can observe framing
561/// and so production can serialize the two child pipes through one file sink.
562pub trait OutputSink {
563    fn write_line(&mut self, line: &[u8]);
564
565    /// Whether each write should begin with its capture time (see
566    /// [`format_capture_stamp`]).
567    ///
568    /// Only a sink that ends every write with a newline may say yes: then each
569    /// write is a whole line of its own, and a stamp at the front of the write
570    /// is at the front of a line. A sink that passes an unterminated piece
571    /// through as-is would have the next piece continue the same line, and a
572    /// stamp there would land in the middle of it.
573    fn stamps_lines(&self) -> bool {
574        false
575    }
576}
577
578struct StderrSink;
579
580impl OutputSink for StderrSink {
581    fn write_line(&mut self, line: &[u8]) {
582        let stderr = std::io::stderr();
583        let mut handle = stderr.lock();
584        let _ = handle.write_all(line);
585    }
586}
587
588impl OutputSink for ChildOutputSink {
589    fn write_line(&mut self, line: &[u8]) {
590        match self {
591            Self::File {
592                sink,
593                path,
594                failure_reported,
595            } => {
596                let result = sink
597                    .lock()
598                    .unwrap_or_else(|poisoned| poisoned.into_inner())
599                    .write_line(line);
600                if let Err(error) = result {
601                    if !failure_reported.swap(true, Ordering::Relaxed) {
602                        tracing::warn!(
603                            path = %path.display(),
604                            error = %error,
605                            "child output capture write failed; later failures are suppressed"
606                        );
607                    }
608                }
609            }
610            Self::Stderr => StderrSink.write_line(line),
611        }
612    }
613
614    // The capture file is read long after it was written, often next to the
615    // daemon's own log, whose lines carry a time. `LineSink::write_line`
616    // appends a newline to any write that lacks one, so every write here is a
617    // whole line and may carry the stamp. The daemon's inherited stderr gets
618    // the module's bytes unchanged: an unterminated piece is continued there by
619    // the next one, and whatever collects that stream stamps it itself.
620    fn stamps_lines(&self) -> bool {
621        matches!(self, Self::File { .. })
622    }
623}
624
625/// Read one process generation's stderr to EOF. `generation` is the value
626/// [`StderrRing::begin_process`] returned for that process; it decides which
627/// section of the ring a line lands in if the reader outlives the process.
628pub(crate) async fn pump_stderr_to<R, S>(
629    source: R,
630    ring: Arc<Mutex<StderrRing>>,
631    generation: u64,
632    mut sink: S,
633) where
634    R: AsyncReadExt + Unpin,
635    S: OutputSink,
636{
637    pump_lines_into(source, Some((&ring, generation)), &mut sink, "stderr").await;
638}
639
640pub(crate) async fn pump_stdout_to<R>(source: R, mut sink: ChildOutputSink)
641where
642    R: AsyncReadExt + Unpin,
643{
644    pump_lines_into(source, None, &mut sink, "stdout").await;
645}
646
647/// Read stderr for whichever process generation is newest when the reader
648/// starts.
649async fn pump_stderr_into<R, S>(source: R, ring: Arc<Mutex<StderrRing>>, sink: &mut S)
650where
651    R: AsyncReadExt + Unpin,
652    S: OutputSink,
653{
654    let generation = lock_ring(&ring).generation();
655    pump_lines_into(source, Some((&ring, generation)), sink, "stderr").await;
656}
657
658async fn pump_lines_into<R, S>(
659    mut source: R,
660    ring: Option<(&Arc<Mutex<StderrRing>>, u64)>,
661    sink: &mut S,
662    stream_name: &str,
663) where
664    R: AsyncReadExt + Unpin,
665    S: OutputSink,
666{
667    if let Some((ring, _)) = ring {
668        lock_ring(ring).mark_captured();
669    }
670
671    let mut pending: Vec<u8> = Vec::new();
672    // Bytes before `scanned_upto` are already known to hold no newline; searching
673    // them again would rescan the whole buffer on every chunk -- for a line with
674    // no newline that is about 64 MiB examined per MiB of module output.
675    let mut scanned_upto = 0usize;
676    // Bytes before `cursor` were emitted as complete lines. They are removed in
677    // one compaction per chunk rather than shifting the buffer once per line.
678    let mut cursor = 0usize;
679    let mut chunk = [0u8; 8192];
680    loop {
681        let read = match source.read(&mut chunk).await {
682            Ok(0) => break,
683            Ok(n) => n,
684            Err(error) => {
685                if let Some((ring, generation)) = ring {
686                    let mut ring = lock_ring(ring);
687                    ring.mark_incomplete(format!("{stream_name} read failed: {error}"));
688                    ring.finish_pump(generation);
689                } else {
690                    tracing::warn!(stream = stream_name, error = %error, "child output capture read failed");
691                }
692                return;
693            }
694        };
695        pending.extend_from_slice(&chunk[..read]);
696
697        while let Some(relative) = find_newline(&pending[scanned_upto..]) {
698            let newline = scanned_upto + relative;
699            emit_line(ring, sink, &pending[cursor..newline], true);
700            cursor = newline + 1;
701            scanned_upto = cursor;
702        }
703        scanned_upto = pending.len();
704
705        if cursor > 0 {
706            pending.drain(..cursor);
707            scanned_upto -= cursor;
708            cursor = 0;
709        }
710
711        if pending.len() >= MAX_PENDING_LINE_BYTES {
712            let line = std::mem::take(&mut pending);
713            emit_line(ring, sink, &line, false);
714            scanned_upto = 0;
715        }
716    }
717
718    if !pending.is_empty() {
719        emit_line(ring, sink, &pending, false);
720    }
721    if let Some((ring, generation)) = ring {
722        lock_ring(ring).finish_pump(generation);
723    }
724}
725
726// Bytes examined by newline searches, summed across a pump. Tests use this to
727// assert the reader does not rescan bytes it already knows contain no newline.
728// Per thread, not process-wide: other tests pump concurrently on their own
729// threads, and a shared counter picked up their searches too, failing the
730// bound by a few bytes under a parallel run. `#[tokio::test]` runs its body
731// and every future it awaits on one thread, so a pump's searches land on the
732// test's own counter.
733#[cfg(test)]
734thread_local! {
735    static SCANNED_BYTES: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
736}
737
738#[cfg(test)]
739fn take_scanned_bytes() -> usize {
740    SCANNED_BYTES.with(|scanned| scanned.replace(0))
741}
742
743/// Locate the next newline in `haystack`, counting the bytes examined so a
744/// test can observe how much of the pending buffer each search walks.
745fn find_newline(haystack: &[u8]) -> Option<usize> {
746    let found = memchr::memchr(b'\n', haystack);
747    #[cfg(test)]
748    SCANNED_BYTES.with(|scanned| {
749        scanned.set(scanned.get() + found.map(|index| index + 1).unwrap_or(haystack.len()));
750    });
751    found
752}
753
754fn emit_line<S: OutputSink>(
755    ring: Option<(&Arc<Mutex<StderrRing>>, u64)>,
756    sink: &mut S,
757    raw: &[u8],
758    terminated: bool,
759) {
760    // One instant for both destinations, taken here because this is where the
761    // line is complete. The module's bytes do not say when the line arrived,
762    // so the time is stored beside it: as `at_ms` in the ring and as the stamp
763    // in the file. A time assigned later, when either is read, would look
764    // recorded while being off by however long the line sat there.
765    let at_ms = unix_ms(SystemTime::now());
766    if let Some((ring, generation)) = ring {
767        lock_ring(ring).push_line_from_at(generation, &String::from_utf8_lossy(raw), at_ms);
768    }
769
770    // Framed and written in ONE call, stamp included. Two writes would let the
771    // other pipe land between the pieces.
772    let stamp = sink.stamps_lines().then(|| format_capture_stamp(at_ms));
773    if stamp.is_none() && !terminated {
774        sink.write_line(raw);
775        return;
776    }
777    let mut framed = Vec::with_capacity(CAPTURE_STAMP_PREFIX_LEN + raw.len() + 1);
778    if let Some(stamp) = stamp {
779        framed.extend_from_slice(stamp.as_bytes());
780        framed.push(b' ');
781    }
782    framed.extend_from_slice(raw);
783    if terminated {
784        framed.push(b'\n');
785    }
786    sink.write_line(&framed);
787}
788
789fn unix_ms(at: SystemTime) -> u64 {
790    // A clock set before 1970 is already wrong about every time it reports;
791    // saturating keeps the stamp well-formed rather than failing the write.
792    at.duration_since(UNIX_EPOCH)
793        .map(|since| u64::try_from(since.as_millis()).unwrap_or(u64::MAX))
794        .unwrap_or(0)
795}
796
797/// Length of a capture stamp, `2026-09-19T07:04:00.685Z`.
798pub const CAPTURE_STAMP_LEN: usize = 24;
799
800/// Length of the prefix on a capture-file line: the stamp and one space.
801pub const CAPTURE_STAMP_PREFIX_LEN: usize = CAPTURE_STAMP_LEN + 1;
802
803/// `at_ms` (Unix milliseconds) as RFC 3339 UTC with milliseconds and `Z`,
804/// the form the daemon's own log lines begin with, so one parser reads the
805/// time of a line in either file. Always [`CAPTURE_STAMP_LEN`] bytes for any
806/// time before the year 10000.
807pub fn format_capture_stamp(at_ms: u64) -> String {
808    let seconds = at_ms / 1000;
809    let millis = at_ms % 1000;
810    let days = seconds / 86_400;
811    let of_day = seconds % 86_400;
812    let (year, month, day) = civil_from_days(days);
813    format!(
814        "{year:04}-{month:02}-{day:02}T{:02}:{:02}:{:02}.{millis:03}Z",
815        of_day / 3600,
816        (of_day % 3600) / 60,
817        of_day % 60,
818    )
819}
820
821/// Split a capture-file line into the time its stamp records (Unix
822/// milliseconds) and the module's bytes after the stamp's space.
823///
824/// `None` when the line does not begin with a well-formed stamp and one space:
825/// a line written before capture lines were stamped has no recorded time, and
826/// the caller must leave it undated rather than guess one.
827pub fn split_capture_stamp(line: &str) -> Option<(u64, &str)> {
828    let bytes = line.as_bytes();
829    if bytes.len() < CAPTURE_STAMP_PREFIX_LEN || bytes[CAPTURE_STAMP_LEN] != b' ' {
830        return None;
831    }
832    let stamp = &bytes[..CAPTURE_STAMP_LEN];
833    for (index, expected) in [
834        (4, b'-'),
835        (7, b'-'),
836        (10, b'T'),
837        (13, b':'),
838        (16, b':'),
839        (19, b'.'),
840        (23, b'Z'),
841    ] {
842        if stamp[index] != expected {
843            return None;
844        }
845    }
846    let number = |from: usize, to: usize| -> Option<u64> {
847        let digits = &stamp[from..to];
848        if !digits.iter().all(u8::is_ascii_digit) {
849            return None;
850        }
851        Some(
852            digits
853                .iter()
854                .fold(0u64, |total, digit| total * 10 + u64::from(digit - b'0')),
855        )
856    };
857    let year = number(0, 4)?;
858    let month = number(5, 7)?;
859    let day = number(8, 10)?;
860    let hour = number(11, 13)?;
861    let minute = number(14, 16)?;
862    let second = number(17, 19)?;
863    let millis = number(20, 23)?;
864    if year < 1970
865        || !(1..=12).contains(&month)
866        || day == 0
867        || day > days_in_month(year, month)
868        || hour > 23
869        || minute > 59
870        || second > 59
871    {
872        return None;
873    }
874    let days = days_from_civil(year, month, day);
875    let seconds = days * 86_400 + hour * 3600 + minute * 60 + second;
876    // Index 25 follows an ASCII space, so it is a char boundary.
877    Some((seconds * 1000 + millis, &line[CAPTURE_STAMP_PREFIX_LEN..]))
878}
879
880fn days_in_month(year: u64, month: u64) -> u64 {
881    match month {
882        2 if year.is_multiple_of(4) && (!year.is_multiple_of(100) || year.is_multiple_of(400)) => {
883            29
884        }
885        2 => 28,
886        4 | 6 | 9 | 11 => 30,
887        _ => 31,
888    }
889}
890
891// Days since 1970-01-01 to a proleptic Gregorian date and back, after Howard
892// Hinnant's `civil_from_days`/`days_from_civil`, restricted to dates from 1970
893// on so the arithmetic stays unsigned.
894fn civil_from_days(days: u64) -> (u64, u64, u64) {
895    let z = days + 719_468;
896    let era = z / 146_097;
897    let doe = z - era * 146_097;
898    let yoe = (doe - doe / 1_460 + doe / 36_524 - doe / 146_096) / 365;
899    let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
900    let mp = (5 * doy + 2) / 153;
901    let day = doy - (153 * mp + 2) / 5 + 1;
902    let month = if mp < 10 { mp + 3 } else { mp - 9 };
903    let year = yoe + era * 400 + u64::from(month <= 2);
904    (year, month, day)
905}
906
907fn days_from_civil(year: u64, month: u64, day: u64) -> u64 {
908    let year = if month <= 2 { year - 1 } else { year };
909    let era = year / 400;
910    let yoe = year - era * 400;
911    let shifted_month = if month > 2 { month - 3 } else { month + 9 };
912    let doy = (153 * shifted_month + 2) / 5 + day - 1;
913    let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
914    era * 146_097 + doe - 719_468
915}
916
917/// `entries` with every capture time removed, for tests that compare what a
918/// pump stored against expected text: the times are whatever the clock read,
919/// and the tests that check them do so on their own.
920#[cfg(test)]
921pub(crate) fn untimed(entries: Vec<TailEntry>) -> Vec<TailEntry> {
922    entries
923        .into_iter()
924        .map(|entry| match entry {
925            TailEntry::Line {
926                text, truncated, ..
927            } => TailEntry::Line {
928                text,
929                truncated,
930                at_ms: None,
931            },
932            TailEntry::ProcessStart => TailEntry::ProcessStart,
933        })
934        .collect()
935}
936
937fn lock_ring(ring: &Arc<Mutex<StderrRing>>) -> std::sync::MutexGuard<'_, StderrRing> {
938    ring.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
939}
940
941/// Cut `line` to at most `max_bytes`, reporting whether it was shortened.
942///
943/// Cuts on a char boundary: slicing a multi-byte sequence would produce invalid
944/// UTF-8, and a panic while capturing a crash message is the worst possible time
945/// to discover that.
946fn truncate_line(line: &str, max_bytes: usize) -> (String, bool) {
947    if line.len() <= max_bytes {
948        return (line.to_string(), false);
949    }
950    let mut end = max_bytes;
951    while end > 0 && !line.is_char_boundary(end) {
952        end -= 1;
953    }
954    (line[..end].to_string(), true)
955}
956
957#[cfg(test)]
958mod tests {
959    use std::{
960        io,
961        pin::Pin,
962        task::{Context, Poll},
963    };
964
965    use super::*;
966    use tokio::io::{AsyncRead, ReadBuf};
967
968    fn ring(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> StderrRing {
969        StderrRing::new(StderrTailConfig::new(max_lines, max_bytes, max_line_bytes))
970    }
971
972    fn lines(snapshot: &StderrTailSnapshot) -> Vec<String> {
973        snapshot
974            .entries
975            .iter()
976            .filter_map(|entry| match entry {
977                TailEntry::Line { text, .. } => Some(text.clone()),
978                TailEntry::ProcessStart => None,
979            })
980            .collect()
981    }
982
983    #[test]
984    fn a_fresh_ring_reports_not_captured_rather_than_empty() {
985        // The distinction this whole module exists for: "nobody was listening"
986        // must not render as "the module said nothing".
987        let ring = ring(10, 1024, 128);
988        let snapshot = ring.snapshot(None, None);
989        assert!(matches!(snapshot.capture, CaptureState::NotCaptured { .. }));
990        assert!(snapshot.entries.is_empty());
991    }
992
993    #[test]
994    fn a_captured_module_that_printed_nothing_is_distinguishable_from_an_uncaptured_one() {
995        let mut captured = ring(10, 1024, 128);
996        captured.mark_captured();
997        let uncaptured = ring(10, 1024, 128);
998
999        let captured = captured.snapshot(None, None);
1000        let uncaptured = uncaptured.snapshot(None, None);
1001
1002        // Both are empty. Only the capture state separates them, which is the
1003        // point -- an assertion on emptiness alone would pass either way.
1004        assert!(captured.entries.is_empty());
1005        assert!(uncaptured.entries.is_empty());
1006        assert_eq!(captured.capture, CaptureState::Captured);
1007        assert!(matches!(
1008            uncaptured.capture,
1009            CaptureState::NotCaptured { .. }
1010        ));
1011    }
1012
1013    #[test]
1014    fn the_line_cap_evicts_oldest_first_and_counts_what_it_dropped() {
1015        let mut ring = ring(3, 10_000, 128);
1016        ring.mark_captured();
1017        for i in 0..6 {
1018            ring.push_line(&format!("line{i}"));
1019        }
1020        let snapshot = ring.snapshot(None, None);
1021        assert_eq!(lines(&snapshot), vec!["line3", "line4", "line5"]);
1022        // Without this the tail silently becomes "the last lines that happened
1023        // to survive" and reads as complete.
1024        assert_eq!(snapshot.dropped_lines, 3);
1025    }
1026
1027    #[test]
1028    fn the_byte_cap_binds_before_the_line_cap_when_lines_are_large() {
1029        // 100 lines allowed, but only ~30 bytes of them.
1030        let mut ring = ring(100, 30, 128);
1031        ring.mark_captured();
1032        for i in 0..10 {
1033            ring.push_line(&format!("{i}--------")); // 9 bytes each
1034        }
1035        let snapshot = ring.snapshot(None, None);
1036        assert!(
1037            snapshot.entries.len() < 10,
1038            "byte cap did not bind: {} entries retained",
1039            snapshot.entries.len()
1040        );
1041        let retained: usize = lines(&snapshot).iter().map(String::len).sum();
1042        assert!(
1043            retained <= 30,
1044            "retained {retained} bytes over a 30 byte cap"
1045        );
1046        assert!(snapshot.dropped_lines > 0);
1047    }
1048
1049    #[test]
1050    fn one_enormous_line_is_truncated_rather_than_evicting_the_tail() {
1051        // The pathological-emitter case: without per-line truncation this single
1052        // line would evict every other line AND be unreadable itself.
1053        let mut ring = ring(10, 10_000, 64);
1054        ring.mark_captured();
1055        ring.push_line("context line that must survive");
1056        ring.push_line(&"x".repeat(40_000));
1057
1058        let snapshot = ring.snapshot(None, None);
1059        let kept = &snapshot.entries;
1060        assert!(matches!(
1061            &kept[0],
1062            TailEntry::Line { text, truncated: false, .. }
1063                if text == "context line that must survive"
1064        ));
1065        let TailEntry::Line {
1066            text, truncated, ..
1067        } = &kept[1]
1068        else {
1069            panic!("expected a truncated line");
1070        };
1071        assert_eq!(text, &"x".repeat(64));
1072        assert!(*truncated);
1073    }
1074
1075    #[test]
1076    fn truncation_is_visible_so_a_cut_line_is_not_mistaken_for_a_short_one() {
1077        let mut ring = ring(10, 10_000, 16);
1078        ring.mark_captured();
1079        ring.push_line("0123456789abcdefghij");
1080        ring.push_line("short");
1081
1082        let snapshot = ring.snapshot(None, None);
1083        let TailEntry::Line { truncated, .. } = &snapshot.entries[0] else {
1084            panic!("expected a line");
1085        };
1086        assert!(truncated);
1087        let TailEntry::Line { truncated, .. } = &snapshot.entries[1] else {
1088            panic!("expected a line");
1089        };
1090        assert!(!truncated, "a short line must not be reported as truncated");
1091    }
1092
1093    #[test]
1094    fn truncation_cuts_on_a_char_boundary_rather_than_splitting_utf8() {
1095        // A panic message with non-ASCII in it is not exotic, and slicing mid
1096        // sequence would panic while capturing a crash.
1097        let mut ring = ring(10, 10_000, 5);
1098        ring.mark_captured();
1099        ring.push_line("aa€€€€");
1100        let snapshot = ring.snapshot(None, None);
1101        let TailEntry::Line {
1102            text, truncated, ..
1103        } = &snapshot.entries[0]
1104        else {
1105            panic!("expected a line");
1106        };
1107        assert!(truncated);
1108        assert!(text.starts_with("aa"));
1109    }
1110
1111    #[test]
1112    fn a_restart_boundary_keeps_generations_distinguishable() {
1113        let mut ring = ring(10, 10_000, 128);
1114        ring.mark_captured();
1115        ring.push_line("before the crash");
1116        ring.push_process_start();
1117        ring.push_line("after the respawn");
1118
1119        let snapshot = ring.snapshot(None, None);
1120        assert_eq!(
1121            snapshot.entries,
1122            vec![
1123                TailEntry::Line {
1124                    text: "before the crash".to_string(),
1125                    truncated: false,
1126                    at_ms: None,
1127                },
1128                TailEntry::ProcessStart,
1129                TailEntry::Line {
1130                    text: "after the respawn".to_string(),
1131                    truncated: false,
1132                    at_ms: None,
1133                },
1134            ]
1135        );
1136    }
1137
1138    #[test]
1139    fn the_ring_survives_respawn_because_the_cause_is_written_before_the_exit() {
1140        // Clearing on restart would discard the lines at the exact moment they
1141        // become the thing being asked for.
1142        let mut ring = ring(10, 10_000, 128);
1143        ring.mark_captured();
1144        ring.push_line("Error: storage section missing");
1145        ring.push_process_start();
1146
1147        let snapshot = ring.snapshot(None, None);
1148        assert!(lines(&snapshot).contains(&"Error: storage section missing".to_string()));
1149    }
1150
1151    #[test]
1152    fn a_caller_limit_returns_the_newest_lines_not_the_oldest() {
1153        let mut ring = ring(100, 100_000, 128);
1154        ring.mark_captured();
1155        for i in 0..10 {
1156            ring.push_line(&format!("line{i}"));
1157        }
1158        let snapshot = ring.snapshot(Some(3), None);
1159        assert_eq!(lines(&snapshot), vec!["line7", "line8", "line9"]);
1160    }
1161
1162    #[test]
1163    fn a_caller_line_limit_keeps_the_boundary_before_the_selected_line() {
1164        let mut ring = ring(100, 100_000, 128);
1165        ring.mark_captured();
1166        ring.push_line("before restart");
1167        ring.push_process_start();
1168        ring.push_line("after restart");
1169
1170        let snapshot = ring.snapshot(Some(1), None);
1171        assert_eq!(
1172            snapshot.entries,
1173            vec![
1174                TailEntry::ProcessStart,
1175                TailEntry::Line {
1176                    text: "after restart".to_string(),
1177                    truncated: false,
1178                    at_ms: None,
1179                },
1180            ]
1181        );
1182    }
1183
1184    #[test]
1185    fn a_caller_line_limit_omits_a_trailing_boundary_after_the_selected_line() {
1186        let mut ring = ring(100, 100_000, 128);
1187        ring.mark_captured();
1188        ring.push_line("before restart");
1189        ring.push_process_start();
1190
1191        let snapshot = ring.snapshot(Some(1), None);
1192        assert_eq!(
1193            snapshot.entries,
1194            vec![TailEntry::Line {
1195                text: "before restart".to_string(),
1196                truncated: false,
1197                at_ms: None,
1198            }]
1199        );
1200    }
1201
1202    #[test]
1203    fn a_caller_limit_reports_what_it_withheld_rather_than_looking_complete() {
1204        let mut ring = ring(100, 100_000, 128);
1205        ring.mark_captured();
1206        for i in 0..10 {
1207            ring.push_line(&format!("line{i}"));
1208        }
1209        // Nothing was evicted; the narrowing is the caller's own. It still has to
1210        // be reported, or a 3-line request reads as a module that wrote 3 lines.
1211        assert_eq!(ring.snapshot(Some(3), None).dropped_lines, 7);
1212        assert_eq!(ring.snapshot(None, None).dropped_lines, 0);
1213    }
1214
1215    #[test]
1216    fn a_caller_limit_cannot_widen_the_rings_own_caps() {
1217        let mut ring = ring(2, 10_000, 128);
1218        ring.mark_captured();
1219        for i in 0..5 {
1220            ring.push_line(&format!("line{i}"));
1221        }
1222        let snapshot = ring.snapshot(Some(1000), Some(1_000_000));
1223        assert_eq!(lines(&snapshot).len(), 2);
1224    }
1225
1226    fn shared(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> Arc<Mutex<StderrRing>> {
1227        Arc::new(Mutex::new(ring(max_lines, max_bytes, max_line_bytes)))
1228    }
1229
1230    /// Records each forwarded write separately, so a test can tell one write of
1231    /// `b"abc\n"` from two writes of `b"abc"` and `b"\n"`.
1232    #[derive(Default)]
1233    struct RecordingSink {
1234        writes: Vec<Vec<u8>>,
1235    }
1236
1237    impl OutputSink for RecordingSink {
1238        fn write_line(&mut self, line: &[u8]) {
1239            self.writes.push(line.to_vec());
1240        }
1241    }
1242
1243    /// Yields predetermined chunks, one per read, so a test controls exactly
1244    /// where the byte stream is split.
1245    struct ChunkedReader {
1246        chunks: VecDeque<Vec<u8>>,
1247    }
1248
1249    impl AsyncRead for ChunkedReader {
1250        fn poll_read(
1251            mut self: Pin<&mut Self>,
1252            _cx: &mut Context<'_>,
1253            buf: &mut ReadBuf<'_>,
1254        ) -> Poll<io::Result<()>> {
1255            match self.chunks.pop_front() {
1256                None => Poll::Ready(Ok(())),
1257                Some(chunk) => {
1258                    buf.put_slice(&chunk);
1259                    Poll::Ready(Ok(()))
1260                }
1261            }
1262        }
1263    }
1264
1265    struct FailingReader {
1266        bytes: Vec<u8>,
1267        emitted: bool,
1268    }
1269
1270    impl AsyncRead for FailingReader {
1271        fn poll_read(
1272            mut self: Pin<&mut Self>,
1273            _cx: &mut Context<'_>,
1274            buf: &mut ReadBuf<'_>,
1275        ) -> Poll<io::Result<()>> {
1276            if self.emitted {
1277                return Poll::Ready(Err(io::Error::other("reader failed")));
1278            }
1279            self.emitted = true;
1280            buf.put_slice(&self.bytes);
1281            Poll::Ready(Ok(()))
1282        }
1283    }
1284
1285    #[tokio::test]
1286    async fn the_pump_splits_on_newlines_and_keeps_a_trailing_fragment() {
1287        let ring = shared(10, 10_000, 128);
1288        // No trailing newline on the last line: a crashing process routinely dies
1289        // mid-line, and that fragment is often the message worth reading.
1290        let source = std::io::Cursor::new(b"one\ntwo\nthree".to_vec());
1291        let mut sink = RecordingSink::default();
1292        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1293
1294        let snapshot = lock_ring(&ring).snapshot(None, None);
1295        assert_eq!(lines(&snapshot), vec!["one", "two", "three"]);
1296        assert_eq!(snapshot.capture, CaptureState::Captured);
1297        assert_eq!(
1298            sink.writes,
1299            vec![b"one\n".to_vec(), b"two\n".to_vec(), b"three".to_vec()]
1300        );
1301    }
1302
1303    #[tokio::test]
1304    async fn a_read_failure_keeps_prior_lines_and_marks_the_capture_incomplete() {
1305        let ring = shared(10, 10_000, 128);
1306        let source = FailingReader {
1307            bytes: b"crash cause\n".to_vec(),
1308            emitted: false,
1309        };
1310        let mut sink = RecordingSink::default();
1311        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1312
1313        let snapshot = lock_ring(&ring).snapshot(None, None);
1314        assert_eq!(lines(&snapshot), vec!["crash cause"]);
1315        assert!(matches!(
1316            snapshot.capture,
1317            CaptureState::Incomplete { ref reason } if reason.contains("reader failed")
1318        ));
1319        assert_eq!(sink.writes, vec![b"crash cause\n".to_vec()]);
1320    }
1321
1322    #[tokio::test]
1323    async fn every_captured_line_is_also_forwarded() {
1324        // Forwarding is not optional. The daemon log is overwhelmingly module
1325        // output; a tap that captured without forwarding would leave it nearly
1326        // empty and every existing reader would report clean on nothing.
1327        let ring = shared(10, 10_000, 128);
1328        let source = std::io::Cursor::new(b"alpha\nbeta\n".to_vec());
1329        let mut sink = RecordingSink::default();
1330        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1331
1332        assert_eq!(sink.writes, vec![b"alpha\n".to_vec(), b"beta\n".to_vec()]);
1333    }
1334
1335    #[tokio::test]
1336    async fn each_forwarded_line_is_exactly_one_write() {
1337        // Inheriting the fd gave line atomicity for free. Reading a pipe and
1338        // re-emitting can split a line that used to be atomic, so the framing
1339        // must be one syscall per complete line -- asserted as one write per
1340        // line, not merely as correct bytes.
1341        let ring = shared(10, 10_000, 128);
1342        let source = std::io::Cursor::new(b"first\nsecond\nthird\n".to_vec());
1343        let mut sink = RecordingSink::default();
1344        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1345
1346        assert_eq!(sink.writes.len(), 3);
1347        for write in &sink.writes {
1348            assert_eq!(
1349                write.iter().filter(|byte| **byte == b'\n').count(),
1350                1,
1351                "a write carried something other than exactly one complete line"
1352            );
1353            assert_eq!(*write.last().unwrap(), b'\n');
1354        }
1355    }
1356
1357    #[test]
1358    fn the_first_process_start_is_not_recorded_because_it_divides_nothing() {
1359        // Otherwise a module that printed nothing renders as a lone boundary
1360        // marker, and every caller has to decide whether that counts as silence.
1361        let mut ring = ring(10, 10_000, 128);
1362        ring.push_process_start();
1363        assert!(ring.snapshot(None, None).entries.is_empty());
1364
1365        ring.push_line("first process said this");
1366        ring.push_process_start();
1367        assert!(
1368            matches!(ring.entries.back(), Some(Slot::ProcessStart { .. })),
1369            "a boundary with output before it must be recorded"
1370        );
1371        ring.push_line("second process said this");
1372        assert_eq!(
1373            ring.snapshot(None, None).entries,
1374            vec![
1375                TailEntry::Line {
1376                    text: "first process said this".to_string(),
1377                    truncated: false,
1378                    at_ms: None,
1379                },
1380                TailEntry::ProcessStart,
1381                TailEntry::Line {
1382                    text: "second process said this".to_string(),
1383                    truncated: false,
1384                    at_ms: None,
1385                },
1386            ],
1387            "only the boundary with output before it may be shown"
1388        );
1389    }
1390
1391    #[test]
1392    fn a_process_start_is_recorded_when_only_dropped_lines_precede_it() {
1393        // The ring can be non-empty in the sense that matters -- lines were
1394        // written and evicted -- while `entries` is empty. Suppressing the
1395        // boundary there would attribute surviving output to the wrong process.
1396        let mut ring = ring(1, 10_000, 128);
1397        ring.push_line("evicted");
1398        ring.push_line("also evicted");
1399        // Emptying `entries` by hand must also zero the running totals kept
1400        // beside it, or the ring holds counts for lines it no longer has and
1401        // any later eviction decision is made against the stale numbers.
1402        ring.entries.clear();
1403        ring.lines = 0;
1404        ring.bytes = 0;
1405        ring.push_process_start();
1406        ring.push_line("survivor");
1407        assert_eq!(
1408            ring.snapshot(None, None).entries,
1409            vec![
1410                TailEntry::ProcessStart,
1411                TailEntry::Line {
1412                    text: "survivor".to_string(),
1413                    truncated: false,
1414                    at_ms: None,
1415                },
1416            ]
1417        );
1418    }
1419
1420    fn line(text: &str) -> TailEntry {
1421        TailEntry::Line {
1422            text: text.to_string(),
1423            truncated: false,
1424            at_ms: None,
1425        }
1426    }
1427
1428    #[test]
1429    fn a_late_line_from_a_retired_process_lands_in_that_processs_section() {
1430        // The reader of a process that already exited may deliver its last
1431        // lines after the next process started. Appending them would put the
1432        // crash's own explanation under its successor's boundary.
1433        let mut ring = ring(10, 10_000, 128);
1434        ring.mark_captured();
1435        let old = ring.begin_process();
1436        ring.push_line_from(old, "old: booting");
1437        ring.retire_pump(old);
1438        let new = ring.begin_process();
1439        ring.push_line_from(new, "new: booting");
1440        ring.push_line_from(old, "old: config error");
1441
1442        assert_eq!(
1443            ring.snapshot(None, None).entries,
1444            vec![
1445                line("old: booting"),
1446                line("old: config error"),
1447                TailEntry::ProcessStart,
1448                line("new: booting"),
1449            ]
1450        );
1451    }
1452
1453    #[test]
1454    fn a_line_from_a_process_that_was_not_retired_is_appended_as_it_arrives() {
1455        // A swap runs the incumbent alongside its candidate; until the
1456        // supervisor retires it, the incumbent is live and its lines are news.
1457        let mut ring = ring(10, 10_000, 128);
1458        ring.mark_captured();
1459        let incumbent = ring.begin_process();
1460        ring.push_line_from(incumbent, "incumbent: before");
1461        let candidate = ring.begin_process();
1462        ring.push_line_from(candidate, "candidate: booting");
1463        ring.push_line_from(incumbent, "incumbent: still serving");
1464
1465        assert_eq!(
1466            ring.snapshot(None, None).entries,
1467            vec![
1468                line("incumbent: before"),
1469                TailEntry::ProcessStart,
1470                line("candidate: booting"),
1471                line("incumbent: still serving"),
1472            ]
1473        );
1474    }
1475
1476    #[test]
1477    fn a_late_line_keeps_its_section_when_the_process_had_printed_nothing_before() {
1478        // A process whose reader had delivered nothing when its successor
1479        // started has a boundary with nothing after it. Dropping that boundary
1480        // as redundant would file the late line under the process before.
1481        let mut ring = ring(10, 10_000, 128);
1482        ring.mark_captured();
1483        let first = ring.begin_process();
1484        ring.push_line_from(first, "first: done");
1485        ring.finish_pump(first);
1486        let old = ring.begin_process();
1487        ring.retire_pump(old);
1488        let new = ring.begin_process();
1489        ring.push_line_from(new, "new: booting");
1490        ring.push_line_from(old, "old: config error");
1491
1492        assert_eq!(
1493            ring.snapshot(None, None).entries,
1494            vec![
1495                line("first: done"),
1496                TailEntry::ProcessStart,
1497                line("old: config error"),
1498                TailEntry::ProcessStart,
1499                line("new: booting"),
1500            ]
1501        );
1502    }
1503
1504    #[test]
1505    fn a_late_line_whose_section_was_evicted_counts_as_dropped() {
1506        let mut ring = ring(2, 10_000, 128);
1507        ring.mark_captured();
1508        let old = ring.begin_process();
1509        ring.push_line_from(old, "old");
1510        ring.retire_pump(old);
1511        let new = ring.begin_process();
1512        for text in ["new 1", "new 2", "new 3"] {
1513            ring.push_line_from(new, text);
1514        }
1515        // The old section and the boundary after it are gone; the late line
1516        // belongs in front of everything retained.
1517        ring.push_line_from(old, "old, late");
1518
1519        let snapshot = ring.snapshot(None, None);
1520        assert_eq!(snapshot.entries, vec![line("new 2"), line("new 3")]);
1521        assert_eq!(snapshot.dropped_lines, 3);
1522    }
1523
1524    #[test]
1525    fn a_late_reader_reads_incomplete_until_its_pipe_reaches_eof() {
1526        let mut ring = ring(10, 10_000, 128);
1527        ring.mark_captured();
1528        let old = ring.begin_process();
1529        ring.retire_pump(old);
1530        ring.mark_pump_late(old, "still open");
1531        ring.begin_process();
1532        assert_eq!(
1533            ring.snapshot(None, None).capture,
1534            CaptureState::Incomplete {
1535                reason: "still open".to_string()
1536            }
1537        );
1538
1539        ring.finish_pump(old);
1540        assert_eq!(ring.snapshot(None, None).capture, CaptureState::Captured);
1541    }
1542
1543    #[test]
1544    fn silent_restarts_do_not_grow_the_ring() {
1545        let mut ring = ring(10, 10_000, 128);
1546        ring.mark_captured();
1547        ring.push_line("once");
1548        for _ in 0..100 {
1549            let generation = ring.begin_process();
1550            ring.finish_pump(generation);
1551        }
1552        assert_eq!(ring.entries.len(), 2);
1553    }
1554
1555    #[tokio::test]
1556    async fn the_pump_marks_captured_even_when_the_module_writes_nothing() {
1557        // Clean EOF with no output is a module that was quiet, not one nobody
1558        // listened to -- and the two must not render alike.
1559        let ring = shared(10, 10_000, 128);
1560        let source = std::io::Cursor::new(Vec::new());
1561        let mut sink = RecordingSink::default();
1562        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1563
1564        let snapshot = lock_ring(&ring).snapshot(None, None);
1565        assert!(snapshot.entries.is_empty());
1566        assert_eq!(snapshot.capture, CaptureState::Captured);
1567        assert!(sink.writes.is_empty());
1568    }
1569
1570    #[tokio::test]
1571    async fn a_line_with_no_newline_cannot_grow_the_reader_without_bound() {
1572        // A module fault must not become a daemon fault: without the pending
1573        // ceiling this buffer grows to whatever the module writes.
1574        let ring = shared(10, 10_000_000, 4 * 1024 * 1024);
1575        let source = std::io::Cursor::new(vec![b'x'; MAX_PENDING_LINE_BYTES + 4096]);
1576        let mut sink = RecordingSink::default();
1577        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1578
1579        let snapshot = lock_ring(&ring).snapshot(None, None);
1580        assert_eq!(
1581            lines(&snapshot).len(),
1582            2,
1583            "expected a forced flush at the ceiling plus the remainder"
1584        );
1585        assert_eq!(
1586            sink.writes,
1587            vec![vec![b'x'; MAX_PENDING_LINE_BYTES], vec![b'x'; 4096],],
1588            "forced flushes and EOF fragments must not invent delimiters"
1589        );
1590    }
1591
1592    #[tokio::test]
1593    async fn boundaries_truncation_and_framing_do_not_depend_on_chunk_splits() {
1594        // The same stream split at hostile boundaries -- mid-line, between a CR
1595        // and its LF, and a line sitting exactly on the per-line cap -- must
1596        // produce the same ring entries and forwarded bytes as any other split.
1597        let ring = shared(100, 100_000, 8);
1598        let source = ChunkedReader {
1599            chunks: vec![
1600                b"fir".to_vec(),
1601                b"st\nsec".to_vec(),
1602                b"ond\ncarry\r".to_vec(),
1603                b"\nover\n".to_vec(),
1604                b"12345678\n".to_vec(),
1605                b"1234567".to_vec(),
1606                b"89\n".to_vec(),
1607                b"tail".to_vec(),
1608            ]
1609            .into_iter()
1610            .collect(),
1611        };
1612        let mut sink = RecordingSink::default();
1613        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1614
1615        let snapshot = lock_ring(&ring).snapshot(None, None);
1616        assert_eq!(snapshot.capture, CaptureState::Captured);
1617        assert_eq!(
1618            untimed(snapshot.entries),
1619            vec![
1620                TailEntry::Line {
1621                    text: "first".to_string(),
1622                    truncated: false,
1623                    at_ms: None,
1624                },
1625                TailEntry::Line {
1626                    text: "second".to_string(),
1627                    truncated: false,
1628                    at_ms: None,
1629                },
1630                // The pump delimits on '\n' alone; a CR belongs to the line body.
1631                TailEntry::Line {
1632                    text: "carry\r".to_string(),
1633                    truncated: false,
1634                    at_ms: None,
1635                },
1636                TailEntry::Line {
1637                    text: "over".to_string(),
1638                    truncated: false,
1639                    at_ms: None,
1640                },
1641                // Exactly at the per-line cap: kept whole.
1642                TailEntry::Line {
1643                    text: "12345678".to_string(),
1644                    truncated: false,
1645                    at_ms: None,
1646                },
1647                // One byte past the cap: cut, and marked as cut.
1648                TailEntry::Line {
1649                    text: "12345678".to_string(),
1650                    truncated: true,
1651                    at_ms: None,
1652                },
1653                TailEntry::Line {
1654                    text: "tail".to_string(),
1655                    truncated: false,
1656                    at_ms: None,
1657                },
1658            ]
1659        );
1660        assert_eq!(
1661            sink.writes,
1662            vec![
1663                b"first\n".to_vec(),
1664                b"second\n".to_vec(),
1665                b"carry\r\n".to_vec(),
1666                b"over\n".to_vec(),
1667                b"12345678\n".to_vec(),
1668                b"123456789\n".to_vec(),
1669                b"tail".to_vec(),
1670            ]
1671        );
1672    }
1673
1674    #[tokio::test]
1675    async fn a_line_with_no_newline_is_not_rescanned_from_byte_zero_on_every_chunk() {
1676        // One 1 MiB line arrives in 8192-byte reads. Searching the whole
1677        // pending buffer for a newline on every chunk scans each byte once per
1678        // chunk that arrived after it -- about 64 MiB examined per MiB of
1679        // output. Searching only the bytes that arrived since the last search
1680        // scans each byte once.
1681        let input = vec![b'x'; MAX_PENDING_LINE_BYTES + 4096];
1682        let ring = shared(10, 10_000_000, 4 * 1024 * 1024);
1683        let source = std::io::Cursor::new(input.clone());
1684        let mut sink = RecordingSink::default();
1685
1686        take_scanned_bytes();
1687        pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1688        let scanned = take_scanned_bytes();
1689
1690        assert!(
1691            scanned <= 2 * input.len(),
1692            "newline searches examined {scanned} bytes for {} bytes of input; \
1693             each chunk must search only newly arrived bytes",
1694            input.len()
1695        );
1696    }
1697
1698    fn now_ms() -> u64 {
1699        unix_ms(SystemTime::now())
1700    }
1701
1702    #[tokio::test]
1703    async fn a_captured_line_carries_the_time_the_reader_framed_it() {
1704        // The ring is read after the fact (`ck module stderr`), and without a
1705        // time a reader cannot tell a crash's last words from a line written
1706        // hours before it.
1707        let ring = shared(10, 10_000, 128);
1708        let before = now_ms();
1709        let source = std::io::Cursor::new(b"first\nsecond".to_vec());
1710        pump_stderr_into(source, Arc::clone(&ring), &mut RecordingSink::default()).await;
1711        let after = now_ms();
1712
1713        let entries = lock_ring(&ring).snapshot(None, None).entries;
1714        assert_eq!(entries.len(), 2);
1715        for entry in entries {
1716            let TailEntry::Line { text, at_ms, .. } = entry else {
1717                panic!("expected only lines, got {entry:?}");
1718            };
1719            let at_ms = at_ms.unwrap_or_else(|| panic!("line {text:?} has no capture time"));
1720            assert!(
1721                (before..=after).contains(&at_ms),
1722                "line {text:?} stamped {at_ms}, outside the pump's run {before}..={after}"
1723            );
1724        }
1725    }
1726
1727    /// Reads a capture file written through the production sink.
1728    async fn capture_through_file_sink(chunks: Vec<Vec<u8>>) -> (String, u64, u64) {
1729        let temp = subc_test_support::TestTempDir::new("stderr-capture-stamp");
1730        let path = temp.path().join("stamped.stderr.log");
1731        let sink = ChildOutputSink::open(&path, cortexkit_log::Retention::default()).unwrap();
1732        let ring = shared(10, 10_000, 128);
1733        let generation = lock_ring(&ring).begin_process();
1734        let before = now_ms();
1735        pump_stderr_to(
1736            ChunkedReader {
1737                chunks: chunks.into_iter().collect(),
1738            },
1739            ring,
1740            generation,
1741            sink,
1742        )
1743        .await;
1744        let after = now_ms();
1745        (std::fs::read_to_string(&path).unwrap(), before, after)
1746    }
1747
1748    /// Checks the prefix's shape by position, independently of
1749    /// [`split_capture_stamp`], so a parser that accepted a malformed stamp
1750    /// could not vouch for the writer that produced it.
1751    fn assert_stamp_shape(line: &str) {
1752        let bytes = line.as_bytes();
1753        assert!(bytes.len() > 25, "line too short for a stamp: {line:?}");
1754        for (index, byte) in bytes[..25].iter().enumerate() {
1755            let expected_separator = match index {
1756                4 | 7 => Some(b'-'),
1757                10 => Some(b'T'),
1758                13 | 16 => Some(b':'),
1759                19 => Some(b'.'),
1760                23 => Some(b'Z'),
1761                24 => Some(b' '),
1762                _ => None,
1763            };
1764            match expected_separator {
1765                Some(separator) => assert_eq!(*byte, separator, "byte {index} of {line:?}"),
1766                None => assert!(byte.is_ascii_digit(), "byte {index} of {line:?}"),
1767            }
1768        }
1769    }
1770
1771    #[tokio::test]
1772    async fn the_capture_file_stamps_each_line_and_keeps_the_module_bytes_verbatim() {
1773        // The second line begins with a stamp of its own (a module that logs
1774        // in the fleet format to stderr). Its bytes must survive untouched:
1775        // the capture stamp goes in front, nothing is parsed or replaced.
1776        let module_lines = [
1777            "plain line",
1778            "2020-01-01T00:00:00.000Z INFO  mymod: own stamp",
1779        ];
1780        let input = format!("{}\n{}\n", module_lines[0], module_lines[1]);
1781        let (contents, before, after) = capture_through_file_sink(vec![input.into_bytes()]).await;
1782
1783        assert!(contents.ends_with('\n'));
1784        let lines: Vec<&str> = contents.lines().collect();
1785        assert_eq!(lines.len(), 2, "capture file: {contents:?}");
1786        for (line, module_line) in lines.iter().zip(module_lines) {
1787            assert_stamp_shape(line);
1788            assert_eq!(&line[25..], module_line, "module bytes changed");
1789            let (at_ms, rest) = split_capture_stamp(line).unwrap();
1790            assert_eq!(rest, module_line);
1791            // Millisecond stamps round the pump's own bounds down.
1792            assert!(
1793                (before..=after).contains(&at_ms),
1794                "stamped {at_ms}, outside the pump's run {before}..={after}"
1795            );
1796        }
1797    }
1798
1799    #[tokio::test]
1800    async fn a_line_split_across_several_writes_is_stamped_once() {
1801        // A module that writes one line in pieces, and a pipe that delivers it
1802        // in pieces, must still produce one stamp at the front of the line and
1803        // none in its middle.
1804        let (contents, _, _) = capture_through_file_sink(vec![
1805            b"par".to_vec(),
1806            b"tial li".to_vec(),
1807            b"ne\nwhole\n".to_vec(),
1808        ])
1809        .await;
1810
1811        let lines: Vec<&str> = contents.lines().collect();
1812        assert_eq!(lines.len(), 2, "capture file: {contents:?}");
1813        for (line, module_line) in lines.iter().zip(["partial line", "whole"]) {
1814            assert_stamp_shape(line);
1815            assert_eq!(&line[25..], module_line, "capture file: {contents:?}");
1816        }
1817    }
1818
1819    #[tokio::test]
1820    async fn a_line_flushed_at_the_ceiling_is_stamped_only_at_the_start_of_each_file_line() {
1821        // Past the reassembly ceiling the pending bytes go out without their
1822        // newline. The file sink ends every write with one, so the rest of
1823        // the line starts a new file line, and that is where its stamp goes;
1824        // a stamp anywhere else would sit in the middle of module bytes.
1825        let mut long = vec![b'x'; MAX_PENDING_LINE_BYTES + 100];
1826        long.push(b'\n');
1827        // The reader's buffer holds 8 KiB, so feed it reads no larger than that.
1828        let chunks = long.chunks(8192).map(<[u8]>::to_vec).collect();
1829        let (contents, _, _) = capture_through_file_sink(chunks).await;
1830
1831        let lines: Vec<&str> = contents.lines().collect();
1832        assert_eq!(lines.len(), 2, "expected the flushed piece and the rest");
1833        for line in &lines {
1834            assert_stamp_shape(line);
1835            assert!(
1836                line[25..].bytes().all(|byte| byte == b'x'),
1837                "a stamp landed inside module bytes"
1838            );
1839        }
1840        let module_bytes: usize = lines.iter().map(|line| line.len() - 25).sum();
1841        assert_eq!(module_bytes, MAX_PENDING_LINE_BYTES + 100);
1842    }
1843
1844    #[test]
1845    fn the_capture_stamp_is_the_daemon_log_timestamp_form() {
1846        for (at_ms, text) in [
1847            (0, "1970-01-01T00:00:00.000Z"),
1848            (951_868_799_999, "2000-02-29T23:59:59.999Z"),
1849            (1_789_801_440_685, "2026-09-19T07:04:00.685Z"),
1850            (4_107_542_400_001, "2100-03-01T00:00:00.001Z"),
1851        ] {
1852            let stamp = format_capture_stamp(at_ms);
1853            assert_eq!(stamp, text);
1854            assert_eq!(stamp.len(), CAPTURE_STAMP_LEN);
1855            // The daemon's own log parser reads the same instant from it, so
1856            // one parser serves both files.
1857            let daemon_line = format!("{stamp} INFO  subc: probe");
1858            let parsed = cortexkit_log::parse_line(&daemon_line).unwrap();
1859            assert_eq!(
1860                parsed.timestamp,
1861                UNIX_EPOCH + std::time::Duration::from_millis(at_ms)
1862            );
1863            assert_eq!(
1864                split_capture_stamp(&format!("{stamp} body")),
1865                Some((at_ms, "body"))
1866            );
1867        }
1868    }
1869
1870    #[test]
1871    fn a_line_without_a_well_formed_stamp_has_no_capture_time() {
1872        for line in [
1873            "",
1874            "plain module output",
1875            "2026-09-19T07:04:00.685Z",
1876            "2026-09-19T07:04:00.685Zbody",
1877            "2026-09-19T07:04:00.685z body",
1878            "2026-09-19 07:04:00.685Z body",
1879            "2026-09-19T07:04:00Z body",
1880            "2026-02-30T07:04:00.685Z body",
1881            "2026-13-19T07:04:00.685Z body",
1882            "2026-09-19T24:04:00.685Z body",
1883            "2026-09-19T07:04:00.6a5Z body",
1884            "1969-12-31T23:59:59.999Z body",
1885        ] {
1886            assert_eq!(split_capture_stamp(line), None, "{line:?}");
1887        }
1888    }
1889
1890    #[test]
1891    fn a_byte_limit_smaller_than_one_line_still_returns_that_line() {
1892        // Returning nothing would be indistinguishable from a quiet module, which
1893        // is the failure this module exists to prevent.
1894        let mut ring = ring(10, 10_000, 128);
1895        ring.mark_captured();
1896        ring.push_line("a line considerably longer than the request limit");
1897        let snapshot = ring.snapshot(None, Some(4));
1898        assert_eq!(snapshot.entries.len(), 1);
1899    }
1900
1901    #[test]
1902    fn an_incoherent_config_clamps_the_line_cap_and_keeps_its_restart_boundary() {
1903        let config = StderrTailConfig::new(2, 10, 100);
1904        assert_eq!(config.max_line_bytes, config.max_bytes);
1905        let mut ring = StderrRing::new(config);
1906        ring.mark_captured();
1907        ring.push_line("old");
1908        ring.push_process_start();
1909        ring.push_line("new process line longer than the ring byte cap");
1910
1911        assert_eq!(
1912            ring.snapshot(None, None).entries,
1913            vec![
1914                TailEntry::ProcessStart,
1915                TailEntry::Line {
1916                    text: "new proces".to_string(),
1917                    truncated: true,
1918                    at_ms: None,
1919                },
1920            ]
1921        );
1922    }
1923}