1use 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
32pub const DEFAULT_MAX_LINE_BYTES: usize = 2048;
39
40pub const DEFAULT_MAX_LINES: usize = 200;
42
43pub const DEFAULT_MAX_BYTES: usize = 64 * 1024;
49
50#[derive(Debug, Clone, PartialEq, Eq)]
58pub enum CaptureState {
59 Captured,
61 Incomplete { reason: String },
66 NotCaptured { reason: String },
68}
69
70#[derive(Debug, Clone, PartialEq, Eq)]
77pub enum TailEntry {
78 Line {
79 text: String,
80 truncated: bool,
82 at_ms: Option<u64>,
86 },
87 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#[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#[derive(Debug, Clone, PartialEq, Eq)]
128enum PumpPhase {
129 Attached,
131 Retired,
134 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 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#[derive(Debug, Clone, PartialEq, Eq)]
170pub struct StderrTailSnapshot {
171 pub capture: CaptureState,
172 pub entries: Vec<TailEntry>,
173 pub dropped_lines: u64,
180}
181
182impl StderrTailSnapshot {
183 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#[derive(Debug)]
208pub struct StderrRing {
209 config: StderrTailConfig,
210 entries: VecDeque<Slot>,
211 lines: usize,
215 bytes: usize,
216 dropped_lines: u64,
217 capture: CaptureState,
218 incomplete: BTreeMap<u64, String>,
221 generation: u64,
223 pumps: BTreeMap<u64, PumpPhase>,
225 evicted_through: u64,
229}
230
231impl StderrRing {
232 pub fn new(config: StderrTailConfig) -> Self {
233 Self {
234 config,
235 entries: VecDeque::new(),
236 lines: 0,
237 bytes: 0,
238 dropped_lines: 0,
239 capture: CaptureState::NotCaptured {
242 reason: "stderr reader has not started".to_string(),
243 },
244 generation: 0,
245 pumps: BTreeMap::new(),
246 incomplete: BTreeMap::new(),
247 evicted_through: 0,
248 }
249 }
250
251 pub fn generation(&self) -> u64 {
253 self.generation
254 }
255
256 pub fn mark_captured(&mut self) {
257 if matches!(self.capture, CaptureState::NotCaptured { .. }) {
258 self.capture = CaptureState::Captured;
259 }
260 }
261
262 pub fn mark_incomplete(&mut self, reason: impl Into<String>) {
263 self.mark_captured();
264 self.mark_incomplete_from(self.generation, reason);
265 }
266
267 pub(crate) fn mark_incomplete_from(&mut self, generation: u64, reason: impl Into<String>) {
268 let oldest_retained = match self.entries.front() {
269 Some(Slot::ProcessStart { generation }) => *generation,
270 Some(Slot::Line { .. }) => self.evicted_through,
271 None => self.generation,
272 };
273 if generation >= oldest_retained {
274 self.incomplete.insert(generation, reason.into());
275 }
276 }
277
278 pub fn mark_not_captured(&mut self, reason: impl Into<String>) {
279 self.capture = CaptureState::NotCaptured {
280 reason: reason.into(),
281 };
282 }
283
284 pub fn push_process_start(&mut self) -> u64 {
296 self.generation += 1;
297 let generation = self.generation;
298 if let Some(Slot::ProcessStart {
304 generation: previous,
305 }) = self.entries.back()
306 {
307 if self.pumps.range(*previous..).next().is_none() {
308 self.entries.pop_back();
309 }
310 }
311 self.push_entry(Slot::ProcessStart { generation });
312 generation
313 }
314
315 pub(crate) fn begin_process(&mut self) -> u64 {
318 let generation = self.push_process_start();
319 self.pumps.insert(generation, PumpPhase::Attached);
320 generation
321 }
322
323 pub(crate) fn retire_pump(&mut self, generation: u64) {
326 if let Some(phase @ PumpPhase::Attached) = self.pumps.get_mut(&generation) {
327 *phase = PumpPhase::Retired;
328 }
329 }
330
331 pub(crate) fn mark_pump_late(&mut self, generation: u64, reason: impl Into<String>) {
336 if let Some(phase) = self.pumps.get_mut(&generation) {
337 *phase = PumpPhase::Late {
338 reason: reason.into(),
339 };
340 }
341 }
342
343 pub(crate) fn finish_pump(&mut self, generation: u64) {
346 self.pumps.remove(&generation);
347 }
348
349 pub fn push_line(&mut self, line: &str) {
356 self.push_line_from(self.generation, line);
357 }
358
359 pub(crate) fn push_line_from(&mut self, generation: u64, line: &str) {
361 self.admit_line(generation, line, None);
362 }
363
364 pub(crate) fn push_line_from_at(&mut self, generation: u64, line: &str, at_ms: u64) {
367 self.admit_line(generation, line, Some(at_ms));
368 }
369
370 fn admit_line(&mut self, generation: u64, line: &str, at_ms: Option<u64>) {
375 let (text, truncated) = truncate_line(line, self.config.max_line_bytes);
376 let slot = Slot::Line {
377 text,
378 truncated,
379 at_ms,
380 };
381 let retired = matches!(
382 self.pumps.get(&generation),
383 Some(PumpPhase::Retired | PumpPhase::Late { .. })
384 );
385 if !retired || generation >= self.generation {
386 self.push_entry(slot);
387 return;
388 }
389 if self.evicted_through > generation {
390 self.dropped_lines += 1;
393 return;
394 }
395 let index = self.entries.iter().position(
396 |slot| matches!(slot, Slot::ProcessStart { generation: start } if *start > generation),
397 );
398 match index {
399 Some(index) => self.insert_entry(index, slot),
400 None => self.push_entry(slot),
401 }
402 }
403
404 fn push_entry(&mut self, entry: Slot) {
405 self.insert_entry(self.entries.len(), entry);
406 }
407
408 fn insert_entry(&mut self, index: usize, entry: Slot) {
409 self.bytes += entry.cost();
410 if matches!(entry, Slot::Line { .. }) {
411 self.lines += 1;
412 }
413 self.entries.insert(index, entry);
414 self.evict_to_fit();
415 }
416
417 fn evict_to_fit(&mut self) {
418 while self.lines > self.config.max_lines
419 || (self.bytes > self.config.max_bytes && self.entries.len() > 1)
420 {
421 let Some(evicted) = self.entries.pop_front() else {
422 break;
423 };
424 self.bytes -= evicted.cost();
425 match evicted {
426 Slot::Line { .. } => {
427 self.lines -= 1;
428 self.dropped_lines += 1;
429 }
430 Slot::ProcessStart { generation } => {
431 self.evicted_through = self.evicted_through.max(generation);
432 }
433 }
434 }
435 let oldest_retained = match self.entries.front() {
436 Some(Slot::ProcessStart { generation }) => *generation,
437 Some(Slot::Line { .. }) => self.evicted_through,
438 None => self.generation,
439 };
440 self.incomplete
441 .retain(|generation, _| *generation >= oldest_retained);
442 }
443
444 pub fn snapshot(
448 &self,
449 max_lines: Option<usize>,
450 max_bytes: Option<usize>,
451 ) -> StderrTailSnapshot {
452 let line_limit = max_lines.unwrap_or(self.config.max_lines);
453 let byte_limit = max_bytes.unwrap_or(self.config.max_bytes);
454
455 let mut visible: Vec<TailEntry> = Vec::with_capacity(self.entries.len());
459 let mut output_before = self.dropped_lines > 0;
460 for slot in &self.entries {
461 match slot {
462 Slot::Line {
463 text,
464 truncated,
465 at_ms,
466 } => {
467 visible.push(TailEntry::Line {
468 text: text.clone(),
469 truncated: *truncated,
470 at_ms: *at_ms,
471 });
472 output_before = true;
473 }
474 Slot::ProcessStart { .. } => {
475 if output_before && !matches!(visible.last(), Some(TailEntry::ProcessStart)) {
476 visible.push(TailEntry::ProcessStart);
477 }
478 }
479 }
480 }
481
482 let mut taken: Vec<TailEntry> = Vec::new();
483 let mut bytes = 0usize;
484 let mut lines = 0usize;
485 for entry in visible.iter().rev() {
488 match entry {
489 TailEntry::Line { .. } => {
490 if lines >= line_limit {
491 break;
492 }
493 let cost = entry.cost();
494 if lines > 0 && bytes + cost > byte_limit {
495 break;
496 }
497 bytes += cost;
498 lines += 1;
499 taken.push(entry.clone());
500 }
501 TailEntry::ProcessStart if lines > 0 => taken.push(entry.clone()),
502 TailEntry::ProcessStart => {}
503 }
504 }
505 taken.reverse();
506
507 let withheld = self.lines.saturating_sub(lines);
508
509 let late = self.pumps.values().find_map(|phase| match phase {
515 PumpPhase::Late { reason } => Some(reason),
516 _ => None,
517 });
518 let incomplete = self.incomplete.values().next().or(late);
519 let capture = match (&self.capture, incomplete) {
520 (CaptureState::Captured, Some(reason)) => CaptureState::Incomplete {
521 reason: reason.clone(),
522 },
523 (capture, _) => capture.clone(),
524 };
525
526 StderrTailSnapshot {
527 capture,
528 entries: taken,
529 dropped_lines: self.dropped_lines + withheld as u64,
534 }
535 }
536}
537
538const MAX_PENDING_LINE_BYTES: usize = 1024 * 1024;
547
548pub async fn pump_stderr<R>(source: R, ring: Arc<Mutex<StderrRing>>)
551where
552 R: AsyncReadExt + Unpin,
553{
554 pump_stderr_into(source, ring, &mut StderrSink).await
555}
556
557#[derive(Clone)]
564pub(crate) enum ChildOutputSink {
565 File {
566 sink: Arc<Mutex<cortexkit_log::LineSink>>,
567 path: Arc<PathBuf>,
568 failure_reported: Arc<AtomicBool>,
569 },
570 Stderr,
571}
572
573impl ChildOutputSink {
574 pub(crate) fn open(path: &Path, retention: cortexkit_log::Retention) -> io::Result<Self> {
575 Ok(Self::File {
576 sink: Arc::new(Mutex::new(cortexkit_log::LineSink::open(path, retention)?)),
577 path: Arc::new(path.to_path_buf()),
578 failure_reported: Arc::new(AtomicBool::new(false)),
579 })
580 }
581}
582
583pub trait OutputSink {
586 fn write_line(&mut self, line: &[u8]);
587
588 fn stamps_lines(&self) -> bool {
597 false
598 }
599}
600
601struct StderrSink;
602
603impl OutputSink for StderrSink {
604 fn write_line(&mut self, line: &[u8]) {
605 let stderr = std::io::stderr();
606 let mut handle = stderr.lock();
607 let _ = handle.write_all(line);
608 }
609}
610
611impl OutputSink for ChildOutputSink {
612 fn write_line(&mut self, line: &[u8]) {
613 match self {
614 Self::File {
615 sink,
616 path,
617 failure_reported,
618 } => {
619 let result = sink
620 .lock()
621 .unwrap_or_else(|poisoned| poisoned.into_inner())
622 .write_line(line);
623 if let Err(error) = result {
624 if !failure_reported.swap(true, Ordering::Relaxed) {
625 tracing::warn!(
626 path = %path.display(),
627 error = %error,
628 "child output capture write failed; later failures are suppressed"
629 );
630 }
631 }
632 }
633 Self::Stderr => StderrSink.write_line(line),
634 }
635 }
636
637 fn stamps_lines(&self) -> bool {
644 matches!(self, Self::File { .. })
645 }
646}
647
648pub(crate) async fn pump_stderr_to<R, S>(
652 source: R,
653 ring: Arc<Mutex<StderrRing>>,
654 generation: u64,
655 mut sink: S,
656) where
657 R: AsyncReadExt + Unpin,
658 S: OutputSink,
659{
660 pump_lines_into(source, Some((&ring, generation)), &mut sink, "stderr").await;
661}
662
663pub(crate) async fn pump_stdout_to<R>(source: R, mut sink: ChildOutputSink)
664where
665 R: AsyncReadExt + Unpin,
666{
667 pump_lines_into(source, None, &mut sink, "stdout").await;
668}
669
670async fn pump_stderr_into<R, S>(source: R, ring: Arc<Mutex<StderrRing>>, sink: &mut S)
673where
674 R: AsyncReadExt + Unpin,
675 S: OutputSink,
676{
677 let generation = lock_ring(&ring).generation();
678 pump_lines_into(source, Some((&ring, generation)), sink, "stderr").await;
679}
680
681async fn pump_lines_into<R, S>(
682 mut source: R,
683 ring: Option<(&Arc<Mutex<StderrRing>>, u64)>,
684 sink: &mut S,
685 stream_name: &str,
686) where
687 R: AsyncReadExt + Unpin,
688 S: OutputSink,
689{
690 if let Some((ring, _)) = ring {
691 lock_ring(ring).mark_captured();
692 }
693
694 let mut pending: Vec<u8> = Vec::new();
695 let mut scanned_upto = 0usize;
699 let mut cursor = 0usize;
702 let mut chunk = [0u8; 8192];
703 loop {
704 let read = match source.read(&mut chunk).await {
705 Ok(0) => break,
706 Ok(n) => n,
707 Err(error) => {
708 if let Some((ring, generation)) = ring {
709 let mut ring = lock_ring(ring);
710 ring.mark_incomplete_from(
711 generation,
712 format!("{stream_name} read failed: {error}"),
713 );
714 ring.finish_pump(generation);
715 } else {
716 tracing::warn!(stream = stream_name, error = %error, "child output capture read failed");
717 }
718 return;
719 }
720 };
721 pending.extend_from_slice(&chunk[..read]);
722
723 while let Some(relative) = find_newline(&pending[scanned_upto..]) {
724 let newline = scanned_upto + relative;
725 emit_line(ring, sink, &pending[cursor..newline], true);
726 cursor = newline + 1;
727 scanned_upto = cursor;
728 }
729 scanned_upto = pending.len();
730
731 if cursor > 0 {
732 pending.drain(..cursor);
733 scanned_upto -= cursor;
734 cursor = 0;
735 }
736
737 if pending.len() >= MAX_PENDING_LINE_BYTES {
738 let line = std::mem::take(&mut pending);
739 emit_line(ring, sink, &line, false);
740 scanned_upto = 0;
741 }
742 }
743
744 if !pending.is_empty() {
745 emit_line(ring, sink, &pending, false);
746 }
747 if let Some((ring, generation)) = ring {
748 lock_ring(ring).finish_pump(generation);
749 }
750}
751
752#[cfg(test)]
760thread_local! {
761 static SCANNED_BYTES: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
762}
763
764#[cfg(test)]
765fn take_scanned_bytes() -> usize {
766 SCANNED_BYTES.with(|scanned| scanned.replace(0))
767}
768
769fn find_newline(haystack: &[u8]) -> Option<usize> {
772 let found = memchr::memchr(b'\n', haystack);
773 #[cfg(test)]
774 SCANNED_BYTES.with(|scanned| {
775 scanned.set(scanned.get() + found.map(|index| index + 1).unwrap_or(haystack.len()));
776 });
777 found
778}
779
780fn emit_line<S: OutputSink>(
781 ring: Option<(&Arc<Mutex<StderrRing>>, u64)>,
782 sink: &mut S,
783 raw: &[u8],
784 terminated: bool,
785) {
786 let at_ms = unix_ms(SystemTime::now());
792 if let Some((ring, generation)) = ring {
793 lock_ring(ring).push_line_from_at(generation, &String::from_utf8_lossy(raw), at_ms);
794 }
795
796 let stamp = sink.stamps_lines().then(|| format_capture_stamp(at_ms));
799 if stamp.is_none() && !terminated {
800 sink.write_line(raw);
801 return;
802 }
803 let mut framed = Vec::with_capacity(CAPTURE_STAMP_PREFIX_LEN + raw.len() + 1);
804 if let Some(stamp) = stamp {
805 framed.extend_from_slice(stamp.as_bytes());
806 framed.push(b' ');
807 }
808 framed.extend_from_slice(raw);
809 if terminated {
810 framed.push(b'\n');
811 }
812 sink.write_line(&framed);
813}
814
815fn unix_ms(at: SystemTime) -> u64 {
816 at.duration_since(UNIX_EPOCH)
819 .map(|since| u64::try_from(since.as_millis()).unwrap_or(u64::MAX))
820 .unwrap_or(0)
821}
822
823pub const CAPTURE_STAMP_LEN: usize = 24;
825
826pub const CAPTURE_STAMP_PREFIX_LEN: usize = CAPTURE_STAMP_LEN + 1;
828
829pub fn format_capture_stamp(at_ms: u64) -> String {
834 let seconds = at_ms / 1000;
835 let millis = at_ms % 1000;
836 let days = seconds / 86_400;
837 let of_day = seconds % 86_400;
838 let (year, month, day) = civil_from_days(days);
839 format!(
840 "{year:04}-{month:02}-{day:02}T{:02}:{:02}:{:02}.{millis:03}Z",
841 of_day / 3600,
842 (of_day % 3600) / 60,
843 of_day % 60,
844 )
845}
846
847pub fn split_capture_stamp(line: &str) -> Option<(u64, &str)> {
854 let bytes = line.as_bytes();
855 if bytes.len() < CAPTURE_STAMP_PREFIX_LEN || bytes[CAPTURE_STAMP_LEN] != b' ' {
856 return None;
857 }
858 let stamp = &bytes[..CAPTURE_STAMP_LEN];
859 for (index, expected) in [
860 (4, b'-'),
861 (7, b'-'),
862 (10, b'T'),
863 (13, b':'),
864 (16, b':'),
865 (19, b'.'),
866 (23, b'Z'),
867 ] {
868 if stamp[index] != expected {
869 return None;
870 }
871 }
872 let number = |from: usize, to: usize| -> Option<u64> {
873 let digits = &stamp[from..to];
874 if !digits.iter().all(u8::is_ascii_digit) {
875 return None;
876 }
877 Some(
878 digits
879 .iter()
880 .fold(0u64, |total, digit| total * 10 + u64::from(digit - b'0')),
881 )
882 };
883 let year = number(0, 4)?;
884 let month = number(5, 7)?;
885 let day = number(8, 10)?;
886 let hour = number(11, 13)?;
887 let minute = number(14, 16)?;
888 let second = number(17, 19)?;
889 let millis = number(20, 23)?;
890 if year < 1970
891 || !(1..=12).contains(&month)
892 || day == 0
893 || day > days_in_month(year, month)
894 || hour > 23
895 || minute > 59
896 || second > 59
897 {
898 return None;
899 }
900 let days = days_from_civil(year, month, day);
901 let seconds = days * 86_400 + hour * 3600 + minute * 60 + second;
902 Some((seconds * 1000 + millis, &line[CAPTURE_STAMP_PREFIX_LEN..]))
904}
905
906fn days_in_month(year: u64, month: u64) -> u64 {
907 match month {
908 2 if year.is_multiple_of(4) && (!year.is_multiple_of(100) || year.is_multiple_of(400)) => {
909 29
910 }
911 2 => 28,
912 4 | 6 | 9 | 11 => 30,
913 _ => 31,
914 }
915}
916
917fn civil_from_days(days: u64) -> (u64, u64, u64) {
921 let z = days + 719_468;
922 let era = z / 146_097;
923 let doe = z - era * 146_097;
924 let yoe = (doe - doe / 1_460 + doe / 36_524 - doe / 146_096) / 365;
925 let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
926 let mp = (5 * doy + 2) / 153;
927 let day = doy - (153 * mp + 2) / 5 + 1;
928 let month = if mp < 10 { mp + 3 } else { mp - 9 };
929 let year = yoe + era * 400 + u64::from(month <= 2);
930 (year, month, day)
931}
932
933fn days_from_civil(year: u64, month: u64, day: u64) -> u64 {
934 let year = if month <= 2 { year - 1 } else { year };
935 let era = year / 400;
936 let yoe = year - era * 400;
937 let shifted_month = if month > 2 { month - 3 } else { month + 9 };
938 let doy = (153 * shifted_month + 2) / 5 + day - 1;
939 let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
940 era * 146_097 + doe - 719_468
941}
942
943#[cfg(test)]
947pub(crate) fn untimed(entries: Vec<TailEntry>) -> Vec<TailEntry> {
948 entries
949 .into_iter()
950 .map(|entry| match entry {
951 TailEntry::Line {
952 text, truncated, ..
953 } => TailEntry::Line {
954 text,
955 truncated,
956 at_ms: None,
957 },
958 TailEntry::ProcessStart => TailEntry::ProcessStart,
959 })
960 .collect()
961}
962
963fn lock_ring(ring: &Arc<Mutex<StderrRing>>) -> std::sync::MutexGuard<'_, StderrRing> {
964 ring.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
965}
966
967fn truncate_line(line: &str, max_bytes: usize) -> (String, bool) {
973 if line.len() <= max_bytes {
974 return (line.to_string(), false);
975 }
976 let mut end = max_bytes;
977 while end > 0 && !line.is_char_boundary(end) {
978 end -= 1;
979 }
980 (line[..end].to_string(), true)
981}
982
983#[cfg(test)]
984mod tests {
985 use std::{
986 io,
987 pin::Pin,
988 task::{Context, Poll},
989 };
990
991 use super::*;
992 use tokio::io::{AsyncRead, ReadBuf};
993
994 fn ring(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> StderrRing {
995 StderrRing::new(StderrTailConfig::new(max_lines, max_bytes, max_line_bytes))
996 }
997
998 fn lines(snapshot: &StderrTailSnapshot) -> Vec<String> {
999 snapshot
1000 .entries
1001 .iter()
1002 .filter_map(|entry| match entry {
1003 TailEntry::Line { text, .. } => Some(text.clone()),
1004 TailEntry::ProcessStart => None,
1005 })
1006 .collect()
1007 }
1008
1009 #[test]
1010 fn a_fresh_ring_reports_not_captured_rather_than_empty() {
1011 let ring = ring(10, 1024, 128);
1014 let snapshot = ring.snapshot(None, None);
1015 assert!(matches!(snapshot.capture, CaptureState::NotCaptured { .. }));
1016 assert!(snapshot.entries.is_empty());
1017 }
1018
1019 #[test]
1020 fn read_failure_does_not_taint_generations_after_its_section_is_evicted() {
1021 let mut ring = StderrRing::new(StderrTailConfig {
1022 max_lines: 1,
1023 ..Default::default()
1024 });
1025 ring.begin_process();
1026 ring.mark_captured();
1027 ring.push_line("old output");
1028 ring.mark_incomplete("read failed");
1029 assert!(matches!(
1030 ring.snapshot(None, None).capture,
1031 CaptureState::Incomplete { .. }
1032 ));
1033 ring.begin_process();
1034 ring.mark_captured();
1035 ring.push_line("new output");
1036 assert_eq!(ring.snapshot(None, None).capture, CaptureState::Captured);
1037 }
1038
1039 #[test]
1040 fn a_captured_module_that_printed_nothing_is_distinguishable_from_an_uncaptured_one() {
1041 let mut captured = ring(10, 1024, 128);
1042 captured.mark_captured();
1043 let uncaptured = ring(10, 1024, 128);
1044
1045 let captured = captured.snapshot(None, None);
1046 let uncaptured = uncaptured.snapshot(None, None);
1047
1048 assert!(captured.entries.is_empty());
1051 assert!(uncaptured.entries.is_empty());
1052 assert_eq!(captured.capture, CaptureState::Captured);
1053 assert!(matches!(
1054 uncaptured.capture,
1055 CaptureState::NotCaptured { .. }
1056 ));
1057 }
1058
1059 #[test]
1060 fn the_line_cap_evicts_oldest_first_and_counts_what_it_dropped() {
1061 let mut ring = ring(3, 10_000, 128);
1062 ring.mark_captured();
1063 for i in 0..6 {
1064 ring.push_line(&format!("line{i}"));
1065 }
1066 let snapshot = ring.snapshot(None, None);
1067 assert_eq!(lines(&snapshot), vec!["line3", "line4", "line5"]);
1068 assert_eq!(snapshot.dropped_lines, 3);
1071 }
1072
1073 #[test]
1074 fn the_byte_cap_binds_before_the_line_cap_when_lines_are_large() {
1075 let mut ring = ring(100, 30, 128);
1077 ring.mark_captured();
1078 for i in 0..10 {
1079 ring.push_line(&format!("{i}--------")); }
1081 let snapshot = ring.snapshot(None, None);
1082 assert!(
1083 snapshot.entries.len() < 10,
1084 "byte cap did not bind: {} entries retained",
1085 snapshot.entries.len()
1086 );
1087 let retained: usize = lines(&snapshot).iter().map(String::len).sum();
1088 assert!(
1089 retained <= 30,
1090 "retained {retained} bytes over a 30 byte cap"
1091 );
1092 assert!(snapshot.dropped_lines > 0);
1093 }
1094
1095 #[test]
1096 fn one_enormous_line_is_truncated_rather_than_evicting_the_tail() {
1097 let mut ring = ring(10, 10_000, 64);
1100 ring.mark_captured();
1101 ring.push_line("context line that must survive");
1102 ring.push_line(&"x".repeat(40_000));
1103
1104 let snapshot = ring.snapshot(None, None);
1105 let kept = &snapshot.entries;
1106 assert!(matches!(
1107 &kept[0],
1108 TailEntry::Line { text, truncated: false, .. }
1109 if text == "context line that must survive"
1110 ));
1111 let TailEntry::Line {
1112 text, truncated, ..
1113 } = &kept[1]
1114 else {
1115 panic!("expected a truncated line");
1116 };
1117 assert_eq!(text, &"x".repeat(64));
1118 assert!(*truncated);
1119 }
1120
1121 #[test]
1122 fn truncation_is_visible_so_a_cut_line_is_not_mistaken_for_a_short_one() {
1123 let mut ring = ring(10, 10_000, 16);
1124 ring.mark_captured();
1125 ring.push_line("0123456789abcdefghij");
1126 ring.push_line("short");
1127
1128 let snapshot = ring.snapshot(None, None);
1129 let TailEntry::Line { truncated, .. } = &snapshot.entries[0] else {
1130 panic!("expected a line");
1131 };
1132 assert!(truncated);
1133 let TailEntry::Line { truncated, .. } = &snapshot.entries[1] else {
1134 panic!("expected a line");
1135 };
1136 assert!(!truncated, "a short line must not be reported as truncated");
1137 }
1138
1139 #[test]
1140 fn truncation_cuts_on_a_char_boundary_rather_than_splitting_utf8() {
1141 let mut ring = ring(10, 10_000, 5);
1144 ring.mark_captured();
1145 ring.push_line("aa€€€€");
1146 let snapshot = ring.snapshot(None, None);
1147 let TailEntry::Line {
1148 text, truncated, ..
1149 } = &snapshot.entries[0]
1150 else {
1151 panic!("expected a line");
1152 };
1153 assert!(truncated);
1154 assert!(text.starts_with("aa"));
1155 }
1156
1157 #[test]
1158 fn a_restart_boundary_keeps_generations_distinguishable() {
1159 let mut ring = ring(10, 10_000, 128);
1160 ring.mark_captured();
1161 ring.push_line("before the crash");
1162 ring.push_process_start();
1163 ring.push_line("after the respawn");
1164
1165 let snapshot = ring.snapshot(None, None);
1166 assert_eq!(
1167 snapshot.entries,
1168 vec![
1169 TailEntry::Line {
1170 text: "before the crash".to_string(),
1171 truncated: false,
1172 at_ms: None,
1173 },
1174 TailEntry::ProcessStart,
1175 TailEntry::Line {
1176 text: "after the respawn".to_string(),
1177 truncated: false,
1178 at_ms: None,
1179 },
1180 ]
1181 );
1182 }
1183
1184 #[test]
1185 fn the_ring_survives_respawn_because_the_cause_is_written_before_the_exit() {
1186 let mut ring = ring(10, 10_000, 128);
1189 ring.mark_captured();
1190 ring.push_line("Error: storage section missing");
1191 ring.push_process_start();
1192
1193 let snapshot = ring.snapshot(None, None);
1194 assert!(lines(&snapshot).contains(&"Error: storage section missing".to_string()));
1195 }
1196
1197 #[test]
1198 fn a_caller_limit_returns_the_newest_lines_not_the_oldest() {
1199 let mut ring = ring(100, 100_000, 128);
1200 ring.mark_captured();
1201 for i in 0..10 {
1202 ring.push_line(&format!("line{i}"));
1203 }
1204 let snapshot = ring.snapshot(Some(3), None);
1205 assert_eq!(lines(&snapshot), vec!["line7", "line8", "line9"]);
1206 }
1207
1208 #[test]
1209 fn a_caller_line_limit_keeps_the_boundary_before_the_selected_line() {
1210 let mut ring = ring(100, 100_000, 128);
1211 ring.mark_captured();
1212 ring.push_line("before restart");
1213 ring.push_process_start();
1214 ring.push_line("after restart");
1215
1216 let snapshot = ring.snapshot(Some(1), None);
1217 assert_eq!(
1218 snapshot.entries,
1219 vec![
1220 TailEntry::ProcessStart,
1221 TailEntry::Line {
1222 text: "after restart".to_string(),
1223 truncated: false,
1224 at_ms: None,
1225 },
1226 ]
1227 );
1228 }
1229
1230 #[test]
1231 fn a_caller_line_limit_omits_a_trailing_boundary_after_the_selected_line() {
1232 let mut ring = ring(100, 100_000, 128);
1233 ring.mark_captured();
1234 ring.push_line("before restart");
1235 ring.push_process_start();
1236
1237 let snapshot = ring.snapshot(Some(1), None);
1238 assert_eq!(
1239 snapshot.entries,
1240 vec![TailEntry::Line {
1241 text: "before restart".to_string(),
1242 truncated: false,
1243 at_ms: None,
1244 }]
1245 );
1246 }
1247
1248 #[test]
1249 fn a_caller_limit_reports_what_it_withheld_rather_than_looking_complete() {
1250 let mut ring = ring(100, 100_000, 128);
1251 ring.mark_captured();
1252 for i in 0..10 {
1253 ring.push_line(&format!("line{i}"));
1254 }
1255 assert_eq!(ring.snapshot(Some(3), None).dropped_lines, 7);
1258 assert_eq!(ring.snapshot(None, None).dropped_lines, 0);
1259 }
1260
1261 #[test]
1262 fn a_caller_limit_cannot_widen_the_rings_own_caps() {
1263 let mut ring = ring(2, 10_000, 128);
1264 ring.mark_captured();
1265 for i in 0..5 {
1266 ring.push_line(&format!("line{i}"));
1267 }
1268 let snapshot = ring.snapshot(Some(1000), Some(1_000_000));
1269 assert_eq!(lines(&snapshot).len(), 2);
1270 }
1271
1272 fn shared(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> Arc<Mutex<StderrRing>> {
1273 Arc::new(Mutex::new(ring(max_lines, max_bytes, max_line_bytes)))
1274 }
1275
1276 #[derive(Default)]
1279 struct RecordingSink {
1280 writes: Vec<Vec<u8>>,
1281 }
1282
1283 impl OutputSink for RecordingSink {
1284 fn write_line(&mut self, line: &[u8]) {
1285 self.writes.push(line.to_vec());
1286 }
1287 }
1288
1289 struct ChunkedReader {
1292 chunks: VecDeque<Vec<u8>>,
1293 }
1294
1295 impl AsyncRead for ChunkedReader {
1296 fn poll_read(
1297 mut self: Pin<&mut Self>,
1298 _cx: &mut Context<'_>,
1299 buf: &mut ReadBuf<'_>,
1300 ) -> Poll<io::Result<()>> {
1301 match self.chunks.pop_front() {
1302 None => Poll::Ready(Ok(())),
1303 Some(chunk) => {
1304 buf.put_slice(&chunk);
1305 Poll::Ready(Ok(()))
1306 }
1307 }
1308 }
1309 }
1310
1311 struct FailingReader {
1312 bytes: Vec<u8>,
1313 emitted: bool,
1314 }
1315
1316 impl AsyncRead for FailingReader {
1317 fn poll_read(
1318 mut self: Pin<&mut Self>,
1319 _cx: &mut Context<'_>,
1320 buf: &mut ReadBuf<'_>,
1321 ) -> Poll<io::Result<()>> {
1322 if self.emitted {
1323 return Poll::Ready(Err(io::Error::other("reader failed")));
1324 }
1325 self.emitted = true;
1326 buf.put_slice(&self.bytes);
1327 Poll::Ready(Ok(()))
1328 }
1329 }
1330
1331 #[tokio::test]
1332 async fn the_pump_splits_on_newlines_and_keeps_a_trailing_fragment() {
1333 let ring = shared(10, 10_000, 128);
1334 let source = std::io::Cursor::new(b"one\ntwo\nthree".to_vec());
1337 let mut sink = RecordingSink::default();
1338 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1339
1340 let snapshot = lock_ring(&ring).snapshot(None, None);
1341 assert_eq!(lines(&snapshot), vec!["one", "two", "three"]);
1342 assert_eq!(snapshot.capture, CaptureState::Captured);
1343 assert_eq!(
1344 sink.writes,
1345 vec![b"one\n".to_vec(), b"two\n".to_vec(), b"three".to_vec()]
1346 );
1347 }
1348
1349 #[tokio::test]
1350 async fn a_read_failure_keeps_prior_lines_and_marks_the_capture_incomplete() {
1351 let ring = shared(10, 10_000, 128);
1352 let source = FailingReader {
1353 bytes: b"crash cause\n".to_vec(),
1354 emitted: false,
1355 };
1356 let mut sink = RecordingSink::default();
1357 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1358
1359 let snapshot = lock_ring(&ring).snapshot(None, None);
1360 assert_eq!(lines(&snapshot), vec!["crash cause"]);
1361 assert!(matches!(
1362 snapshot.capture,
1363 CaptureState::Incomplete { ref reason } if reason.contains("reader failed")
1364 ));
1365 assert_eq!(sink.writes, vec![b"crash cause\n".to_vec()]);
1366 }
1367
1368 #[tokio::test]
1369 async fn every_captured_line_is_also_forwarded() {
1370 let ring = shared(10, 10_000, 128);
1374 let source = std::io::Cursor::new(b"alpha\nbeta\n".to_vec());
1375 let mut sink = RecordingSink::default();
1376 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1377
1378 assert_eq!(sink.writes, vec![b"alpha\n".to_vec(), b"beta\n".to_vec()]);
1379 }
1380
1381 #[tokio::test]
1382 async fn each_forwarded_line_is_exactly_one_write() {
1383 let ring = shared(10, 10_000, 128);
1388 let source = std::io::Cursor::new(b"first\nsecond\nthird\n".to_vec());
1389 let mut sink = RecordingSink::default();
1390 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1391
1392 assert_eq!(sink.writes.len(), 3);
1393 for write in &sink.writes {
1394 assert_eq!(
1395 write.iter().filter(|byte| **byte == b'\n').count(),
1396 1,
1397 "a write carried something other than exactly one complete line"
1398 );
1399 assert_eq!(*write.last().unwrap(), b'\n');
1400 }
1401 }
1402
1403 #[test]
1404 fn the_first_process_start_is_not_recorded_because_it_divides_nothing() {
1405 let mut ring = ring(10, 10_000, 128);
1408 ring.push_process_start();
1409 assert!(ring.snapshot(None, None).entries.is_empty());
1410
1411 ring.push_line("first process said this");
1412 ring.push_process_start();
1413 assert!(
1414 matches!(ring.entries.back(), Some(Slot::ProcessStart { .. })),
1415 "a boundary with output before it must be recorded"
1416 );
1417 ring.push_line("second process said this");
1418 assert_eq!(
1419 ring.snapshot(None, None).entries,
1420 vec![
1421 TailEntry::Line {
1422 text: "first process said this".to_string(),
1423 truncated: false,
1424 at_ms: None,
1425 },
1426 TailEntry::ProcessStart,
1427 TailEntry::Line {
1428 text: "second process said this".to_string(),
1429 truncated: false,
1430 at_ms: None,
1431 },
1432 ],
1433 "only the boundary with output before it may be shown"
1434 );
1435 }
1436
1437 #[test]
1438 fn a_process_start_is_recorded_when_only_dropped_lines_precede_it() {
1439 let mut ring = ring(1, 10_000, 128);
1443 ring.push_line("evicted");
1444 ring.push_line("also evicted");
1445 ring.entries.clear();
1449 ring.lines = 0;
1450 ring.bytes = 0;
1451 ring.push_process_start();
1452 ring.push_line("survivor");
1453 assert_eq!(
1454 ring.snapshot(None, None).entries,
1455 vec![
1456 TailEntry::ProcessStart,
1457 TailEntry::Line {
1458 text: "survivor".to_string(),
1459 truncated: false,
1460 at_ms: None,
1461 },
1462 ]
1463 );
1464 }
1465
1466 fn line(text: &str) -> TailEntry {
1467 TailEntry::Line {
1468 text: text.to_string(),
1469 truncated: false,
1470 at_ms: None,
1471 }
1472 }
1473
1474 #[test]
1475 fn a_late_line_from_a_retired_process_lands_in_that_processs_section() {
1476 let mut ring = ring(10, 10_000, 128);
1480 ring.mark_captured();
1481 let old = ring.begin_process();
1482 ring.push_line_from(old, "old: booting");
1483 ring.retire_pump(old);
1484 let new = ring.begin_process();
1485 ring.push_line_from(new, "new: booting");
1486 ring.push_line_from(old, "old: config error");
1487
1488 assert_eq!(
1489 ring.snapshot(None, None).entries,
1490 vec![
1491 line("old: booting"),
1492 line("old: config error"),
1493 TailEntry::ProcessStart,
1494 line("new: booting"),
1495 ]
1496 );
1497 }
1498
1499 #[test]
1500 fn a_line_from_a_process_that_was_not_retired_is_appended_as_it_arrives() {
1501 let mut ring = ring(10, 10_000, 128);
1504 ring.mark_captured();
1505 let incumbent = ring.begin_process();
1506 ring.push_line_from(incumbent, "incumbent: before");
1507 let candidate = ring.begin_process();
1508 ring.push_line_from(candidate, "candidate: booting");
1509 ring.push_line_from(incumbent, "incumbent: still serving");
1510
1511 assert_eq!(
1512 ring.snapshot(None, None).entries,
1513 vec![
1514 line("incumbent: before"),
1515 TailEntry::ProcessStart,
1516 line("candidate: booting"),
1517 line("incumbent: still serving"),
1518 ]
1519 );
1520 }
1521
1522 #[test]
1523 fn a_late_line_keeps_its_section_when_the_process_had_printed_nothing_before() {
1524 let mut ring = ring(10, 10_000, 128);
1528 ring.mark_captured();
1529 let first = ring.begin_process();
1530 ring.push_line_from(first, "first: done");
1531 ring.finish_pump(first);
1532 let old = ring.begin_process();
1533 ring.retire_pump(old);
1534 let new = ring.begin_process();
1535 ring.push_line_from(new, "new: booting");
1536 ring.push_line_from(old, "old: config error");
1537
1538 assert_eq!(
1539 ring.snapshot(None, None).entries,
1540 vec![
1541 line("first: done"),
1542 TailEntry::ProcessStart,
1543 line("old: config error"),
1544 TailEntry::ProcessStart,
1545 line("new: booting"),
1546 ]
1547 );
1548 }
1549
1550 #[test]
1551 fn a_late_line_whose_section_was_evicted_counts_as_dropped() {
1552 let mut ring = ring(2, 10_000, 128);
1553 ring.mark_captured();
1554 let old = ring.begin_process();
1555 ring.push_line_from(old, "old");
1556 ring.retire_pump(old);
1557 let new = ring.begin_process();
1558 for text in ["new 1", "new 2", "new 3"] {
1559 ring.push_line_from(new, text);
1560 }
1561 ring.push_line_from(old, "old, late");
1564
1565 let snapshot = ring.snapshot(None, None);
1566 assert_eq!(snapshot.entries, vec![line("new 2"), line("new 3")]);
1567 assert_eq!(snapshot.dropped_lines, 3);
1568 }
1569
1570 #[test]
1571 fn a_late_reader_reads_incomplete_until_its_pipe_reaches_eof() {
1572 let mut ring = ring(10, 10_000, 128);
1573 ring.mark_captured();
1574 let old = ring.begin_process();
1575 ring.retire_pump(old);
1576 ring.mark_pump_late(old, "still open");
1577 ring.begin_process();
1578 assert_eq!(
1579 ring.snapshot(None, None).capture,
1580 CaptureState::Incomplete {
1581 reason: "still open".to_string()
1582 }
1583 );
1584
1585 ring.finish_pump(old);
1586 assert_eq!(ring.snapshot(None, None).capture, CaptureState::Captured);
1587 }
1588
1589 #[test]
1590 fn silent_restarts_do_not_grow_the_ring() {
1591 let mut ring = ring(10, 10_000, 128);
1592 ring.mark_captured();
1593 ring.push_line("once");
1594 for _ in 0..100 {
1595 let generation = ring.begin_process();
1596 ring.finish_pump(generation);
1597 }
1598 assert_eq!(ring.entries.len(), 2);
1599 }
1600
1601 #[tokio::test]
1602 async fn the_pump_marks_captured_even_when_the_module_writes_nothing() {
1603 let ring = shared(10, 10_000, 128);
1606 let source = std::io::Cursor::new(Vec::new());
1607 let mut sink = RecordingSink::default();
1608 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1609
1610 let snapshot = lock_ring(&ring).snapshot(None, None);
1611 assert!(snapshot.entries.is_empty());
1612 assert_eq!(snapshot.capture, CaptureState::Captured);
1613 assert!(sink.writes.is_empty());
1614 }
1615
1616 #[tokio::test]
1617 async fn a_line_with_no_newline_cannot_grow_the_reader_without_bound() {
1618 let ring = shared(10, 10_000_000, 4 * 1024 * 1024);
1621 let source = std::io::Cursor::new(vec![b'x'; MAX_PENDING_LINE_BYTES + 4096]);
1622 let mut sink = RecordingSink::default();
1623 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1624
1625 let snapshot = lock_ring(&ring).snapshot(None, None);
1626 assert_eq!(
1627 lines(&snapshot).len(),
1628 2,
1629 "expected a forced flush at the ceiling plus the remainder"
1630 );
1631 assert_eq!(
1632 sink.writes,
1633 vec![vec![b'x'; MAX_PENDING_LINE_BYTES], vec![b'x'; 4096],],
1634 "forced flushes and EOF fragments must not invent delimiters"
1635 );
1636 }
1637
1638 #[tokio::test]
1639 async fn boundaries_truncation_and_framing_do_not_depend_on_chunk_splits() {
1640 let ring = shared(100, 100_000, 8);
1644 let source = ChunkedReader {
1645 chunks: vec![
1646 b"fir".to_vec(),
1647 b"st\nsec".to_vec(),
1648 b"ond\ncarry\r".to_vec(),
1649 b"\nover\n".to_vec(),
1650 b"12345678\n".to_vec(),
1651 b"1234567".to_vec(),
1652 b"89\n".to_vec(),
1653 b"tail".to_vec(),
1654 ]
1655 .into_iter()
1656 .collect(),
1657 };
1658 let mut sink = RecordingSink::default();
1659 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1660
1661 let snapshot = lock_ring(&ring).snapshot(None, None);
1662 assert_eq!(snapshot.capture, CaptureState::Captured);
1663 assert_eq!(
1664 untimed(snapshot.entries),
1665 vec![
1666 TailEntry::Line {
1667 text: "first".to_string(),
1668 truncated: false,
1669 at_ms: None,
1670 },
1671 TailEntry::Line {
1672 text: "second".to_string(),
1673 truncated: false,
1674 at_ms: None,
1675 },
1676 TailEntry::Line {
1678 text: "carry\r".to_string(),
1679 truncated: false,
1680 at_ms: None,
1681 },
1682 TailEntry::Line {
1683 text: "over".to_string(),
1684 truncated: false,
1685 at_ms: None,
1686 },
1687 TailEntry::Line {
1689 text: "12345678".to_string(),
1690 truncated: false,
1691 at_ms: None,
1692 },
1693 TailEntry::Line {
1695 text: "12345678".to_string(),
1696 truncated: true,
1697 at_ms: None,
1698 },
1699 TailEntry::Line {
1700 text: "tail".to_string(),
1701 truncated: false,
1702 at_ms: None,
1703 },
1704 ]
1705 );
1706 assert_eq!(
1707 sink.writes,
1708 vec![
1709 b"first\n".to_vec(),
1710 b"second\n".to_vec(),
1711 b"carry\r\n".to_vec(),
1712 b"over\n".to_vec(),
1713 b"12345678\n".to_vec(),
1714 b"123456789\n".to_vec(),
1715 b"tail".to_vec(),
1716 ]
1717 );
1718 }
1719
1720 #[tokio::test]
1721 async fn a_line_with_no_newline_is_not_rescanned_from_byte_zero_on_every_chunk() {
1722 let input = vec![b'x'; MAX_PENDING_LINE_BYTES + 4096];
1728 let ring = shared(10, 10_000_000, 4 * 1024 * 1024);
1729 let source = std::io::Cursor::new(input.clone());
1730 let mut sink = RecordingSink::default();
1731
1732 take_scanned_bytes();
1733 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1734 let scanned = take_scanned_bytes();
1735
1736 assert!(
1737 scanned <= 2 * input.len(),
1738 "newline searches examined {scanned} bytes for {} bytes of input; \
1739 each chunk must search only newly arrived bytes",
1740 input.len()
1741 );
1742 }
1743
1744 fn now_ms() -> u64 {
1745 unix_ms(SystemTime::now())
1746 }
1747
1748 #[tokio::test]
1749 async fn a_captured_line_carries_the_time_the_reader_framed_it() {
1750 let ring = shared(10, 10_000, 128);
1754 let before = now_ms();
1755 let source = std::io::Cursor::new(b"first\nsecond".to_vec());
1756 pump_stderr_into(source, Arc::clone(&ring), &mut RecordingSink::default()).await;
1757 let after = now_ms();
1758
1759 let entries = lock_ring(&ring).snapshot(None, None).entries;
1760 assert_eq!(entries.len(), 2);
1761 for entry in entries {
1762 let TailEntry::Line { text, at_ms, .. } = entry else {
1763 panic!("expected only lines, got {entry:?}");
1764 };
1765 let at_ms = at_ms.unwrap_or_else(|| panic!("line {text:?} has no capture time"));
1766 assert!(
1767 (before..=after).contains(&at_ms),
1768 "line {text:?} stamped {at_ms}, outside the pump's run {before}..={after}"
1769 );
1770 }
1771 }
1772
1773 async fn capture_through_file_sink(chunks: Vec<Vec<u8>>) -> (String, u64, u64) {
1775 let temp = subc_test_support::TestTempDir::new("stderr-capture-stamp");
1776 let path = temp.path().join("stamped.stderr.log");
1777 let sink = ChildOutputSink::open(&path, cortexkit_log::Retention::default()).unwrap();
1778 let ring = shared(10, 10_000, 128);
1779 let generation = lock_ring(&ring).begin_process();
1780 let before = now_ms();
1781 pump_stderr_to(
1782 ChunkedReader {
1783 chunks: chunks.into_iter().collect(),
1784 },
1785 ring,
1786 generation,
1787 sink,
1788 )
1789 .await;
1790 let after = now_ms();
1791 (std::fs::read_to_string(&path).unwrap(), before, after)
1792 }
1793
1794 fn assert_stamp_shape(line: &str) {
1798 let bytes = line.as_bytes();
1799 assert!(bytes.len() > 25, "line too short for a stamp: {line:?}");
1800 for (index, byte) in bytes[..25].iter().enumerate() {
1801 let expected_separator = match index {
1802 4 | 7 => Some(b'-'),
1803 10 => Some(b'T'),
1804 13 | 16 => Some(b':'),
1805 19 => Some(b'.'),
1806 23 => Some(b'Z'),
1807 24 => Some(b' '),
1808 _ => None,
1809 };
1810 match expected_separator {
1811 Some(separator) => assert_eq!(*byte, separator, "byte {index} of {line:?}"),
1812 None => assert!(byte.is_ascii_digit(), "byte {index} of {line:?}"),
1813 }
1814 }
1815 }
1816
1817 #[tokio::test]
1818 async fn the_capture_file_stamps_each_line_and_keeps_the_module_bytes_verbatim() {
1819 let module_lines = [
1823 "plain line",
1824 "2020-01-01T00:00:00.000Z INFO mymod: own stamp",
1825 ];
1826 let input = format!("{}\n{}\n", module_lines[0], module_lines[1]);
1827 let (contents, before, after) = capture_through_file_sink(vec![input.into_bytes()]).await;
1828
1829 assert!(contents.ends_with('\n'));
1830 let lines: Vec<&str> = contents.lines().collect();
1831 assert_eq!(lines.len(), 2, "capture file: {contents:?}");
1832 for (line, module_line) in lines.iter().zip(module_lines) {
1833 assert_stamp_shape(line);
1834 assert_eq!(&line[25..], module_line, "module bytes changed");
1835 let (at_ms, rest) = split_capture_stamp(line).unwrap();
1836 assert_eq!(rest, module_line);
1837 assert!(
1839 (before..=after).contains(&at_ms),
1840 "stamped {at_ms}, outside the pump's run {before}..={after}"
1841 );
1842 }
1843 }
1844
1845 #[tokio::test]
1846 async fn a_line_split_across_several_writes_is_stamped_once() {
1847 let (contents, _, _) = capture_through_file_sink(vec![
1851 b"par".to_vec(),
1852 b"tial li".to_vec(),
1853 b"ne\nwhole\n".to_vec(),
1854 ])
1855 .await;
1856
1857 let lines: Vec<&str> = contents.lines().collect();
1858 assert_eq!(lines.len(), 2, "capture file: {contents:?}");
1859 for (line, module_line) in lines.iter().zip(["partial line", "whole"]) {
1860 assert_stamp_shape(line);
1861 assert_eq!(&line[25..], module_line, "capture file: {contents:?}");
1862 }
1863 }
1864
1865 #[tokio::test]
1866 async fn a_line_flushed_at_the_ceiling_is_stamped_only_at_the_start_of_each_file_line() {
1867 let mut long = vec![b'x'; MAX_PENDING_LINE_BYTES + 100];
1872 long.push(b'\n');
1873 let chunks = long.chunks(8192).map(<[u8]>::to_vec).collect();
1875 let (contents, _, _) = capture_through_file_sink(chunks).await;
1876
1877 let lines: Vec<&str> = contents.lines().collect();
1878 assert_eq!(lines.len(), 2, "expected the flushed piece and the rest");
1879 for line in &lines {
1880 assert_stamp_shape(line);
1881 assert!(
1882 line[25..].bytes().all(|byte| byte == b'x'),
1883 "a stamp landed inside module bytes"
1884 );
1885 }
1886 let module_bytes: usize = lines.iter().map(|line| line.len() - 25).sum();
1887 assert_eq!(module_bytes, MAX_PENDING_LINE_BYTES + 100);
1888 }
1889
1890 #[test]
1891 fn the_capture_stamp_is_the_daemon_log_timestamp_form() {
1892 for (at_ms, text) in [
1893 (0, "1970-01-01T00:00:00.000Z"),
1894 (951_868_799_999, "2000-02-29T23:59:59.999Z"),
1895 (1_789_801_440_685, "2026-09-19T07:04:00.685Z"),
1896 (4_107_542_400_001, "2100-03-01T00:00:00.001Z"),
1897 ] {
1898 let stamp = format_capture_stamp(at_ms);
1899 assert_eq!(stamp, text);
1900 assert_eq!(stamp.len(), CAPTURE_STAMP_LEN);
1901 let daemon_line = format!("{stamp} INFO subc: probe");
1904 let parsed = cortexkit_log::parse_line(&daemon_line).unwrap();
1905 assert_eq!(
1906 parsed.timestamp,
1907 UNIX_EPOCH + std::time::Duration::from_millis(at_ms)
1908 );
1909 assert_eq!(
1910 split_capture_stamp(&format!("{stamp} body")),
1911 Some((at_ms, "body"))
1912 );
1913 }
1914 }
1915
1916 #[test]
1917 fn a_line_without_a_well_formed_stamp_has_no_capture_time() {
1918 for line in [
1919 "",
1920 "plain module output",
1921 "2026-09-19T07:04:00.685Z",
1922 "2026-09-19T07:04:00.685Zbody",
1923 "2026-09-19T07:04:00.685z body",
1924 "2026-09-19 07:04:00.685Z body",
1925 "2026-09-19T07:04:00Z body",
1926 "2026-02-30T07:04:00.685Z body",
1927 "2026-13-19T07:04:00.685Z body",
1928 "2026-09-19T24:04:00.685Z body",
1929 "2026-09-19T07:04:00.6a5Z body",
1930 "1969-12-31T23:59:59.999Z body",
1931 ] {
1932 assert_eq!(split_capture_stamp(line), None, "{line:?}");
1933 }
1934 }
1935
1936 #[test]
1937 fn a_byte_limit_smaller_than_one_line_still_returns_that_line() {
1938 let mut ring = ring(10, 10_000, 128);
1941 ring.mark_captured();
1942 ring.push_line("a line considerably longer than the request limit");
1943 let snapshot = ring.snapshot(None, Some(4));
1944 assert_eq!(snapshot.entries.len(), 1);
1945 }
1946
1947 #[test]
1948 fn an_incoherent_config_clamps_the_line_cap_and_keeps_its_restart_boundary() {
1949 let config = StderrTailConfig::new(2, 10, 100);
1950 assert_eq!(config.max_line_bytes, config.max_bytes);
1951 let mut ring = StderrRing::new(config);
1952 ring.mark_captured();
1953 ring.push_line("old");
1954 ring.push_process_start();
1955 ring.push_line("new process line longer than the ring byte cap");
1956
1957 assert_eq!(
1958 ring.snapshot(None, None).entries,
1959 vec![
1960 TailEntry::ProcessStart,
1961 TailEntry::Line {
1962 text: "new proces".to_string(),
1963 truncated: true,
1964 at_ms: None,
1965 },
1966 ]
1967 );
1968 }
1969}