1use std::collections::VecDeque;
22use std::io::{self, Write};
23use std::path::{Path, PathBuf};
24use std::sync::{
25 atomic::{AtomicBool, Ordering},
26 Arc, Mutex,
27};
28
29use tokio::io::AsyncReadExt;
30
31pub const DEFAULT_MAX_LINE_BYTES: usize = 2048;
38
39pub const DEFAULT_MAX_LINES: usize = 200;
41
42pub const DEFAULT_MAX_BYTES: usize = 64 * 1024;
48
49#[derive(Debug, Clone, PartialEq, Eq)]
57pub enum CaptureState {
58 Captured,
60 Incomplete { reason: String },
62 NotCaptured { reason: String },
64}
65
66#[derive(Debug, Clone, PartialEq, Eq)]
73pub enum TailEntry {
74 Line {
75 text: String,
76 truncated: bool,
78 },
79 ProcessStart,
82}
83
84impl TailEntry {
85 fn cost(&self) -> usize {
86 match self {
87 Self::Line { text, .. } => text.len(),
88 Self::ProcessStart => 0,
89 }
90 }
91}
92
93#[derive(Debug, Clone, Copy, PartialEq, Eq)]
94pub struct StderrTailConfig {
95 max_lines: usize,
96 max_bytes: usize,
97 max_line_bytes: usize,
98}
99
100impl StderrTailConfig {
101 pub const fn new(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> Self {
104 Self {
105 max_lines,
106 max_bytes,
107 max_line_bytes: if max_line_bytes > max_bytes {
108 max_bytes
109 } else {
110 max_line_bytes
111 },
112 }
113 }
114}
115
116impl Default for StderrTailConfig {
117 fn default() -> Self {
118 Self::new(DEFAULT_MAX_LINES, DEFAULT_MAX_BYTES, DEFAULT_MAX_LINE_BYTES)
119 }
120}
121
122#[derive(Debug, Clone, PartialEq, Eq)]
124pub struct StderrTailSnapshot {
125 pub capture: CaptureState,
126 pub entries: Vec<TailEntry>,
127 pub dropped_lines: u64,
134}
135
136impl StderrTailSnapshot {
137 pub fn not_captured(reason: impl Into<String>) -> Self {
139 Self {
140 capture: CaptureState::NotCaptured {
141 reason: reason.into(),
142 },
143 entries: Vec::new(),
144 dropped_lines: 0,
145 }
146 }
147}
148
149#[derive(Debug)]
156pub struct StderrRing {
157 config: StderrTailConfig,
158 entries: VecDeque<TailEntry>,
159 lines: usize,
163 bytes: usize,
164 dropped_lines: u64,
165 capture: CaptureState,
166}
167
168impl StderrRing {
169 pub fn new(config: StderrTailConfig) -> Self {
170 Self {
171 config,
172 entries: VecDeque::new(),
173 lines: 0,
174 bytes: 0,
175 dropped_lines: 0,
176 capture: CaptureState::NotCaptured {
179 reason: "stderr reader has not started".to_string(),
180 },
181 }
182 }
183
184 pub fn mark_captured(&mut self) {
185 if matches!(self.capture, CaptureState::NotCaptured { .. }) {
186 self.capture = CaptureState::Captured;
187 }
188 }
189
190 pub fn mark_incomplete(&mut self, reason: impl Into<String>) {
191 self.capture = CaptureState::Incomplete {
192 reason: reason.into(),
193 };
194 }
195
196 pub fn mark_not_captured(&mut self, reason: impl Into<String>) {
197 self.capture = CaptureState::NotCaptured {
198 reason: reason.into(),
199 };
200 }
201
202 pub fn push_process_start(&mut self) {
211 if self.entries.is_empty() && self.dropped_lines == 0 {
212 return;
213 }
214 if matches!(self.entries.back(), Some(TailEntry::ProcessStart)) {
215 return;
216 }
217 self.push_entry(TailEntry::ProcessStart);
218 }
219
220 pub fn push_line(&mut self, line: &str) {
225 let (text, truncated) = truncate_line(line, self.config.max_line_bytes);
226 self.push_entry(TailEntry::Line { text, truncated });
227 }
228
229 fn push_entry(&mut self, entry: TailEntry) {
230 self.bytes += entry.cost();
231 if matches!(entry, TailEntry::Line { .. }) {
232 self.lines += 1;
233 }
234 self.entries.push_back(entry);
235 self.evict_to_fit();
236 }
237
238 fn evict_to_fit(&mut self) {
239 while self.lines > self.config.max_lines
240 || (self.bytes > self.config.max_bytes && self.entries.len() > 1)
241 {
242 let Some(evicted) = self.entries.pop_front() else {
243 break;
244 };
245 self.bytes -= evicted.cost();
246 if matches!(evicted, TailEntry::Line { .. }) {
247 self.lines -= 1;
248 self.dropped_lines += 1;
249 }
250 }
251 }
252
253 pub fn snapshot(
257 &self,
258 max_lines: Option<usize>,
259 max_bytes: Option<usize>,
260 ) -> StderrTailSnapshot {
261 let line_limit = max_lines.unwrap_or(self.config.max_lines);
262 let byte_limit = max_bytes.unwrap_or(self.config.max_bytes);
263
264 let mut taken: Vec<TailEntry> = Vec::new();
265 let mut bytes = 0usize;
266 let mut lines = 0usize;
267 for entry in self.entries.iter().rev() {
270 match entry {
271 TailEntry::Line { .. } => {
272 if lines >= line_limit {
273 break;
274 }
275 let cost = entry.cost();
276 if lines > 0 && bytes + cost > byte_limit {
277 break;
278 }
279 bytes += cost;
280 lines += 1;
281 taken.push(entry.clone());
282 }
283 TailEntry::ProcessStart if lines > 0 => taken.push(entry.clone()),
284 TailEntry::ProcessStart => {}
285 }
286 }
287 taken.reverse();
288
289 let withheld = self
290 .entries
291 .iter()
292 .filter(|entry| matches!(entry, TailEntry::Line { .. }))
293 .count()
294 .saturating_sub(
295 taken
296 .iter()
297 .filter(|entry| matches!(entry, TailEntry::Line { .. }))
298 .count(),
299 );
300
301 StderrTailSnapshot {
302 capture: self.capture.clone(),
303 entries: taken,
304 dropped_lines: self.dropped_lines + withheld as u64,
309 }
310 }
311}
312
313const MAX_PENDING_LINE_BYTES: usize = 1024 * 1024;
322
323pub async fn pump_stderr<R>(source: R, ring: Arc<Mutex<StderrRing>>)
326where
327 R: AsyncReadExt + Unpin,
328{
329 pump_stderr_into(source, ring, &mut StderrSink).await
330}
331
332#[derive(Clone)]
339pub(crate) enum ChildOutputSink {
340 File {
341 sink: Arc<Mutex<cortexkit_log::LineSink>>,
342 path: Arc<PathBuf>,
343 failure_reported: Arc<AtomicBool>,
344 },
345 Stderr,
346}
347
348impl ChildOutputSink {
349 pub(crate) fn open(path: &Path, retention: cortexkit_log::Retention) -> io::Result<Self> {
350 Ok(Self::File {
351 sink: Arc::new(Mutex::new(cortexkit_log::LineSink::open(path, retention)?)),
352 path: Arc::new(path.to_path_buf()),
353 failure_reported: Arc::new(AtomicBool::new(false)),
354 })
355 }
356}
357
358pub trait OutputSink {
361 fn write_line(&mut self, line: &[u8]);
362}
363
364struct StderrSink;
365
366impl OutputSink for StderrSink {
367 fn write_line(&mut self, line: &[u8]) {
368 let stderr = std::io::stderr();
369 let mut handle = stderr.lock();
370 let _ = handle.write_all(line);
371 }
372}
373
374impl OutputSink for ChildOutputSink {
375 fn write_line(&mut self, line: &[u8]) {
376 match self {
377 Self::File {
378 sink,
379 path,
380 failure_reported,
381 } => {
382 let result = sink
383 .lock()
384 .unwrap_or_else(|poisoned| poisoned.into_inner())
385 .write_line(line);
386 if let Err(error) = result {
387 if !failure_reported.swap(true, Ordering::Relaxed) {
388 tracing::warn!(
389 path = %path.display(),
390 error = %error,
391 "child output capture write failed; later failures are suppressed"
392 );
393 }
394 }
395 }
396 Self::Stderr => StderrSink.write_line(line),
397 }
398 }
399}
400
401pub(crate) async fn pump_stderr_to<R>(
402 source: R,
403 ring: Arc<Mutex<StderrRing>>,
404 mut sink: ChildOutputSink,
405) where
406 R: AsyncReadExt + Unpin,
407{
408 pump_stderr_into(source, ring, &mut sink).await;
409}
410
411pub(crate) async fn pump_stdout_to<R>(source: R, mut sink: ChildOutputSink)
412where
413 R: AsyncReadExt + Unpin,
414{
415 pump_lines_into(source, None, &mut sink, "stdout").await;
416}
417
418async fn pump_stderr_into<R, S>(source: R, ring: Arc<Mutex<StderrRing>>, sink: &mut S)
419where
420 R: AsyncReadExt + Unpin,
421 S: OutputSink,
422{
423 pump_lines_into(source, Some(&ring), sink, "stderr").await;
424}
425
426async fn pump_lines_into<R, S>(
427 mut source: R,
428 ring: Option<&Arc<Mutex<StderrRing>>>,
429 sink: &mut S,
430 stream_name: &str,
431) where
432 R: AsyncReadExt + Unpin,
433 S: OutputSink,
434{
435 if let Some(ring) = ring {
436 lock_ring(ring).mark_captured();
437 }
438
439 let mut pending: Vec<u8> = Vec::new();
440 let mut scanned_upto = 0usize;
444 let mut cursor = 0usize;
447 let mut chunk = [0u8; 8192];
448 loop {
449 let read = match source.read(&mut chunk).await {
450 Ok(0) => break,
451 Ok(n) => n,
452 Err(error) => {
453 if let Some(ring) = ring {
454 lock_ring(ring).mark_incomplete(format!("{stream_name} read failed: {error}"));
455 } else {
456 tracing::warn!(stream = stream_name, error = %error, "child output capture read failed");
457 }
458 return;
459 }
460 };
461 pending.extend_from_slice(&chunk[..read]);
462
463 while let Some(relative) = find_newline(&pending[scanned_upto..]) {
464 let newline = scanned_upto + relative;
465 emit_line(ring, sink, &pending[cursor..newline], true);
466 cursor = newline + 1;
467 scanned_upto = cursor;
468 }
469 scanned_upto = pending.len();
470
471 if cursor > 0 {
472 pending.drain(..cursor);
473 scanned_upto -= cursor;
474 cursor = 0;
475 }
476
477 if pending.len() >= MAX_PENDING_LINE_BYTES {
478 let line = std::mem::take(&mut pending);
479 emit_line(ring, sink, &line, false);
480 scanned_upto = 0;
481 }
482 }
483
484 if !pending.is_empty() {
485 emit_line(ring, sink, &pending, false);
486 }
487}
488
489#[cfg(test)]
497thread_local! {
498 static SCANNED_BYTES: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
499}
500
501#[cfg(test)]
502fn take_scanned_bytes() -> usize {
503 SCANNED_BYTES.with(|scanned| scanned.replace(0))
504}
505
506fn find_newline(haystack: &[u8]) -> Option<usize> {
509 let found = memchr::memchr(b'\n', haystack);
510 #[cfg(test)]
511 SCANNED_BYTES.with(|scanned| {
512 scanned.set(scanned.get() + found.map(|index| index + 1).unwrap_or(haystack.len()));
513 });
514 found
515}
516
517fn emit_line<S: OutputSink>(
518 ring: Option<&Arc<Mutex<StderrRing>>>,
519 sink: &mut S,
520 raw: &[u8],
521 terminated: bool,
522) {
523 if let Some(ring) = ring {
524 lock_ring(ring).push_line(&String::from_utf8_lossy(raw));
525 }
526
527 if terminated {
530 let mut framed = Vec::with_capacity(raw.len() + 1);
531 framed.extend_from_slice(raw);
532 framed.push(b'\n');
533 sink.write_line(&framed);
534 } else {
535 sink.write_line(raw);
536 }
537}
538
539fn lock_ring(ring: &Arc<Mutex<StderrRing>>) -> std::sync::MutexGuard<'_, StderrRing> {
540 ring.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
541}
542
543fn truncate_line(line: &str, max_bytes: usize) -> (String, bool) {
549 if line.len() <= max_bytes {
550 return (line.to_string(), false);
551 }
552 let mut end = max_bytes;
553 while end > 0 && !line.is_char_boundary(end) {
554 end -= 1;
555 }
556 (line[..end].to_string(), true)
557}
558
559#[cfg(test)]
560mod tests {
561 use std::{
562 io,
563 pin::Pin,
564 task::{Context, Poll},
565 };
566
567 use super::*;
568 use tokio::io::{AsyncRead, ReadBuf};
569
570 fn ring(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> StderrRing {
571 StderrRing::new(StderrTailConfig::new(max_lines, max_bytes, max_line_bytes))
572 }
573
574 fn lines(snapshot: &StderrTailSnapshot) -> Vec<String> {
575 snapshot
576 .entries
577 .iter()
578 .filter_map(|entry| match entry {
579 TailEntry::Line { text, .. } => Some(text.clone()),
580 TailEntry::ProcessStart => None,
581 })
582 .collect()
583 }
584
585 #[test]
586 fn a_fresh_ring_reports_not_captured_rather_than_empty() {
587 let ring = ring(10, 1024, 128);
590 let snapshot = ring.snapshot(None, None);
591 assert!(matches!(snapshot.capture, CaptureState::NotCaptured { .. }));
592 assert!(snapshot.entries.is_empty());
593 }
594
595 #[test]
596 fn a_captured_module_that_printed_nothing_is_distinguishable_from_an_uncaptured_one() {
597 let mut captured = ring(10, 1024, 128);
598 captured.mark_captured();
599 let uncaptured = ring(10, 1024, 128);
600
601 let captured = captured.snapshot(None, None);
602 let uncaptured = uncaptured.snapshot(None, None);
603
604 assert!(captured.entries.is_empty());
607 assert!(uncaptured.entries.is_empty());
608 assert_eq!(captured.capture, CaptureState::Captured);
609 assert!(matches!(
610 uncaptured.capture,
611 CaptureState::NotCaptured { .. }
612 ));
613 }
614
615 #[test]
616 fn the_line_cap_evicts_oldest_first_and_counts_what_it_dropped() {
617 let mut ring = ring(3, 10_000, 128);
618 ring.mark_captured();
619 for i in 0..6 {
620 ring.push_line(&format!("line{i}"));
621 }
622 let snapshot = ring.snapshot(None, None);
623 assert_eq!(lines(&snapshot), vec!["line3", "line4", "line5"]);
624 assert_eq!(snapshot.dropped_lines, 3);
627 }
628
629 #[test]
630 fn the_byte_cap_binds_before_the_line_cap_when_lines_are_large() {
631 let mut ring = ring(100, 30, 128);
633 ring.mark_captured();
634 for i in 0..10 {
635 ring.push_line(&format!("{i}--------")); }
637 let snapshot = ring.snapshot(None, None);
638 assert!(
639 snapshot.entries.len() < 10,
640 "byte cap did not bind: {} entries retained",
641 snapshot.entries.len()
642 );
643 let retained: usize = lines(&snapshot).iter().map(String::len).sum();
644 assert!(
645 retained <= 30,
646 "retained {retained} bytes over a 30 byte cap"
647 );
648 assert!(snapshot.dropped_lines > 0);
649 }
650
651 #[test]
652 fn one_enormous_line_is_truncated_rather_than_evicting_the_tail() {
653 let mut ring = ring(10, 10_000, 64);
656 ring.mark_captured();
657 ring.push_line("context line that must survive");
658 ring.push_line(&"x".repeat(40_000));
659
660 let snapshot = ring.snapshot(None, None);
661 let kept = &snapshot.entries;
662 assert!(matches!(
663 &kept[0],
664 TailEntry::Line { text, truncated: false }
665 if text == "context line that must survive"
666 ));
667 let TailEntry::Line { text, truncated } = &kept[1] else {
668 panic!("expected a truncated line");
669 };
670 assert_eq!(text, &"x".repeat(64));
671 assert!(*truncated);
672 }
673
674 #[test]
675 fn truncation_is_visible_so_a_cut_line_is_not_mistaken_for_a_short_one() {
676 let mut ring = ring(10, 10_000, 16);
677 ring.mark_captured();
678 ring.push_line("0123456789abcdefghij");
679 ring.push_line("short");
680
681 let snapshot = ring.snapshot(None, None);
682 let TailEntry::Line { truncated, .. } = &snapshot.entries[0] else {
683 panic!("expected a line");
684 };
685 assert!(truncated);
686 let TailEntry::Line { truncated, .. } = &snapshot.entries[1] else {
687 panic!("expected a line");
688 };
689 assert!(!truncated, "a short line must not be reported as truncated");
690 }
691
692 #[test]
693 fn truncation_cuts_on_a_char_boundary_rather_than_splitting_utf8() {
694 let mut ring = ring(10, 10_000, 5);
697 ring.mark_captured();
698 ring.push_line("aa€€€€");
699 let snapshot = ring.snapshot(None, None);
700 let TailEntry::Line { text, truncated } = &snapshot.entries[0] else {
701 panic!("expected a line");
702 };
703 assert!(truncated);
704 assert!(text.starts_with("aa"));
705 }
706
707 #[test]
708 fn a_restart_boundary_keeps_generations_distinguishable() {
709 let mut ring = ring(10, 10_000, 128);
710 ring.mark_captured();
711 ring.push_line("before the crash");
712 ring.push_process_start();
713 ring.push_line("after the respawn");
714
715 let snapshot = ring.snapshot(None, None);
716 assert_eq!(
717 snapshot.entries,
718 vec![
719 TailEntry::Line {
720 text: "before the crash".to_string(),
721 truncated: false
722 },
723 TailEntry::ProcessStart,
724 TailEntry::Line {
725 text: "after the respawn".to_string(),
726 truncated: false
727 },
728 ]
729 );
730 }
731
732 #[test]
733 fn the_ring_survives_respawn_because_the_cause_is_written_before_the_exit() {
734 let mut ring = ring(10, 10_000, 128);
737 ring.mark_captured();
738 ring.push_line("Error: storage section missing");
739 ring.push_process_start();
740
741 let snapshot = ring.snapshot(None, None);
742 assert!(lines(&snapshot).contains(&"Error: storage section missing".to_string()));
743 }
744
745 #[test]
746 fn a_caller_limit_returns_the_newest_lines_not_the_oldest() {
747 let mut ring = ring(100, 100_000, 128);
748 ring.mark_captured();
749 for i in 0..10 {
750 ring.push_line(&format!("line{i}"));
751 }
752 let snapshot = ring.snapshot(Some(3), None);
753 assert_eq!(lines(&snapshot), vec!["line7", "line8", "line9"]);
754 }
755
756 #[test]
757 fn a_caller_line_limit_keeps_the_boundary_before_the_selected_line() {
758 let mut ring = ring(100, 100_000, 128);
759 ring.mark_captured();
760 ring.push_line("before restart");
761 ring.push_process_start();
762 ring.push_line("after restart");
763
764 let snapshot = ring.snapshot(Some(1), None);
765 assert_eq!(
766 snapshot.entries,
767 vec![
768 TailEntry::ProcessStart,
769 TailEntry::Line {
770 text: "after restart".to_string(),
771 truncated: false,
772 },
773 ]
774 );
775 }
776
777 #[test]
778 fn a_caller_line_limit_omits_a_trailing_boundary_after_the_selected_line() {
779 let mut ring = ring(100, 100_000, 128);
780 ring.mark_captured();
781 ring.push_line("before restart");
782 ring.push_process_start();
783
784 let snapshot = ring.snapshot(Some(1), None);
785 assert_eq!(
786 snapshot.entries,
787 vec![TailEntry::Line {
788 text: "before restart".to_string(),
789 truncated: false,
790 }]
791 );
792 }
793
794 #[test]
795 fn a_caller_limit_reports_what_it_withheld_rather_than_looking_complete() {
796 let mut ring = ring(100, 100_000, 128);
797 ring.mark_captured();
798 for i in 0..10 {
799 ring.push_line(&format!("line{i}"));
800 }
801 assert_eq!(ring.snapshot(Some(3), None).dropped_lines, 7);
804 assert_eq!(ring.snapshot(None, None).dropped_lines, 0);
805 }
806
807 #[test]
808 fn a_caller_limit_cannot_widen_the_rings_own_caps() {
809 let mut ring = ring(2, 10_000, 128);
810 ring.mark_captured();
811 for i in 0..5 {
812 ring.push_line(&format!("line{i}"));
813 }
814 let snapshot = ring.snapshot(Some(1000), Some(1_000_000));
815 assert_eq!(lines(&snapshot).len(), 2);
816 }
817
818 fn shared(max_lines: usize, max_bytes: usize, max_line_bytes: usize) -> Arc<Mutex<StderrRing>> {
819 Arc::new(Mutex::new(ring(max_lines, max_bytes, max_line_bytes)))
820 }
821
822 #[derive(Default)]
825 struct RecordingSink {
826 writes: Vec<Vec<u8>>,
827 }
828
829 impl OutputSink for RecordingSink {
830 fn write_line(&mut self, line: &[u8]) {
831 self.writes.push(line.to_vec());
832 }
833 }
834
835 struct ChunkedReader {
838 chunks: VecDeque<Vec<u8>>,
839 }
840
841 impl AsyncRead for ChunkedReader {
842 fn poll_read(
843 mut self: Pin<&mut Self>,
844 _cx: &mut Context<'_>,
845 buf: &mut ReadBuf<'_>,
846 ) -> Poll<io::Result<()>> {
847 match self.chunks.pop_front() {
848 None => Poll::Ready(Ok(())),
849 Some(chunk) => {
850 buf.put_slice(&chunk);
851 Poll::Ready(Ok(()))
852 }
853 }
854 }
855 }
856
857 struct FailingReader {
858 bytes: Vec<u8>,
859 emitted: bool,
860 }
861
862 impl AsyncRead for FailingReader {
863 fn poll_read(
864 mut self: Pin<&mut Self>,
865 _cx: &mut Context<'_>,
866 buf: &mut ReadBuf<'_>,
867 ) -> Poll<io::Result<()>> {
868 if self.emitted {
869 return Poll::Ready(Err(io::Error::other("reader failed")));
870 }
871 self.emitted = true;
872 buf.put_slice(&self.bytes);
873 Poll::Ready(Ok(()))
874 }
875 }
876
877 #[tokio::test]
878 async fn the_pump_splits_on_newlines_and_keeps_a_trailing_fragment() {
879 let ring = shared(10, 10_000, 128);
880 let source = std::io::Cursor::new(b"one\ntwo\nthree".to_vec());
883 let mut sink = RecordingSink::default();
884 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
885
886 let snapshot = lock_ring(&ring).snapshot(None, None);
887 assert_eq!(lines(&snapshot), vec!["one", "two", "three"]);
888 assert_eq!(snapshot.capture, CaptureState::Captured);
889 assert_eq!(
890 sink.writes,
891 vec![b"one\n".to_vec(), b"two\n".to_vec(), b"three".to_vec()]
892 );
893 }
894
895 #[tokio::test]
896 async fn a_read_failure_keeps_prior_lines_and_marks_the_capture_incomplete() {
897 let ring = shared(10, 10_000, 128);
898 let source = FailingReader {
899 bytes: b"crash cause\n".to_vec(),
900 emitted: false,
901 };
902 let mut sink = RecordingSink::default();
903 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
904
905 let snapshot = lock_ring(&ring).snapshot(None, None);
906 assert_eq!(lines(&snapshot), vec!["crash cause"]);
907 assert!(matches!(
908 snapshot.capture,
909 CaptureState::Incomplete { ref reason } if reason.contains("reader failed")
910 ));
911 assert_eq!(sink.writes, vec![b"crash cause\n".to_vec()]);
912 }
913
914 #[tokio::test]
915 async fn every_captured_line_is_also_forwarded() {
916 let ring = shared(10, 10_000, 128);
920 let source = std::io::Cursor::new(b"alpha\nbeta\n".to_vec());
921 let mut sink = RecordingSink::default();
922 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
923
924 assert_eq!(sink.writes, vec![b"alpha\n".to_vec(), b"beta\n".to_vec()]);
925 }
926
927 #[tokio::test]
928 async fn each_forwarded_line_is_exactly_one_write() {
929 let ring = shared(10, 10_000, 128);
934 let source = std::io::Cursor::new(b"first\nsecond\nthird\n".to_vec());
935 let mut sink = RecordingSink::default();
936 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
937
938 assert_eq!(sink.writes.len(), 3);
939 for write in &sink.writes {
940 assert_eq!(
941 write.iter().filter(|byte| **byte == b'\n').count(),
942 1,
943 "a write carried something other than exactly one complete line"
944 );
945 assert_eq!(*write.last().unwrap(), b'\n');
946 }
947 }
948
949 #[test]
950 fn the_first_process_start_is_not_recorded_because_it_divides_nothing() {
951 let mut ring = ring(10, 10_000, 128);
954 ring.push_process_start();
955 assert!(ring.snapshot(None, None).entries.is_empty());
956
957 ring.push_line("first process said this");
958 ring.push_process_start();
959 assert!(
960 matches!(ring.entries.back(), Some(TailEntry::ProcessStart)),
961 "a boundary with output before it must be recorded"
962 );
963 }
964
965 #[test]
966 fn a_process_start_is_recorded_when_only_dropped_lines_precede_it() {
967 let mut ring = ring(1, 10_000, 128);
971 ring.push_line("evicted");
972 ring.push_line("also evicted");
973 ring.entries.clear();
977 ring.lines = 0;
978 ring.bytes = 0;
979 ring.push_process_start();
980 assert!(matches!(
981 ring.entries.front(),
982 Some(TailEntry::ProcessStart)
983 ));
984 }
985
986 #[tokio::test]
987 async fn the_pump_marks_captured_even_when_the_module_writes_nothing() {
988 let ring = shared(10, 10_000, 128);
991 let source = std::io::Cursor::new(Vec::new());
992 let mut sink = RecordingSink::default();
993 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
994
995 let snapshot = lock_ring(&ring).snapshot(None, None);
996 assert!(snapshot.entries.is_empty());
997 assert_eq!(snapshot.capture, CaptureState::Captured);
998 assert!(sink.writes.is_empty());
999 }
1000
1001 #[tokio::test]
1002 async fn a_line_with_no_newline_cannot_grow_the_reader_without_bound() {
1003 let ring = shared(10, 10_000_000, 4 * 1024 * 1024);
1006 let source = std::io::Cursor::new(vec![b'x'; MAX_PENDING_LINE_BYTES + 4096]);
1007 let mut sink = RecordingSink::default();
1008 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1009
1010 let snapshot = lock_ring(&ring).snapshot(None, None);
1011 assert_eq!(
1012 lines(&snapshot).len(),
1013 2,
1014 "expected a forced flush at the ceiling plus the remainder"
1015 );
1016 assert_eq!(
1017 sink.writes,
1018 vec![vec![b'x'; MAX_PENDING_LINE_BYTES], vec![b'x'; 4096],],
1019 "forced flushes and EOF fragments must not invent delimiters"
1020 );
1021 }
1022
1023 #[tokio::test]
1024 async fn boundaries_truncation_and_framing_do_not_depend_on_chunk_splits() {
1025 let ring = shared(100, 100_000, 8);
1029 let source = ChunkedReader {
1030 chunks: vec![
1031 b"fir".to_vec(),
1032 b"st\nsec".to_vec(),
1033 b"ond\ncarry\r".to_vec(),
1034 b"\nover\n".to_vec(),
1035 b"12345678\n".to_vec(),
1036 b"1234567".to_vec(),
1037 b"89\n".to_vec(),
1038 b"tail".to_vec(),
1039 ]
1040 .into_iter()
1041 .collect(),
1042 };
1043 let mut sink = RecordingSink::default();
1044 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1045
1046 let snapshot = lock_ring(&ring).snapshot(None, None);
1047 assert_eq!(snapshot.capture, CaptureState::Captured);
1048 assert_eq!(
1049 snapshot.entries,
1050 vec![
1051 TailEntry::Line {
1052 text: "first".to_string(),
1053 truncated: false
1054 },
1055 TailEntry::Line {
1056 text: "second".to_string(),
1057 truncated: false
1058 },
1059 TailEntry::Line {
1061 text: "carry\r".to_string(),
1062 truncated: false
1063 },
1064 TailEntry::Line {
1065 text: "over".to_string(),
1066 truncated: false
1067 },
1068 TailEntry::Line {
1070 text: "12345678".to_string(),
1071 truncated: false
1072 },
1073 TailEntry::Line {
1075 text: "12345678".to_string(),
1076 truncated: true
1077 },
1078 TailEntry::Line {
1079 text: "tail".to_string(),
1080 truncated: false
1081 },
1082 ]
1083 );
1084 assert_eq!(
1085 sink.writes,
1086 vec![
1087 b"first\n".to_vec(),
1088 b"second\n".to_vec(),
1089 b"carry\r\n".to_vec(),
1090 b"over\n".to_vec(),
1091 b"12345678\n".to_vec(),
1092 b"123456789\n".to_vec(),
1093 b"tail".to_vec(),
1094 ]
1095 );
1096 }
1097
1098 #[tokio::test]
1099 async fn a_line_with_no_newline_is_not_rescanned_from_byte_zero_on_every_chunk() {
1100 let input = vec![b'x'; MAX_PENDING_LINE_BYTES + 4096];
1106 let ring = shared(10, 10_000_000, 4 * 1024 * 1024);
1107 let source = std::io::Cursor::new(input.clone());
1108 let mut sink = RecordingSink::default();
1109
1110 take_scanned_bytes();
1111 pump_stderr_into(source, Arc::clone(&ring), &mut sink).await;
1112 let scanned = take_scanned_bytes();
1113
1114 assert!(
1115 scanned <= 2 * input.len(),
1116 "newline searches examined {scanned} bytes for {} bytes of input; \
1117 each chunk must search only newly arrived bytes",
1118 input.len()
1119 );
1120 }
1121
1122 #[test]
1123 fn a_byte_limit_smaller_than_one_line_still_returns_that_line() {
1124 let mut ring = ring(10, 10_000, 128);
1127 ring.mark_captured();
1128 ring.push_line("a line considerably longer than the request limit");
1129 let snapshot = ring.snapshot(None, Some(4));
1130 assert_eq!(snapshot.entries.len(), 1);
1131 }
1132
1133 #[test]
1134 fn an_incoherent_config_clamps_the_line_cap_and_keeps_its_restart_boundary() {
1135 let config = StderrTailConfig::new(2, 10, 100);
1136 assert_eq!(config.max_line_bytes, config.max_bytes);
1137 let mut ring = StderrRing::new(config);
1138 ring.mark_captured();
1139 ring.push_line("old");
1140 ring.push_process_start();
1141 ring.push_line("new process line longer than the ring byte cap");
1142
1143 assert_eq!(
1144 ring.snapshot(None, None).entries,
1145 vec![
1146 TailEntry::ProcessStart,
1147 TailEntry::Line {
1148 text: "new proces".to_string(),
1149 truncated: true,
1150 },
1151 ]
1152 );
1153 }
1154}