use std::sync::atomic::{AtomicBool, Ordering};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
#[non_exhaustive]
pub enum StdioMode {
#[default]
Piped,
Inherit,
Null,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
#[non_exhaustive]
pub enum LineTerminator {
#[default]
Newline,
CarriageReturn,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum OverflowMode {
DropOldest,
DropNewest,
Error,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct OutputBufferPolicy {
pub max_lines: Option<usize>,
pub max_bytes: Option<usize>,
pub overflow: OverflowMode,
}
impl OutputBufferPolicy {
pub fn unbounded() -> Self {
Self {
max_lines: None,
max_bytes: None,
overflow: OverflowMode::DropOldest,
}
}
pub fn bounded(max_lines: usize) -> Self {
Self {
max_lines: Some(max_lines),
max_bytes: None,
overflow: OverflowMode::DropOldest,
}
}
pub fn fail_loud(max_lines: usize) -> Self {
Self {
max_lines: Some(max_lines),
max_bytes: None,
overflow: OverflowMode::Error,
}
}
#[must_use]
pub fn with_max_bytes(mut self, max_bytes: usize) -> Self {
self.max_bytes = Some(max_bytes);
self
}
#[must_use]
pub fn with_overflow(mut self, overflow: OverflowMode) -> Self {
self.overflow = overflow;
self
}
}
impl Default for OutputBufferPolicy {
fn default() -> Self {
Self::unbounded()
}
}
pub(crate) fn push_capped_bytes(
buf: &mut Vec<u8>,
chunk: &[u8],
cap: Option<usize>,
mode: OverflowMode,
overflowed: &AtomicBool,
truncated: &AtomicBool,
) {
let Some(cap) = cap else {
buf.extend_from_slice(chunk);
return;
};
match mode {
OverflowMode::Error => {
if buf.len() + chunk.len() > cap {
overflowed.store(true, Ordering::Relaxed);
let room = cap.saturating_sub(buf.len());
buf.extend_from_slice(&chunk[..room.min(chunk.len())]);
} else {
buf.extend_from_slice(chunk);
}
}
OverflowMode::DropNewest => {
let room = cap.saturating_sub(buf.len());
let take = room.min(chunk.len());
if take < chunk.len() {
truncated.store(true, Ordering::Relaxed);
}
buf.extend_from_slice(&chunk[..take]);
}
OverflowMode::DropOldest => {
buf.extend_from_slice(chunk);
if buf.len() > cap {
truncated.store(true, Ordering::Relaxed);
if buf.len() > cap.saturating_mul(2) {
let excess = buf.len() - cap;
buf.drain(..excess);
}
}
}
}
}
#[mutants::skip]
pub(crate) fn clamp_dropoldest_tail(buf: &mut Vec<u8>, cap: Option<usize>, mode: OverflowMode) {
if let (OverflowMode::DropOldest, Some(cap)) = (mode, cap)
&& buf.len() > cap
{
let excess = buf.len() - cap;
buf.drain(..excess);
}
}
#[cfg(test)]
mod tests {
use super::*;
fn flags() -> (AtomicBool, AtomicBool) {
(AtomicBool::new(false), AtomicBool::new(false))
}
#[test]
fn error_mode_under_cap_does_not_overflow() {
let (overflowed, truncated) = flags();
let mut buf = Vec::new();
push_capped_bytes(
&mut buf,
b"abcd",
Some(5),
OverflowMode::Error,
&overflowed,
&truncated,
);
assert_eq!(buf, b"abcd", "cap-1 input must fit whole");
assert!(!overflowed.load(Ordering::Relaxed));
assert!(!truncated.load(Ordering::Relaxed));
}
#[test]
fn error_mode_exactly_at_cap_does_not_overflow() {
let (overflowed, truncated) = flags();
let mut buf = Vec::new();
push_capped_bytes(
&mut buf,
b"abcde",
Some(5),
OverflowMode::Error,
&overflowed,
&truncated,
);
assert_eq!(buf, b"abcde");
assert!(
!overflowed.load(Ordering::Relaxed),
"buf.len() + chunk.len() == cap must not overflow"
);
assert!(!truncated.load(Ordering::Relaxed));
}
#[test]
fn error_mode_one_over_cap_overflows_and_writes_partial_room() {
let (overflowed, truncated) = flags();
let mut buf = Vec::new();
push_capped_bytes(
&mut buf,
b"abcdef",
Some(5),
OverflowMode::Error,
&overflowed,
&truncated,
);
assert_eq!(
buf, b"abcde",
"only room = cap - buf.len() bytes are retained"
);
assert!(overflowed.load(Ordering::Relaxed));
}
#[test]
fn error_mode_partial_room_when_buffer_already_holds_some_bytes() {
let (overflowed, truncated) = flags();
let mut buf = b"ab".to_vec();
push_capped_bytes(
&mut buf,
b"cdefgh",
Some(5),
OverflowMode::Error,
&overflowed,
&truncated,
);
assert_eq!(buf, b"abcde");
assert!(overflowed.load(Ordering::Relaxed));
}
#[test]
fn error_mode_zero_room_appends_nothing_further() {
let (overflowed, truncated) = flags();
let mut buf = b"abcde".to_vec();
push_capped_bytes(
&mut buf,
b"f",
Some(5),
OverflowMode::Error,
&overflowed,
&truncated,
);
assert_eq!(buf, b"abcde", "room is 0 once already at cap");
assert!(overflowed.load(Ordering::Relaxed));
}
#[test]
fn drop_newest_under_cap_does_not_truncate() {
let (overflowed, truncated) = flags();
let mut buf = Vec::new();
push_capped_bytes(
&mut buf,
b"abcd",
Some(5),
OverflowMode::DropNewest,
&overflowed,
&truncated,
);
assert_eq!(buf, b"abcd");
assert!(!truncated.load(Ordering::Relaxed));
assert!(!overflowed.load(Ordering::Relaxed));
}
#[test]
fn drop_newest_exact_fit_does_not_truncate() {
let (overflowed, truncated) = flags();
let mut buf = Vec::new();
push_capped_bytes(
&mut buf,
b"abcde",
Some(5),
OverflowMode::DropNewest,
&overflowed,
&truncated,
);
assert_eq!(buf, b"abcde");
assert!(
!truncated.load(Ordering::Relaxed),
"take == chunk.len() must not truncate"
);
}
#[test]
fn drop_newest_one_over_cap_truncates_and_keeps_head() {
let (overflowed, truncated) = flags();
let mut buf = Vec::new();
push_capped_bytes(
&mut buf,
b"abcdef",
Some(5),
OverflowMode::DropNewest,
&overflowed,
&truncated,
);
assert_eq!(buf, b"abcde");
assert!(
truncated.load(Ordering::Relaxed),
"take < chunk.len() must truncate"
);
}
#[test]
fn drop_newest_zero_room_still_truncates() {
let (overflowed, truncated) = flags();
let mut buf = b"abcde".to_vec();
push_capped_bytes(
&mut buf,
b"f",
Some(5),
OverflowMode::DropNewest,
&overflowed,
&truncated,
);
assert_eq!(buf, b"abcde");
assert!(truncated.load(Ordering::Relaxed));
}
#[test]
fn drop_oldest_under_cap_does_not_truncate() {
let (overflowed, truncated) = flags();
let mut buf = Vec::new();
push_capped_bytes(
&mut buf,
b"abcd",
Some(5),
OverflowMode::DropOldest,
&overflowed,
&truncated,
);
assert_eq!(buf, b"abcd");
assert!(!truncated.load(Ordering::Relaxed));
}
#[test]
fn drop_oldest_exactly_at_cap_does_not_truncate() {
let (overflowed, truncated) = flags();
let mut buf = Vec::new();
push_capped_bytes(
&mut buf,
b"abcde",
Some(5),
OverflowMode::DropOldest,
&overflowed,
&truncated,
);
assert_eq!(buf, b"abcde");
assert!(
!truncated.load(Ordering::Relaxed),
"buf.len() == cap must not truncate"
);
}
#[test]
fn drop_oldest_one_over_cap_truncates_without_compacting() {
let (overflowed, truncated) = flags();
let mut buf = Vec::new();
push_capped_bytes(
&mut buf,
b"abcdef",
Some(5),
OverflowMode::DropOldest,
&overflowed,
&truncated,
);
assert_eq!(
buf, b"abcdef",
"below the 2*cap compaction threshold, the tail is left as-is"
);
assert!(truncated.load(Ordering::Relaxed));
}
#[test]
fn drop_oldest_exactly_at_double_cap_does_not_compact() {
let (overflowed, truncated) = flags();
let mut buf = Vec::new();
push_capped_bytes(
&mut buf,
b"0123456789",
Some(5),
OverflowMode::DropOldest,
&overflowed,
&truncated,
);
assert_eq!(
buf, b"0123456789",
"buf.len() == 2*cap must not trigger compaction"
);
assert!(truncated.load(Ordering::Relaxed));
}
#[test]
fn drop_oldest_one_over_double_cap_compacts_to_exact_excess() {
let (overflowed, truncated) = flags();
let mut buf = Vec::new();
push_capped_bytes(
&mut buf,
b"abcdefghijk",
Some(5),
OverflowMode::DropOldest,
&overflowed,
&truncated,
);
assert_eq!(buf, b"ghijk");
assert!(truncated.load(Ordering::Relaxed));
}
#[test]
fn clamp_dropoldest_tail_under_cap_is_unchanged() {
let mut buf = b"abcd".to_vec();
clamp_dropoldest_tail(&mut buf, Some(5), OverflowMode::DropOldest);
assert_eq!(buf, b"abcd");
}
#[test]
fn clamp_dropoldest_tail_exactly_at_cap_is_unchanged() {
let mut buf = b"abcde".to_vec();
clamp_dropoldest_tail(&mut buf, Some(5), OverflowMode::DropOldest);
assert_eq!(buf, b"abcde", "buf.len() == cap must not be clamped");
}
#[test]
fn clamp_dropoldest_tail_one_over_cap_drops_exact_excess() {
let mut buf = b"abcdef".to_vec();
clamp_dropoldest_tail(&mut buf, Some(5), OverflowMode::DropOldest);
assert_eq!(
buf, b"bcdef",
"excess = 6 - 5 = 1 byte dropped from the front"
);
}
#[test]
fn clamp_dropoldest_tail_at_double_cap_drops_exact_excess() {
let mut buf = b"0123456789".to_vec();
clamp_dropoldest_tail(&mut buf, Some(5), OverflowMode::DropOldest);
assert_eq!(buf, b"56789");
}
#[test]
fn clamp_dropoldest_tail_well_over_cap_drops_exact_excess() {
let mut buf = b"abcdefghijk".to_vec();
clamp_dropoldest_tail(&mut buf, Some(5), OverflowMode::DropOldest);
assert_eq!(buf, b"ghijk");
}
#[test]
fn clamp_dropoldest_tail_ignores_other_modes() {
let mut buf = b"abcdefghijk".to_vec();
clamp_dropoldest_tail(&mut buf, Some(5), OverflowMode::DropNewest);
assert_eq!(buf, b"abcdefghijk", "only DropOldest is clamped");
}
#[test]
fn clamp_dropoldest_tail_no_cap_is_a_no_op() {
let mut buf = b"abcdefghijk".to_vec();
clamp_dropoldest_tail(&mut buf, None, OverflowMode::DropOldest);
assert_eq!(buf, b"abcdefghijk");
}
mod proptests {
use super::*;
use crate::pump::SharedLines;
use proptest::prelude::*;
#[derive(Default)]
struct ExpectedLines {
lines: Vec<String>,
count: usize,
seen_bytes: usize,
dropped: usize,
overflowed: bool,
}
impl ExpectedLines {
fn push(&mut self, line: String, policy: OutputBufferPolicy) {
self.count += 1;
self.seen_bytes += line.len();
match policy.overflow {
OverflowMode::Error => {
let over = match (policy.max_lines, policy.max_bytes) {
(None, None) => true,
(max_lines, max_bytes) => {
max_lines.is_some_and(|cap| self.count > cap)
|| max_bytes.is_some_and(|cap| {
self.seen_bytes > cap || self.count > cap
})
}
};
if over {
self.overflowed = true;
self.dropped += 1;
} else {
self.lines.push(line);
}
}
OverflowMode::DropOldest => {
self.lines.push(line);
let mut dropped = false;
while policy.max_lines.is_some_and(|cap| self.lines.len() > cap)
|| policy.max_bytes.is_some_and(|cap| {
self.lines.iter().map(String::len).sum::<usize>() > cap
|| self.lines.len() > cap
})
{
self.lines.remove(0);
dropped = true;
}
if dropped {
self.dropped += 1;
}
}
OverflowMode::DropNewest => {
let bytes: usize = self.lines.iter().map(String::len).sum();
let fits = policy.max_lines.is_none_or(|cap| self.lines.len() < cap)
&& policy.max_bytes.is_none_or(|cap| {
bytes + line.len() <= cap && self.lines.len() < cap
});
if fits {
self.lines.push(line);
} else {
self.dropped += 1;
}
}
}
}
}
fn arb_line() -> impl Strategy<Value = String> {
prop::collection::vec(any::<char>(), 0..16)
.prop_map(|chars| chars.into_iter().collect())
}
fn arb_cap(limit: usize) -> impl Strategy<Value = Option<usize>> {
prop_oneof![Just(None), (0usize..=limit).prop_map(Some)]
}
fn arb_overflow_mode() -> impl Strategy<Value = OverflowMode> {
prop_oneof![
Just(OverflowMode::DropOldest),
Just(OverflowMode::DropNewest),
Just(OverflowMode::Error),
]
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(100))]
#[test]
fn line_buffer_policy_matches_its_model(
input in prop::collection::vec(arb_line(), 0..24),
max_lines in arb_cap(12), max_bytes in arb_cap(96), overflow in arb_overflow_mode(),
) {
let policy = OutputBufferPolicy { max_lines, max_bytes, overflow };
let sink = SharedLines::new(&policy);
let mut expected = ExpectedLines::default();
for line in input {
expected.push(line.clone(), policy);
sink.push(line);
prop_assert_eq!(sink.count(), expected.count);
prop_assert_eq!(sink.seen_bytes(), expected.seen_bytes);
prop_assert_eq!(sink.dropped(), expected.dropped);
prop_assert_eq!(sink.overflowed(), expected.overflowed);
}
let retained = sink.drain();
let retained_bytes: usize = retained.iter().map(String::len).sum();
if let Some(cap) = policy.max_lines { prop_assert!(retained.len() <= cap); }
if let Some(cap) = policy.max_bytes {
prop_assert!(retained_bytes <= cap);
prop_assert!(retained.len() <= cap);
}
prop_assert_eq!(retained, expected.lines);
prop_assert_eq!(sink.dropped() > 0, expected.dropped > 0);
}
#[test]
fn raw_byte_buffer_policy_matches_its_model(
chunks in prop::collection::vec(prop::collection::vec(any::<u8>(), 0..80), 0..24),
max_lines in arb_cap(12), max_bytes in arb_cap(96), overflow in arb_overflow_mode(),
) {
let policy = OutputBufferPolicy { max_lines, max_bytes, overflow };
let input: Vec<u8> = chunks.iter().flatten().copied().collect();
let (overflowed, truncated) = (AtomicBool::new(false), AtomicBool::new(false));
let mut actual = Vec::new();
for chunk in &chunks {
let was_overflowed = overflowed.load(Ordering::Relaxed);
let was_truncated = truncated.load(Ordering::Relaxed);
push_capped_bytes(&mut actual, chunk, policy.max_bytes, policy.overflow, &overflowed, &truncated);
prop_assert!(!was_overflowed || overflowed.load(Ordering::Relaxed));
prop_assert!(!was_truncated || truncated.load(Ordering::Relaxed));
}
clamp_dropoldest_tail(&mut actual, policy.max_bytes, policy.overflow);
let (expected, expected_overflowed, expected_truncated) = match policy.max_bytes {
None => (input, false, false),
Some(cap) => match policy.overflow {
OverflowMode::DropOldest => {
let start = input.len().saturating_sub(cap);
(input[start..].to_vec(), false, input.len() > cap)
}
OverflowMode::DropNewest => (input[..input.len().min(cap)].to_vec(), false, input.len() > cap),
OverflowMode::Error => (input[..input.len().min(cap)].to_vec(), input.len() > cap, false),
},
};
if let Some(cap) = policy.max_bytes { prop_assert!(actual.len() <= cap); }
prop_assert_eq!(actual, expected);
prop_assert_eq!(overflowed.load(Ordering::Relaxed), expected_overflowed);
prop_assert_eq!(truncated.load(Ordering::Relaxed), expected_truncated);
}
}
}
}