use std::collections::VecDeque;
use std::io::Read;
const MAX_LINE_BYTES: usize = 64 * 1024;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum Keep {
Top,
Bottom,
}
pub struct LineCap {
max_lines: usize,
keep: Keep,
retained: VecDeque<Vec<u8>>,
dropped: usize,
partial: Vec<u8>,
overlong: bool,
}
impl LineCap {
pub fn new(max_lines: usize, keep: Keep) -> LineCap {
LineCap {
max_lines,
keep,
retained: VecDeque::new(),
dropped: 0,
partial: Vec::new(),
overlong: false,
}
}
pub fn feed(&mut self, chunk: &[u8]) {
let mut rest = chunk;
while let Some(nl) = rest.iter().position(|&b| b == b'\n') {
self.extend_line(&rest[..nl]);
self.end_line();
rest = &rest[nl + 1..];
}
self.extend_line(rest);
}
pub fn finish(mut self) -> (Vec<Vec<u8>>, usize) {
if !self.partial.is_empty() {
let last = std::mem::take(&mut self.partial);
self.accept(&last);
}
(self.retained.into(), self.dropped)
}
pub fn snapshot(&self) -> (Vec<Vec<u8>>, usize) {
(self.retained.iter().cloned().collect(), self.dropped)
}
pub fn retained_len(&self) -> usize {
self.retained.len()
}
pub fn dropped_count(&self) -> usize {
self.dropped
}
fn extend_line(&mut self, bytes: &[u8]) {
if self.overlong {
return;
}
if self.partial.len() + bytes.len() > MAX_LINE_BYTES {
self.overlong = true;
self.dropped += 1;
self.partial = Vec::new();
return;
}
self.partial.extend_from_slice(bytes);
}
fn end_line(&mut self) {
if self.overlong {
self.overlong = false;
return;
}
let mut line = std::mem::take(&mut self.partial);
if line.last() == Some(&b'\r') {
line.pop();
}
line.push(b'\n');
self.accept(&line);
line.clear();
self.partial = line;
}
fn accept(&mut self, line: &[u8]) {
match self.keep {
Keep::Bottom => {
self.retained.push_back(line.to_vec());
while self.retained.len() > self.max_lines {
self.retained.pop_front();
self.dropped += 1;
}
}
Keep::Top => {
if self.retained.len() < self.max_lines {
self.retained.push_back(line.to_vec());
} else {
self.dropped += 1;
}
}
}
}
}
#[derive(Clone, Copy)]
pub struct Retention {
pub max_lines: usize,
pub keep: Keep,
}
pub fn read_all<R: Read>(pipe: Option<R>, retention: Retention) -> (Vec<Vec<u8>>, usize) {
let mut cap = LineCap::new(retention.max_lines, retention.keep);
if let Some(mut pipe) = pipe {
let mut buf = [0u8; 8 * 1024];
loop {
match pipe.read(&mut buf) {
Ok(0) => break,
Ok(n) => cap.feed(&buf[..n]),
Err(err) if err.kind() == std::io::ErrorKind::Interrupted => continue,
Err(_) => break,
}
}
}
cap.finish()
}
pub fn compact_count(n: usize) -> String {
const K: usize = 1_000;
const M: usize = 1_000_000;
match n {
n if n < K => n.to_string(),
n if n < M => format!("{}.{}k", n / K, (n % K) / 100),
n => format!("{}.{}M", n / M, (n % M) / 100_000),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn under_the_bound_nothing_is_dropped_and_the_lines_are_verbatim() {
let mut cap = LineCap::new(10, Keep::Bottom);
cap.feed(b"a\nb\nc\n");
assert_eq!(
cap.finish(),
(vec![b"a\n".to_vec(), b"b\n".to_vec(), b"c\n".to_vec()], 0)
);
}
#[test]
fn keep_bottom_retains_the_newest_and_counts_the_rest() {
let mut cap = LineCap::new(2, Keep::Bottom);
cap.feed(b"1\n2\n3\n4\n");
assert_eq!(cap.finish(), (vec![b"3\n".to_vec(), b"4\n".to_vec()], 2));
}
#[test]
fn keep_top_retains_the_oldest_and_counts_the_rest() {
let mut cap = LineCap::new(2, Keep::Top);
cap.feed(b"1\n2\n3\n4\n");
assert_eq!(cap.finish(), (vec![b"1\n".to_vec(), b"2\n".to_vec()], 2));
}
#[test]
fn a_trailing_line_without_a_newline_is_still_a_line() {
let mut cap = LineCap::new(10, Keep::Bottom);
cap.feed(b"a\nb");
assert_eq!(cap.finish(), (vec![b"a\n".to_vec(), b"b".to_vec()], 0));
}
#[test]
fn a_crlf_terminator_is_recognised_whole() {
let mut cap = LineCap::new(10, Keep::Bottom);
cap.feed(b"a\r\nb\r\n");
assert_eq!(cap.finish(), (vec![b"a\n".to_vec(), b"b\n".to_vec()], 0));
}
#[test]
fn a_crlf_split_across_two_reads_is_still_one_terminator() {
let mut cap = LineCap::new(10, Keep::Bottom);
cap.feed(b"a\r");
cap.feed(b"\nb\r\n");
assert_eq!(cap.finish(), (vec![b"a\n".to_vec(), b"b\n".to_vec()], 0));
}
#[test]
fn a_bare_carriage_return_is_content_and_survives() {
let mut cap = LineCap::new(10, Keep::Bottom);
cap.feed(b"10%\r50%\r100%\n");
assert_eq!(cap.finish(), (vec![b"10%\r50%\r100%\n".to_vec()], 0));
}
#[test]
fn an_unterminated_line_keeps_its_trailing_carriage_return() {
let mut cap = LineCap::new(10, Keep::Bottom);
cap.feed(b"a\r\nb\r");
assert_eq!(cap.finish(), (vec![b"a\n".to_vec(), b"b\r".to_vec()], 0));
}
#[test]
fn invalid_utf8_survives_the_accumulator_untouched() {
let mut cap = LineCap::new(10, Keep::Bottom);
cap.feed(b"\xff\n");
assert_eq!(cap.finish(), (vec![b"\xff\n".to_vec()], 0));
}
#[test]
fn the_retained_set_never_exceeds_the_bound_however_much_is_fed() {
let mut cap = LineCap::new(4, Keep::Bottom);
for _ in 0..10_000 {
cap.feed(b"x\n");
}
let (lines, dropped) = cap.finish();
assert_eq!(lines.len(), 4);
assert_eq!(dropped, 9_996);
}
#[test]
fn a_line_split_across_feeds_is_one_line() {
let mut cap = LineCap::new(10, Keep::Bottom);
cap.feed(b"hel");
cap.feed(b"lo\n");
assert_eq!(cap.finish(), (vec![b"hello\n".to_vec()], 0));
}
#[test]
fn a_split_line_still_counts_once_against_the_bound() {
let mut cap = LineCap::new(2, Keep::Bottom);
cap.feed(b"1\n2\n3");
cap.feed(b"\n");
assert_eq!(cap.finish(), (vec![b"2\n".to_vec(), b"3\n".to_vec()], 1));
}
#[test]
fn a_multibyte_char_split_across_feeds_survives() {
let mut cap = LineCap::new(10, Keep::Bottom);
let s = "héllo\n".as_bytes();
let cut = 2; cap.feed(&s[..cut]);
cap.feed(&s[cut..]);
assert_eq!(cap.finish(), (vec!["héllo\n".as_bytes().to_vec()], 0));
}
#[test]
fn many_tiny_feeds_bound_the_partial_buffer_too() {
let mut cap = LineCap::new(4, Keep::Bottom);
for _ in 0..50_000 {
cap.feed(b"x");
}
let (lines, _) = cap.finish();
assert_eq!(lines.len(), 1);
}
#[test]
fn one_enormous_line_does_not_grow_without_limit() {
let mut cap = LineCap::new(10, Keep::Bottom);
for _ in 0..(MAX_LINE_BYTES / 8 + 100) {
cap.feed(b"xxxxxxxx");
}
let (lines, dropped) = cap.finish();
assert!(
lines.iter().all(|l| l.len() <= MAX_LINE_BYTES + 1),
"a line grew past the backstop"
);
assert!(dropped > 0, "an over-long line must be counted as dropped");
}
#[test]
fn an_over_long_line_counts_once_however_many_reads_it_spans() {
let mut cap = LineCap::new(10, Keep::Bottom);
for _ in 0..40 {
cap.feed(&vec![b'x'; MAX_LINE_BYTES / 2]);
}
cap.feed(b"\n");
let (lines, dropped) = cap.finish();
assert!(lines.is_empty());
assert_eq!(dropped, 1);
}
#[test]
fn a_line_at_exactly_the_backstop_is_kept_whole() {
let mut cap = LineCap::new(10, Keep::Bottom);
cap.feed(&vec![b'x'; MAX_LINE_BYTES]);
cap.feed(b"\n");
let (lines, dropped) = cap.finish();
assert_eq!(lines.len(), 1);
assert_eq!(lines[0].len(), MAX_LINE_BYTES + 1); assert_eq!(dropped, 0);
}
#[test]
fn the_line_after_an_over_long_one_is_unaffected() {
let mut cap = LineCap::new(10, Keep::Bottom);
cap.feed(&vec![b'x'; MAX_LINE_BYTES * 2]);
cap.feed(b"\nafter\n");
let (lines, _) = cap.finish();
assert_eq!(lines.last().unwrap(), b"after\n");
}
#[test]
fn a_zero_bound_retains_nothing_and_counts_everything() {
let mut cap = LineCap::new(0, Keep::Bottom);
cap.feed(b"a\nb\n");
assert_eq!(cap.finish(), (Vec::<Vec<u8>>::new(), 2));
}
#[test]
fn a_snapshot_reads_the_retained_set_without_consuming_it() {
let mut cap = LineCap::new(10, Keep::Bottom);
cap.feed(b"a\nb\n");
assert_eq!(cap.snapshot(), (vec![b"a\n".to_vec(), b"b\n".to_vec()], 0));
assert_eq!(cap.snapshot(), (vec![b"a\n".to_vec(), b"b\n".to_vec()], 0));
cap.feed(b"c\n");
assert_eq!(cap.snapshot().0.len(), 3);
}
#[test]
fn a_snapshot_respects_the_bound_and_reports_the_same_drop_count() {
let mut cap = LineCap::new(2, Keep::Bottom);
cap.feed(b"1\n2\n3\n4\n");
assert_eq!(cap.snapshot(), (vec![b"3\n".to_vec(), b"4\n".to_vec()], 2));
}
#[test]
fn a_snapshot_excludes_the_unterminated_line_that_finish_would_flush() {
let mut cap = LineCap::new(10, Keep::Bottom);
cap.feed(b"a\nb");
assert_eq!(cap.snapshot(), (vec![b"a\n".to_vec()], 0));
assert_eq!(cap.finish(), (vec![b"a\n".to_vec(), b"b".to_vec()], 0));
}
#[test]
fn a_snapshot_sees_a_line_only_once_its_terminator_arrives() {
let mut cap = LineCap::new(10, Keep::Bottom);
cap.feed(b"par");
assert_eq!(cap.snapshot().0, Vec::<Vec<u8>>::new());
cap.feed(b"tial\n");
assert_eq!(cap.snapshot().0, vec![b"partial\n".to_vec()]);
}
#[test]
fn the_cheap_accessors_agree_with_the_snapshot() {
let mut cap = LineCap::new(2, Keep::Bottom);
cap.feed(b"1\n2\n3\n");
let (lines, dropped) = cap.snapshot();
assert_eq!(cap.retained_len(), lines.len());
assert_eq!(cap.dropped_count(), dropped);
}
}