use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use crate::core::retain::{LineCap, Retention};
fn cap(policy: Retention) -> LineCap {
LineCap::new(policy.max_lines, policy.keep)
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum Stream {
Stdout,
Stderr,
}
#[derive(Debug, Eq, PartialEq)]
pub struct Emission {
pub stdout: Vec<Vec<u8>>,
pub stderr: Vec<Vec<u8>>,
pub dropped: usize,
pub at: jiff::Timestamp,
}
struct State {
policy: (Retention, Retention),
stdout: LineCap,
stderr: LineCap,
pending: Option<Emission>,
}
impl State {
fn shape(&self) -> (usize, usize, usize, usize) {
(
self.stdout.retained_len(),
self.stdout.dropped_count(),
self.stderr.retained_len(),
self.stderr.dropped_count(),
)
}
}
#[derive(Clone)]
pub struct Emissions(Arc<Mutex<State>>);
impl Emissions {
pub fn new(stdout: Retention, stderr: Retention) -> Emissions {
Emissions(Arc::new(Mutex::new(State {
policy: (stdout, stderr),
stdout: cap(stdout),
stderr: cap(stderr),
pending: None,
})))
}
pub fn feed(&self, stream: Stream, chunk: &[u8], at: jiff::Timestamp) -> bool {
let mut state = self.lock();
let before = state.shape();
match stream {
Stream::Stdout => state.stdout.feed(chunk),
Stream::Stderr => state.stderr.feed(chunk),
}
if state.shape() == before {
return false;
}
let (stdout, out_dropped) = state.stdout.snapshot();
let (stderr, err_dropped) = state.stderr.snapshot();
let was_empty = state.pending.is_none();
state.pending = Some(Emission {
stdout,
stderr,
dropped: out_dropped + err_dropped,
at,
});
was_empty
}
pub fn take(&self) -> Option<Emission> {
self.lock().pending.take()
}
pub fn finish(&self) -> (Vec<Vec<u8>>, Vec<Vec<u8>>, usize) {
let mut state = self.lock();
let (out_policy, err_policy) = state.policy;
let stdout = std::mem::replace(&mut state.stdout, cap(out_policy));
let stderr = std::mem::replace(&mut state.stderr, cap(err_policy));
let (out, out_dropped) = stdout.finish();
let (err, err_dropped) = stderr.finish();
(out, err, out_dropped + err_dropped)
}
fn lock(&self) -> MutexGuard<'_, State> {
self.0.lock().unwrap_or_else(PoisonError::into_inner)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::core::retain::Keep;
fn ample() -> Retention {
Retention {
max_lines: 100,
keep: Keep::Bottom,
}
}
fn bound(max_lines: usize) -> Retention {
Retention {
max_lines,
keep: Keep::Bottom,
}
}
fn at() -> jiff::Timestamp {
jiff::Timestamp::now()
}
#[test]
fn the_first_visible_line_transitions_the_outbox_and_says_so() {
let e = Emissions::new(ample(), ample());
assert!(e.feed(Stream::Stdout, b"one\n", at()), "empty -> occupied");
let body = e.take().expect("a body");
assert_eq!(body.stdout, vec![b"one\n".to_vec()]);
assert!(body.stderr.is_empty());
}
#[test]
fn feeding_while_a_body_is_unread_updates_it_and_sends_no_second_wake() {
let e = Emissions::new(ample(), ample());
assert!(e.feed(Stream::Stdout, b"one\n", at()));
for i in 0..500 {
assert!(
!e.feed(Stream::Stdout, format!("line-{i}\n").as_bytes(), at()),
"a wake was sent while a body was still pending"
);
}
let body = e.take().expect("a body");
assert_eq!(body.stdout.last().unwrap(), b"line-499\n");
assert!(e.take().is_none(), "one slot, never a backlog");
}
#[test]
fn taking_re_arms_the_wake() {
let e = Emissions::new(ample(), ample());
assert!(e.feed(Stream::Stdout, b"a\n", at()));
e.take();
assert!(e.feed(Stream::Stdout, b"b\n", at()), "a taken slot re-arms");
}
#[test]
fn a_partial_line_moves_nothing_and_wakes_nobody() {
let e = Emissions::new(ample(), ample());
assert!(!e.feed(Stream::Stdout, b"no-terminator", at()));
assert!(e.take().is_none());
assert!(
e.feed(Stream::Stdout, b"-yet\n", at()),
"the terminator moves it"
);
}
#[test]
fn the_two_streams_stay_apart_and_a_published_body_carries_both() {
let e = Emissions::new(ample(), ample());
e.feed(Stream::Stdout, b"out-1\n", at());
e.feed(Stream::Stderr, b"err-1\n", at());
let body = e.take().expect("a body");
assert_eq!(body.stdout, vec![b"out-1\n".to_vec()]);
assert_eq!(body.stderr, vec![b"err-1\n".to_vec()]);
}
#[test]
fn a_stderr_feed_publishes_the_current_stdout_too() {
let e = Emissions::new(ample(), ample());
e.feed(Stream::Stdout, b"out-1\n", at());
e.take();
e.feed(Stream::Stderr, b"err-1\n", at());
let body = e.take().expect("a body");
assert_eq!(
body.stdout,
vec![b"out-1\n".to_vec()],
"stdout must not vanish"
);
assert_eq!(body.stderr, vec![b"err-1\n".to_vec()]);
}
#[test]
fn each_stream_keeps_its_own_bound_and_the_drop_count_is_the_sum() {
let e = Emissions::new(bound(1), bound(1));
e.feed(Stream::Stdout, b"o1\no2\no3\n", at());
e.feed(Stream::Stderr, b"e1\ne2\n", at());
let body = e.take().expect("a body");
assert_eq!(body.stdout, vec![b"o3\n".to_vec()]);
assert_eq!(body.stderr, vec![b"e2\n".to_vec()]);
assert_eq!(body.dropped, 3, "2 dropped from stdout + 1 from stderr");
}
#[test]
fn eviction_alone_counts_as_movement() {
let e = Emissions::new(bound(1), bound(1));
assert!(e.feed(Stream::Stdout, b"first\n", at()));
e.take();
assert!(
e.feed(Stream::Stdout, b"second\n", at()),
"eviction is movement"
);
}
#[test]
fn finishing_flushes_the_line_the_child_never_terminated() {
let e = Emissions::new(ample(), ample());
e.feed(Stream::Stdout, b"whole\npartial", at());
let body = e.take().expect("a body");
assert_eq!(body.stdout, vec![b"whole\n".to_vec()], "not while it lives");
let (stdout, _, _) = e.finish();
assert_eq!(
stdout,
vec![b"whole\n".to_vec(), b"partial".to_vec()],
"but yes once it is gone, and verbatim — no terminator invented"
);
}
#[test]
fn finishing_reports_what_both_bounds_dropped() {
let e = Emissions::new(bound(1), bound(1));
e.feed(Stream::Stdout, b"o1\no2\no3\n", at());
e.feed(Stream::Stderr, b"e1\ne2\n", at());
let (stdout, stderr, dropped) = e.finish();
assert_eq!(stdout, vec![b"o3\n".to_vec()]);
assert_eq!(stderr, vec![b"e2\n".to_vec()]);
assert_eq!(dropped, 3, "2 dropped from stdout + 1 from stderr");
}
#[test]
fn finishing_re_arms_the_caps_for_the_next_child() {
let e = Emissions::new(bound(10), bound(10));
e.feed(Stream::Stdout, b"first-child\n", at());
assert_eq!(e.finish().0, vec![b"first-child\n".to_vec()]);
e.take();
assert!(
e.feed(Stream::Stdout, b"second-child\n", at()),
"a re-armed cap still wakes the loop"
);
let body = e.take().expect("a body");
assert_eq!(
body.stdout,
vec![b"second-child\n".to_vec()],
"the next child starts clean, neither appending nor bounded to nothing"
);
assert_eq!(body.dropped, 0, "and its drop count starts at zero");
}
#[test]
fn a_clone_shares_the_state() {
let a = Emissions::new(ample(), ample());
let b = a.clone();
b.feed(Stream::Stdout, b"through\n", at());
assert!(a.take().is_some(), "a clone must share, not copy");
}
#[test]
fn a_panicking_feeder_does_not_poison_the_state() {
let e = Emissions::new(ample(), ample());
let c = e.clone();
let _ = std::thread::spawn(move || {
c.feed(Stream::Stdout, b"before the panic\n", at());
panic!("worker died");
})
.join();
assert!(e.take().is_some(), "a panic must not wedge the state");
}
}