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 generation: u64,
220 pumps: BTreeMap<u64, PumpPhase>,
222 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 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 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 pub fn push_process_start(&mut self) -> u64 {
282 self.generation += 1;
283 let generation = self.generation;
284 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 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 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 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 pub(crate) fn finish_pump(&mut self, generation: u64) {
332 self.pumps.remove(&generation);
333 }
334
335 pub fn push_line(&mut self, line: &str) {
342 self.push_line_from(self.generation, line);
343 }
344
345 pub(crate) fn push_line_from(&mut self, generation: u64, line: &str) {
347 self.admit_line(generation, line, None);
348 }
349
350 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 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 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 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 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 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 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 dropped_lines: self.dropped_lines + withheld as u64,
511 }
512 }
513}
514
515const MAX_PENDING_LINE_BYTES: usize = 1024 * 1024;
524
525pub 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#[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
560pub trait OutputSink {
563 fn write_line(&mut self, line: &[u8]);
564
565 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 fn stamps_lines(&self) -> bool {
621 matches!(self, Self::File { .. })
622 }
623}
624
625pub(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
647async 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 let mut scanned_upto = 0usize;
676 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#[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
743fn 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 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 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 at.duration_since(UNIX_EPOCH)
793 .map(|since| u64::try_from(since.as_millis()).unwrap_or(u64::MAX))
794 .unwrap_or(0)
795}
796
797pub const CAPTURE_STAMP_LEN: usize = 24;
799
800pub const CAPTURE_STAMP_PREFIX_LEN: usize = CAPTURE_STAMP_LEN + 1;
802
803pub 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
821pub 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 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
891fn 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#[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
941fn 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 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 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 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 let mut ring = ring(100, 30, 128);
1031 ring.mark_captured();
1032 for i in 0..10 {
1033 ring.push_line(&format!("{i}--------")); }
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 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 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 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 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 #[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 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 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 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 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 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 let mut ring = ring(1, 10_000, 128);
1397 ring.push_line("evicted");
1398 ring.push_line("also evicted");
1399 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 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 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 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 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 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 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 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 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 TailEntry::Line {
1643 text: "12345678".to_string(),
1644 truncated: false,
1645 at_ms: None,
1646 },
1647 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 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 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 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 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 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 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 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 let mut long = vec![b'x'; MAX_PENDING_LINE_BYTES + 100];
1826 long.push(b'\n');
1827 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 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 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}