#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Boundary {
Starts,
Ends,
Neither,
}
pub fn indented(line: &str) -> Boundary {
match line.as_bytes().first() {
Some(b' ' | b'\t') => Boundary::Neither,
_ => Boundary::Starts,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Line<'a> {
pub text: &'a str,
pub start_offset: u64,
pub end_offset: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub enum Completion {
Boundary,
Deadline,
Oversized,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct Record {
pub body: String,
pub start_offset: u64,
pub end_offset: u64,
pub lines: usize,
pub completion: Completion,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct Limits {
pub max_lines: usize,
pub max_bytes: usize,
}
impl Limits {
pub fn new(max_lines: usize, max_bytes: usize) -> Self {
Self {
max_lines: max_lines.max(1),
max_bytes: max_bytes.max(1),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub enum Start {
WaitForStart,
Collect,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct Delimiter {
limits: Limits,
start: Start,
buffer: Option<Record>,
dropped_lines: usize,
dropped_bytes: usize,
emitted: usize,
}
impl Delimiter {
pub fn new(limits: Limits, start: Start) -> Self {
Self {
limits,
start,
buffer: None,
dropped_lines: 0,
dropped_bytes: 0,
emitted: 0,
}
}
pub fn resume(limits: Limits, start: Start, record: Record) -> Self {
Self {
limits,
start,
buffer: Some(record),
dropped_lines: 0,
dropped_bytes: 0,
emitted: 0,
}
}
pub fn push(&mut self, line: Line<'_>, signal: impl Fn(&str) -> Boundary) -> Option<Record> {
match signal(line.text) {
Boundary::Starts => {
let sealed = self.seal(Completion::Boundary);
self.begin(line);
sealed
}
Boundary::Ends => self.seal(Completion::Boundary),
Boundary::Neither => {
let verdict = match self.buffer.as_ref() {
Some(record) => Some(over_limit(self.limits, record, line.text)),
None => None,
};
match verdict {
Some(true) => {
let sealed = self.seal(Completion::Oversized);
self.drop_line(line);
sealed
}
Some(false) => {
let record = self.buffer.as_mut().expect("accumulating");
record.body.push_str(line.text);
record.end_offset = line.end_offset;
record.lines += 1;
None
}
None if self.start == Start::Collect => {
self.begin(line);
None
}
None => {
self.drop_line(line);
None
}
}
}
}
}
pub fn flush(&mut self) -> Option<Record> {
self.seal(Completion::Deadline)
}
pub fn is_accumulating(&self) -> bool {
self.buffer.is_some()
}
pub fn pending(&self) -> Option<&Record> {
self.buffer.as_ref()
}
pub fn into_pending(mut self) -> Option<Record> {
self.buffer.take()
}
pub fn dropped_lines(&self) -> usize {
self.dropped_lines
}
pub fn dropped_bytes(&self) -> usize {
self.dropped_bytes
}
pub fn emitted(&self) -> usize {
self.emitted
}
fn begin(&mut self, line: Line<'_>) {
self.buffer = Some(Record {
body: line.text.to_string(),
start_offset: line.start_offset,
end_offset: line.end_offset,
lines: 1,
completion: Completion::Boundary,
});
}
fn drop_line(&mut self, line: Line<'_>) {
self.dropped_lines += 1;
self.dropped_bytes += line.text.len();
}
fn seal(&mut self, completion: Completion) -> Option<Record> {
let mut record = self.buffer.take()?;
record.completion = completion;
self.emitted += 1;
Some(record)
}
}
fn over_limit(limits: Limits, record: &Record, incoming: &str) -> bool {
record.lines + 1 > limits.max_lines || record.body.len() + incoming.len() > limits.max_bytes
}
#[cfg(test)]
mod tests {
use super::*;
const LIMITS: Limits = Limits {
max_lines: 1000,
max_bytes: 1 << 20,
};
fn anchored() -> Delimiter {
Delimiter::new(LIMITS, Start::WaitForStart)
}
fn line(text: &str, start: u64) -> Line<'_> {
Line {
text,
start_offset: start,
end_offset: start + text.len() as u64,
}
}
fn anchor(line: &str) -> Boundary {
if line.starts_with("20") && line.contains("+08 ") {
Boundary::Starts
} else {
Boundary::Neither
}
}
fn install_log_anchor(line: &str) -> Boundary {
let iso = line.starts_with("20") && line.contains("+08 ");
let bsd = line.starts_with("Jul ") || line.starts_with("Aug ");
if iso || bsd {
Boundary::Starts
} else {
Boundary::Neither
}
}
fn blank_separated(line: &str) -> Boundary {
if line.trim().is_empty() {
Boundary::Ends
} else {
Boundary::Neither
}
}
fn run(delimiter: &mut Delimiter, lines: &[&str]) -> Vec<Record> {
let mut offset = 0;
let mut out = Vec::new();
for text in lines {
if let Some(record) = delimiter.push(line(text, offset), anchor) {
out.push(record);
}
offset += text.len() as u64;
}
out
}
fn run_and_flush(delimiter: &mut Delimiter, lines: &[&str]) -> Vec<Record> {
let mut records = run(delimiter, lines);
records.extend(delimiter.flush());
records
}
#[test]
fn the_next_anchor_seals_the_previous_so_the_last_one_stays_open() {
let mut delimiter = anchored();
let out = run(
&mut delimiter,
&[
"2026-09-23 20:24:24+08 host a: one\n",
"\tcont\n",
"2026-09-23 20:24:25+08 host a: two\n",
],
);
assert_eq!(out.len(), 1);
assert_eq!(out[0].body, "2026-09-23 20:24:24+08 host a: one\n\tcont\n");
assert_eq!(out[0].lines, 2);
assert_eq!(out[0].completion, Completion::Boundary);
assert_eq!(out[0].start_offset, 0);
assert_eq!(out[0].end_offset, 41);
assert!(delimiter.is_accumulating());
let tail = delimiter.flush().expect("flush");
assert_eq!(tail.body, "2026-09-23 20:24:25+08 host a: two\n");
assert_eq!(tail.completion, Completion::Deadline);
assert!(!delimiter.is_accumulating());
}
#[test]
fn adjacent_anchors_are_adjacent_single_line_records() {
let mut delimiter = anchored();
let records = run_and_flush(
&mut delimiter,
&[
"2026-09-23 20:24:24+08 host a: one\n",
"2026-09-23 20:24:25+08 host a: two\n",
"2026-09-23 20:24:26+08 host a: three\n",
],
);
assert_eq!(records.len(), 3);
assert!(records.iter().all(|record| record.lines == 1));
assert_eq!(
records
.iter()
.map(|record| record.body.split_once(": ").expect("body").1)
.collect::<Vec<_>>(),
vec!["one\n", "two\n", "three\n"]
);
assert_eq!(records[0].end_offset, records[1].start_offset);
}
#[test]
fn a_start_signal_then_an_end_signal_yields_a_single_line_record() {
let mut delimiter = anchored();
let signal = |line: &str| match anchor(line) {
Boundary::Starts => Boundary::Starts,
_ if line.trim().is_empty() => Boundary::Ends,
_ => Boundary::Neither,
};
let first = delimiter
.push(line("2026-09-23 20:24:24+08 host a: one\n", 0), signal)
.is_none();
assert!(first);
let sealed = delimiter
.push(line("\n", 37), signal)
.expect("空行把上一条封上");
assert_eq!(sealed.body, "2026-09-23 20:24:24+08 host a: one\n");
assert_eq!(sealed.lines, 1);
assert_eq!(sealed.completion, Completion::Boundary);
assert!(!delimiter.is_accumulating());
assert_eq!(delimiter.dropped_lines(), 0);
}
#[test]
fn an_end_signal_with_nothing_open_is_neither_a_record_nor_a_drop() {
let mut delimiter = Delimiter::new(LIMITS, Start::Collect);
assert!(delimiter.push(line("\n", 0), blank_separated).is_none());
assert!(delimiter.push(line(" \n", 1), blank_separated).is_none());
assert_eq!(delimiter.emitted(), 0);
assert_eq!(delimiter.dropped_lines(), 0, "空行是边界,不是丢弃");
assert!(!delimiter.is_accumulating());
}
#[test]
fn an_empty_line_under_the_indented_reader_becomes_a_record_of_nothing() {
let mut delimiter = anchored();
assert!(delimiter.push(line("\n", 0), indented).is_none());
let sealed = delimiter.flush().expect("flush");
assert_eq!(sealed.body, "\n");
assert_eq!(sealed.lines, 1);
assert_eq!(delimiter.dropped_lines(), 0);
}
#[test]
fn an_indented_first_line_lands_mid_record_and_is_dropped() {
let mut delimiter = anchored();
let out = run(
&mut delimiter,
&[
"\tcont of an earlier record\n",
"\tmore of it\n",
"2026-09-23 20:24:24+08 host a: real\n",
],
);
assert!(out.is_empty());
assert_eq!(delimiter.dropped_lines(), 2);
assert_eq!(
delimiter.dropped_bytes(),
"\tcont of an earlier record\n".len() + "\tmore of it\n".len()
);
assert_eq!(
delimiter.flush().expect("flush").body,
"2026-09-23 20:24:24+08 host a: real\n"
);
}
#[test]
fn nothing_ever_arriving_never_becomes_a_record() {
for start in [Start::WaitForStart, Start::Collect] {
let mut delimiter = Delimiter::new(LIMITS, start);
assert!(delimiter.flush().is_none());
assert_eq!(delimiter.emitted(), 0);
}
let mut delimiter = anchored();
run(&mut delimiter, &["\tno head\n"]);
assert!(delimiter.flush().is_none());
assert_eq!(delimiter.dropped_lines(), 1);
}
#[test]
fn a_record_spanning_several_batches_is_emitted_exactly_once() {
let mut delimiter = anchored();
assert!(run(&mut delimiter, &["2026-09-23 20:24:24+08 host a: one\n"]).is_empty());
assert!(run(&mut delimiter, &["\tcont A\n"]).is_empty());
assert!(run(&mut delimiter, &["\tcont B\n"]).is_empty());
let out = run(&mut delimiter, &["2026-09-23 20:24:25+08 host a: two\n"]);
assert_eq!(out.len(), 1);
assert_eq!(
out[0].body,
"2026-09-23 20:24:24+08 host a: one\n\tcont A\n\tcont B\n"
);
assert_eq!(out[0].lines, 3);
}
#[test]
fn a_collecting_start_drops_nothing_up_front() {
let mut delimiter = Delimiter::new(LIMITS, Start::Collect);
let mut records = Vec::new();
let mut offset = 0;
for text in ["first\n", "second\n"] {
if let Some(record) = delimiter.push(line(text, offset), blank_separated) {
records.push(record);
}
offset += text.len() as u64;
}
records.extend(delimiter.flush());
assert_eq!(delimiter.dropped_lines(), 0);
assert_eq!(records.len(), 1);
assert_eq!(records[0].body, "first\nsecond\n");
assert_eq!(records[0].start_offset, 0);
assert_eq!(records[0].end_offset, 13);
}
#[test]
fn a_start_signal_also_works_in_a_collecting_format() {
let mut delimiter = Delimiter::new(LIMITS, Start::Collect);
let signal = |line: &str| match anchor(line) {
Boundary::Starts => Boundary::Starts,
_ if line.trim().is_empty() => Boundary::Ends,
_ => Boundary::Neither,
};
let mut records = Vec::new();
let mut offset = 0;
for text in [
"junk without an anchor\n",
"2026-09-23 20:24:24+08 host a: one\n",
"\n",
"2026-09-23 20:24:25+08 host a: two\n",
] {
if let Some(record) = delimiter.push(line(text, offset), signal) {
records.push(record);
}
offset += text.len() as u64;
}
records.extend(delimiter.flush());
assert_eq!(
records
.iter()
.map(|record| record.body.as_str())
.collect::<Vec<_>>(),
vec![
"junk without an anchor\n",
"2026-09-23 20:24:24+08 host a: one\n",
"2026-09-23 20:24:25+08 host a: two\n",
]
);
assert_eq!(delimiter.dropped_lines(), 0);
}
#[test]
fn the_line_limit_is_inclusive_and_the_next_line_is_refused() {
let mut delimiter = Delimiter::new(Limits::new(3, 1 << 20), Start::WaitForStart);
let out = run(
&mut delimiter,
&[
"2026-09-23 20:24:24+08 host a: one\n",
"\tcont 1\n",
"\tcont 2\n",
"\tcont 3\n",
"2026-09-23 20:24:25+08 host a: two\n",
],
);
assert_eq!(out.len(), 1);
assert_eq!(out[0].lines, 3, "正好等于上限要放行");
assert_eq!(out[0].completion, Completion::Oversized);
assert_eq!(delimiter.dropped_lines(), 1);
}
#[test]
fn the_byte_limit_is_inclusive() {
let one = "2026-09-23 20:24:24+08 host a: one\n";
let cont = "\tcont\n";
let mut delimiter = Delimiter::new(
Limits::new(1000, one.len() + cont.len()),
Start::WaitForStart,
);
assert!(delimiter.push(line(one, 0), anchor).is_none());
assert!(
delimiter
.push(line(cont, one.len() as u64), anchor)
.is_none(),
"正好等于上限要放行"
);
let sealed = delimiter
.push(line(cont, (one.len() + cont.len()) as u64), anchor)
.expect("第三行越限");
assert_eq!(sealed.completion, Completion::Oversized);
assert_eq!(sealed.lines, 2);
assert_eq!(sealed.body.len(), one.len() + cont.len());
}
#[test]
fn an_oversized_record_is_sealed_and_marked_and_the_rest_is_dropped() {
let mut delimiter = Delimiter::new(Limits::new(3, 4096), Start::WaitForStart);
let out = run(
&mut delimiter,
&[
"2026-09-23 20:24:24+08 host a: one\n",
"\tcont 1\n",
"\tcont 2\n",
"\tcont 3\n",
"\tcont 4\n",
"2026-09-23 20:24:25+08 host a: two\n",
],
);
assert_eq!(out.len(), 1);
assert_eq!(out[0].completion, Completion::Oversized);
assert_eq!(
out[0].body,
"2026-09-23 20:24:24+08 host a: one\n\tcont 1\n\tcont 2\n"
);
assert_eq!(delimiter.dropped_lines(), 2);
assert_eq!(
delimiter.flush().expect("flush").body,
"2026-09-23 20:24:25+08 host a: two\n"
);
}
#[test]
fn a_single_line_over_the_limit_is_not_cut_in_half() {
let mut delimiter = Delimiter::new(Limits::new(10, 8), Start::WaitForStart);
let long = "2026-09-23 20:24:24+08 host a: a very long line\n";
assert!(
delimiter.push(line(long, 0), anchor).is_none(),
"单行本身就超限:先收下,不假装能截"
);
let sealed = delimiter
.push(line("\tcont\n", long.len() as u64), anchor)
.expect("第二条续行到限,把上一条封上");
assert_eq!(sealed.completion, Completion::Oversized);
assert_eq!(sealed.body, long, "单行原样保留,没被截");
assert_eq!(sealed.lines, 1);
assert_eq!(delimiter.dropped_lines(), 1);
}
#[test]
fn after_an_oversized_seal_a_collecting_format_restarts_on_the_next_line() {
let mut delimiter = Delimiter::new(Limits::new(2, 1 << 20), Start::Collect);
let mut records = Vec::new();
let mut offset = 0;
for text in ["r1 a\n", "r1 b\n", "r1 c\n", "r1 d\n", "\n", "r2\n"] {
if let Some(record) = delimiter.push(line(text, offset), blank_separated) {
records.push(record);
}
offset += text.len() as u64;
}
records.extend(delimiter.flush());
assert_eq!(records.len(), 3);
assert_eq!(records[0].completion, Completion::Oversized);
assert_eq!(records[0].body, "r1 a\nr1 b\n");
assert_eq!(records[1].body, "r1 d\n", "无头的那截");
assert_eq!(records[2].body, "r2\n");
assert_eq!(delimiter.dropped_lines(), 1);
}
#[test]
fn a_delimiter_survives_a_round_trip_through_json() {
let junk = "\tno head\n";
let one = "2026-09-23 20:24:24+08 host a: one\n";
let cont = "\tcont A\n";
let mut delimiter = anchored();
assert!(delimiter.push(line(junk, 0), anchor).is_none());
let base = junk.len() as u64;
assert!(delimiter.push(line(one, base), anchor).is_none());
assert!(
delimiter
.push(line(cont, base + one.len() as u64), anchor)
.is_none()
);
assert_eq!(delimiter.dropped_lines(), 1);
let json = serde_json::to_string(&delimiter).expect("serialize");
let mut resumed: Delimiter = serde_json::from_str(&json).expect("deserialize");
assert_eq!(
resumed.pending().expect("pending").body,
format!("{one}{cont}")
);
assert_eq!(resumed.dropped_lines(), 1);
assert_eq!(resumed.dropped_bytes(), junk.len());
let next = base + one.len() as u64 + cont.len() as u64;
let sealed = resumed
.push(line("2026-09-23 20:24:25+08 host a: two\n", next), anchor)
.expect("sealed");
assert_eq!(sealed.body, format!("{one}{cont}"));
assert_eq!(sealed.lines, 2);
assert_eq!(sealed.completion, Completion::Boundary);
assert_eq!(sealed.start_offset, base);
assert_eq!(sealed.end_offset, next);
assert_eq!(resumed.emitted(), 1);
}
#[test]
fn offsets_tile_the_input_without_gaps_or_overlap() {
let mut delimiter = anchored();
let lines = [
"junk\n",
"2026-09-23 20:24:24+08 host a: one\n",
"\tcont\n",
"2026-09-23 20:24:25+08 host a: two\n",
"2026-09-23 20:24:26+08 host a: three\n",
];
let records = run_and_flush(&mut delimiter, &lines);
assert_eq!(records[0].start_offset, "junk\n".len() as u64);
assert_eq!(delimiter.dropped_bytes(), "junk\n".len());
for pair in records.windows(2) {
assert_eq!(
pair[0].end_offset, pair[1].start_offset,
"记录之间不许有洞或重叠:{pair:?}"
);
}
let total: u64 = lines.iter().map(|text| text.len() as u64).sum();
assert_eq!(records.last().expect("records").end_offset, total);
}
#[test]
fn records_come_out_in_order() {
let mut delimiter = anchored();
let records = run_and_flush(
&mut delimiter,
&[
"2026-09-23 20:24:24+08 host a: one\n",
"2026-09-23 20:24:25+08 host a: two\n",
"2026-09-23 20:24:26+08 host a: three\n",
],
);
let starts: Vec<u64> = records.iter().map(|record| record.start_offset).collect();
let mut sorted = starts.clone();
sorted.sort_unstable();
assert_eq!(starts, sorted);
assert!(starts.windows(2).all(|pair| pair[0] < pair[1]));
assert_eq!(
records
.iter()
.map(|record| record.body.split_once(": ").expect("body").1)
.collect::<Vec<_>>(),
vec!["one\n", "two\n", "three\n"]
);
}
struct Rng(u64);
impl Rng {
fn next(&mut self) -> u64 {
self.0 ^= self.0 << 13;
self.0 ^= self.0 >> 7;
self.0 ^= self.0 << 17;
self.0
}
fn pick(&mut self, bound: usize) -> usize {
(self.next() % bound as u64) as usize
}
}
#[test]
fn no_byte_is_ever_silently_lost() {
let shapes = [
"2026-09-23 20:24:24+08 host a: anchor line\n",
"\tcontinuation\n",
" another continuation\n",
"plain line without an anchor\n",
];
let mut rng = Rng(0x5eed_1234_5678_9abc);
let mut input = String::new();
let mut lines: Vec<(&str, u64)> = Vec::new();
for _ in 0..400 {
let text = shapes[rng.pick(shapes.len())];
lines.push((text, input.len() as u64));
input.push_str(text);
}
lines.push((shapes[0], input.len() as u64));
input.push_str(shapes[0]);
let mut delimiter = Delimiter::new(Limits::new(3, 64), Start::WaitForStart);
let mut records = Vec::new();
let mut index = 0;
while index < lines.len() {
let batch = 1 + rng.pick(5);
for (text, start) in &lines[index..(index + batch).min(lines.len())] {
if let Some(record) = delimiter.push(line(text, *start), anchor) {
records.push(record);
}
}
index += batch;
}
records.extend(delimiter.flush());
assert!(
records.len() > 10,
"流里应当确实产出记录:{}",
records.len()
);
assert!(delimiter.dropped_bytes() > 0, "应当走到过丢弃路径");
assert!(
records
.iter()
.any(|record| record.completion == Completion::Oversized),
"应当走到过超限路径"
);
assert!(
records
.iter()
.any(|record| record.completion == Completion::Deadline),
"应当走到过到期封口路径"
);
for record in &records {
assert_eq!(
record.body,
input[record.start_offset as usize..record.end_offset as usize],
"记录的正文必须与它的区间逐字节对得上"
);
assert!(record.start_offset < record.end_offset, "{record:?}");
}
for pair in records.windows(2) {
assert!(
pair[0].end_offset <= pair[1].start_offset,
"记录不许重叠或倒序:{pair:?}"
);
}
let covered: u64 = records
.iter()
.map(|record| record.end_offset - record.start_offset)
.sum();
assert_eq!(
covered + delimiter.dropped_bytes() as u64,
input.len() as u64,
"有字节去向不明(既不在记录里,也没被计入丢弃)"
);
}
#[test]
fn the_indented_reader_says_what_it_means() {
assert_eq!(indented("plain line\n"), Boundary::Starts);
assert_eq!(indented(" spaced\n"), Boundary::Neither);
assert_eq!(indented("\ttabbed\n"), Boundary::Neither);
assert_eq!(indented("\n"), Boundary::Starts);
}
#[test]
fn an_install_log_shaped_stream_folds_by_its_two_anchors() {
let mut delimiter = Delimiter::new(LIMITS, Start::WaitForStart);
let lines = [
"2026-09-23 20:24:24+08 MBP softwareupdated[565]: Setting up (\n",
"\t\"<SUOSUProduct: MSU>\",\n",
"\t)\n",
"Jul 17 12:00:47 MBP Installer Progress[66]: phases set to (\n",
"\t\"phase one\",\n",
"\t)\n",
"2026-09-23 20:24:25+08 MBP loginwindow[428]: policy = 0\n",
];
let mut offset = 0;
let mut records = Vec::new();
for text in lines {
if let Some(record) = delimiter.push(line(text, offset), install_log_anchor) {
records.push(record);
}
offset += text.len() as u64;
}
records.extend(delimiter.flush());
assert_eq!(records.len(), 3);
assert_eq!(records[0].lines, 3);
assert_eq!(records[1].lines, 3, "BSD 锚也要认(它同样开一条新记录)");
assert_eq!(records[2].lines, 1);
assert_eq!(delimiter.dropped_lines(), 0);
assert_eq!(
records[0].body,
"2026-09-23 20:24:24+08 MBP softwareupdated[565]: Setting up (\n\t\"<SUOSUProduct: MSU>\",\n\t)\n"
);
}
}