use std::collections::VecDeque;
use std::io::{self, Write};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex, PoisonError};
use tracing_subscriber::fmt::MakeWriter;
const CAPACITY: usize = 1024;
#[derive(Debug, Default)]
struct CaptureInner {
lines: VecDeque<String>,
partial: String,
armed: bool,
}
#[derive(Debug, Clone, Default)]
pub struct LogCapture {
inner: Arc<Mutex<CaptureInner>>,
dropped: Arc<AtomicU64>,
}
impl LogCapture {
pub fn new() -> Self {
Self::default()
}
pub fn arm(&self) {
self.lock().armed = true;
}
pub fn disarm(&self) {
let mut inner = self.lock();
inner.armed = false;
inner.lines.clear();
inner.partial.clear();
}
pub fn is_armed(&self) -> bool {
self.lock().armed
}
pub fn drain(&self) -> Vec<String> {
std::mem::take(&mut self.lock().lines).into()
}
pub fn dropped(&self) -> u64 {
self.dropped.load(Ordering::Relaxed)
}
fn push(&self, buf: &[u8]) {
let text = String::from_utf8_lossy(buf);
let mut inner = self.lock();
for ch in text.chars() {
if ch != '\n' {
inner.partial.push(ch);
continue;
}
let mut line = std::mem::take(&mut inner.partial);
if line.ends_with('\r') {
line.pop();
}
while inner.lines.len() >= CAPACITY {
inner.lines.pop_front();
self.dropped.fetch_add(1, Ordering::Relaxed);
}
inner.lines.push_back(line);
}
}
fn lock(&self) -> std::sync::MutexGuard<'_, CaptureInner> {
self.inner.lock().unwrap_or_else(PoisonError::into_inner)
}
}
#[derive(Debug)]
pub enum CaptureWriter {
Stderr(io::Stderr),
Buffered(LogCapture),
}
impl CaptureWriter {
#[cfg(test)]
pub fn is_buffered(&self) -> bool {
matches!(self, Self::Buffered(_))
}
}
impl Write for CaptureWriter {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
match self {
Self::Stderr(w) => w.write(buf),
Self::Buffered(c) => {
c.push(buf);
Ok(buf.len())
}
}
}
fn flush(&mut self) -> io::Result<()> {
match self {
Self::Stderr(w) => w.flush(),
Self::Buffered(_) => Ok(()),
}
}
}
impl<'a> MakeWriter<'a> for LogCapture {
type Writer = CaptureWriter;
fn make_writer(&'a self) -> Self::Writer {
if self.is_armed() {
CaptureWriter::Buffered(self.clone())
} else {
CaptureWriter::Stderr(io::stderr())
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn disarmed_capture_writes_to_stderr() {
let capture = LogCapture::new();
assert!(!capture.is_armed());
assert!(
!capture.make_writer().is_buffered(),
"a disarmed capture must hand back the ordinary stderr writer"
);
}
#[test]
fn arming_diverts_and_disarming_restores() {
let capture = LogCapture::new();
capture.arm();
assert!(capture.is_armed());
assert!(capture.make_writer().is_buffered());
capture
.make_writer()
.write_all(b"buffered\n")
.expect("write");
capture.disarm();
assert!(!capture.is_armed());
assert!(!capture.make_writer().is_buffered());
assert!(
capture.drain().is_empty(),
"disarm discards what the pane already showed"
);
}
#[test]
fn armed_capture_buffers_whole_lines() {
let capture = LogCapture::new();
capture.arm();
let mut w = capture.make_writer();
w.write_all(b"2026-08-08T16:53:27Z WARN fetch failed\n")
.expect("write");
w.write_all(b"second line\n").expect("write");
assert_eq!(
capture.drain(),
vec![
"2026-08-08T16:53:27Z WARN fetch failed".to_string(),
"second line".to_string()
]
);
assert!(capture.drain().is_empty(), "drain empties the buffer");
}
#[test]
fn a_partial_line_is_held_until_its_newline_arrives() {
let capture = LogCapture::new();
capture.arm();
let mut w = capture.make_writer();
w.write_all(b"half a li").expect("write");
assert!(
capture.drain().is_empty(),
"an unterminated line must not render as a truncated one"
);
w.write_all(b"ne\n").expect("write");
assert_eq!(capture.drain(), vec!["half a line".to_string()]);
}
#[test]
fn overflow_drops_oldest_and_counts() {
let capture = LogCapture::new();
capture.arm();
let mut w = capture.make_writer();
for i in 0..CAPACITY + 5 {
w.write_all(format!("line {i}\n").as_bytes())
.expect("write");
}
let lines = capture.drain();
assert_eq!(lines.len(), CAPACITY);
assert_eq!(lines[0], format!("line {}", 5), "the oldest are evicted");
assert_eq!(capture.dropped(), 5);
}
#[test]
fn clones_share_one_buffer() {
let capture = LogCapture::new();
let clone = capture.clone();
clone.arm();
assert!(capture.is_armed(), "arming is visible through every handle");
capture
.make_writer()
.write_all(b"from the original\n")
.expect("write");
assert_eq!(clone.drain(), vec!["from the original".to_string()]);
}
}