1#[derive(Debug, Clone, Copy, PartialEq, Eq)]
41pub enum Boundary {
42 Starts,
44 Ends,
46 Neither,
48}
49
50pub fn indented(line: &str) -> Boundary {
59 match line.as_bytes().first() {
60 Some(b' ' | b'\t') => Boundary::Neither,
61 _ => Boundary::Starts,
62 }
63}
64
65#[derive(Debug, Clone, Copy, PartialEq, Eq)]
70pub struct Line<'a> {
71 pub text: &'a str,
72 pub start_offset: u64,
73 pub end_offset: u64,
74}
75
76#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
78pub enum Completion {
79 Boundary,
81 Deadline,
83 Oversized,
85}
86
87#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
89pub struct Record {
90 pub body: String,
92 pub start_offset: u64,
94 pub end_offset: u64,
95 pub lines: usize,
97 pub completion: Completion,
99}
100
101#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
107pub struct Limits {
108 pub max_lines: usize,
109 pub max_bytes: usize,
110}
111
112impl Limits {
113 pub fn new(max_lines: usize, max_bytes: usize) -> Self {
115 Self {
116 max_lines: max_lines.max(1),
117 max_bytes: max_bytes.max(1),
118 }
119 }
120}
121
122#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
124pub enum Start {
125 WaitForStart,
129 Collect,
134}
135
136#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
145pub struct Delimiter {
146 limits: Limits,
147 start: Start,
148 buffer: Option<Record>,
150 dropped_lines: usize,
151 dropped_bytes: usize,
152 emitted: usize,
153}
154
155impl Delimiter {
156 pub fn new(limits: Limits, start: Start) -> Self {
157 Self {
158 limits,
159 start,
160 buffer: None,
161 dropped_lines: 0,
162 dropped_bytes: 0,
163 emitted: 0,
164 }
165 }
166
167 pub fn resume(limits: Limits, start: Start, record: Record) -> Self {
172 Self {
173 limits,
174 start,
175 buffer: Some(record),
176 dropped_lines: 0,
177 dropped_bytes: 0,
178 emitted: 0,
179 }
180 }
181
182 pub fn push(&mut self, line: Line<'_>, signal: impl Fn(&str) -> Boundary) -> Option<Record> {
190 match signal(line.text) {
191 Boundary::Starts => {
192 let sealed = self.seal(Completion::Boundary);
193 self.begin(line);
194 sealed
195 }
196 Boundary::Ends => self.seal(Completion::Boundary),
199 Boundary::Neither => {
200 let verdict = match self.buffer.as_ref() {
203 Some(record) => Some(over_limit(self.limits, record, line.text)),
204 None => None,
205 };
206 match verdict {
207 Some(true) => {
208 let sealed = self.seal(Completion::Oversized);
209 self.drop_line(line);
210 sealed
211 }
212 Some(false) => {
213 let record = self.buffer.as_mut().expect("accumulating");
214 record.body.push_str(line.text);
215 record.end_offset = line.end_offset;
216 record.lines += 1;
217 None
218 }
219 None if self.start == Start::Collect => {
220 self.begin(line);
221 None
222 }
223 None => {
225 self.drop_line(line);
226 None
227 }
228 }
229 }
230 }
231 }
232
233 pub fn flush(&mut self) -> Option<Record> {
237 self.seal(Completion::Deadline)
238 }
239
240 pub fn is_accumulating(&self) -> bool {
242 self.buffer.is_some()
243 }
244
245 pub fn pending(&self) -> Option<&Record> {
249 self.buffer.as_ref()
250 }
251
252 pub fn into_pending(mut self) -> Option<Record> {
257 self.buffer.take()
258 }
259
260 pub fn dropped_lines(&self) -> usize {
264 self.dropped_lines
265 }
266
267 pub fn dropped_bytes(&self) -> usize {
271 self.dropped_bytes
272 }
273
274 pub fn emitted(&self) -> usize {
276 self.emitted
277 }
278
279 fn begin(&mut self, line: Line<'_>) {
282 self.buffer = Some(Record {
283 body: line.text.to_string(),
284 start_offset: line.start_offset,
285 end_offset: line.end_offset,
286 lines: 1,
287 completion: Completion::Boundary,
288 });
289 }
290
291 fn drop_line(&mut self, line: Line<'_>) {
292 self.dropped_lines += 1;
293 self.dropped_bytes += line.text.len();
294 }
295
296 fn seal(&mut self, completion: Completion) -> Option<Record> {
297 let mut record = self.buffer.take()?;
298 record.completion = completion;
299 self.emitted += 1;
300 Some(record)
301 }
302}
303
304fn over_limit(limits: Limits, record: &Record, incoming: &str) -> bool {
306 record.lines + 1 > limits.max_lines || record.body.len() + incoming.len() > limits.max_bytes
307}
308
309#[cfg(test)]
310mod tests {
311 use super::*;
312
313 const LIMITS: Limits = Limits {
314 max_lines: 1000,
315 max_bytes: 1 << 20,
316 };
317
318 fn anchored() -> Delimiter {
319 Delimiter::new(LIMITS, Start::WaitForStart)
320 }
321
322 fn line(text: &str, start: u64) -> Line<'_> {
324 Line {
325 text,
326 start_offset: start,
327 end_offset: start + text.len() as u64,
328 }
329 }
330
331 fn anchor(line: &str) -> Boundary {
333 if line.starts_with("20") && line.contains("+08 ") {
334 Boundary::Starts
335 } else {
336 Boundary::Neither
337 }
338 }
339
340 fn install_log_anchor(line: &str) -> Boundary {
342 let iso = line.starts_with("20") && line.contains("+08 ");
343 let bsd = line.starts_with("Jul ") || line.starts_with("Aug ");
344 if iso || bsd {
345 Boundary::Starts
346 } else {
347 Boundary::Neither
348 }
349 }
350
351 fn blank_separated(line: &str) -> Boundary {
353 if line.trim().is_empty() {
354 Boundary::Ends
355 } else {
356 Boundary::Neither
357 }
358 }
359
360 fn run(delimiter: &mut Delimiter, lines: &[&str]) -> Vec<Record> {
362 let mut offset = 0;
363 let mut out = Vec::new();
364 for text in lines {
365 if let Some(record) = delimiter.push(line(text, offset), anchor) {
366 out.push(record);
367 }
368 offset += text.len() as u64;
369 }
370 out
371 }
372
373 fn run_and_flush(delimiter: &mut Delimiter, lines: &[&str]) -> Vec<Record> {
375 let mut records = run(delimiter, lines);
376 records.extend(delimiter.flush());
377 records
378 }
379
380 #[test]
383 fn the_next_anchor_seals_the_previous_so_the_last_one_stays_open() {
384 let mut delimiter = anchored();
385 let out = run(
386 &mut delimiter,
387 &[
388 "2026-09-23 20:24:24+08 host a: one\n",
389 "\tcont\n",
390 "2026-09-23 20:24:25+08 host a: two\n",
391 ],
392 );
393 assert_eq!(out.len(), 1);
395 assert_eq!(out[0].body, "2026-09-23 20:24:24+08 host a: one\n\tcont\n");
396 assert_eq!(out[0].lines, 2);
397 assert_eq!(out[0].completion, Completion::Boundary);
398 assert_eq!(out[0].start_offset, 0);
399 assert_eq!(out[0].end_offset, 41);
400 assert!(delimiter.is_accumulating());
401
402 let tail = delimiter.flush().expect("flush");
404 assert_eq!(tail.body, "2026-09-23 20:24:25+08 host a: two\n");
405 assert_eq!(tail.completion, Completion::Deadline);
406 assert!(!delimiter.is_accumulating());
407 }
408
409 #[test]
410 fn adjacent_anchors_are_adjacent_single_line_records() {
411 let mut delimiter = anchored();
413 let records = run_and_flush(
414 &mut delimiter,
415 &[
416 "2026-09-23 20:24:24+08 host a: one\n",
417 "2026-09-23 20:24:25+08 host a: two\n",
418 "2026-09-23 20:24:26+08 host a: three\n",
419 ],
420 );
421 assert_eq!(records.len(), 3);
422 assert!(records.iter().all(|record| record.lines == 1));
423 assert_eq!(
424 records
425 .iter()
426 .map(|record| record.body.split_once(": ").expect("body").1)
427 .collect::<Vec<_>>(),
428 vec!["one\n", "two\n", "three\n"]
429 );
430 assert_eq!(records[0].end_offset, records[1].start_offset);
432 }
433
434 #[test]
435 fn a_start_signal_then_an_end_signal_yields_a_single_line_record() {
436 let mut delimiter = anchored();
438 let signal = |line: &str| match anchor(line) {
439 Boundary::Starts => Boundary::Starts,
440 _ if line.trim().is_empty() => Boundary::Ends,
441 _ => Boundary::Neither,
442 };
443 let first = delimiter
444 .push(line("2026-09-23 20:24:24+08 host a: one\n", 0), signal)
445 .is_none();
446 assert!(first);
447 let sealed = delimiter
448 .push(line("\n", 37), signal)
449 .expect("空行把上一条封上");
450 assert_eq!(sealed.body, "2026-09-23 20:24:24+08 host a: one\n");
451 assert_eq!(sealed.lines, 1);
452 assert_eq!(sealed.completion, Completion::Boundary);
453 assert!(!delimiter.is_accumulating());
454 assert_eq!(delimiter.dropped_lines(), 0);
456 }
457
458 #[test]
459 fn an_end_signal_with_nothing_open_is_neither_a_record_nor_a_drop() {
460 let mut delimiter = Delimiter::new(LIMITS, Start::Collect);
462 assert!(delimiter.push(line("\n", 0), blank_separated).is_none());
463 assert!(delimiter.push(line(" \n", 1), blank_separated).is_none());
464 assert_eq!(delimiter.emitted(), 0);
465 assert_eq!(delimiter.dropped_lines(), 0, "空行是边界,不是丢弃");
466 assert!(!delimiter.is_accumulating());
467 }
468
469 #[test]
470 fn an_empty_line_under_the_indented_reader_becomes_a_record_of_nothing() {
471 let mut delimiter = anchored();
474 assert!(delimiter.push(line("\n", 0), indented).is_none());
475 let sealed = delimiter.flush().expect("flush");
476 assert_eq!(sealed.body, "\n");
477 assert_eq!(sealed.lines, 1);
478 assert_eq!(delimiter.dropped_lines(), 0);
479 }
480
481 #[test]
484 fn an_indented_first_line_lands_mid_record_and_is_dropped() {
485 let mut delimiter = anchored();
487 let out = run(
488 &mut delimiter,
489 &[
490 "\tcont of an earlier record\n",
491 "\tmore of it\n",
492 "2026-09-23 20:24:24+08 host a: real\n",
493 ],
494 );
495 assert!(out.is_empty());
496 assert_eq!(delimiter.dropped_lines(), 2);
497 assert_eq!(
498 delimiter.dropped_bytes(),
499 "\tcont of an earlier record\n".len() + "\tmore of it\n".len()
500 );
501 assert_eq!(
502 delimiter.flush().expect("flush").body,
503 "2026-09-23 20:24:24+08 host a: real\n"
504 );
505 }
506
507 #[test]
508 fn nothing_ever_arriving_never_becomes_a_record() {
509 for start in [Start::WaitForStart, Start::Collect] {
511 let mut delimiter = Delimiter::new(LIMITS, start);
512 assert!(delimiter.flush().is_none());
513 assert_eq!(delimiter.emitted(), 0);
514 }
515 let mut delimiter = anchored();
516 run(&mut delimiter, &["\tno head\n"]);
517 assert!(delimiter.flush().is_none());
518 assert_eq!(delimiter.dropped_lines(), 1);
519 }
520
521 #[test]
522 fn a_record_spanning_several_batches_is_emitted_exactly_once() {
523 let mut delimiter = anchored();
525 assert!(run(&mut delimiter, &["2026-09-23 20:24:24+08 host a: one\n"]).is_empty());
526 assert!(run(&mut delimiter, &["\tcont A\n"]).is_empty());
527 assert!(run(&mut delimiter, &["\tcont B\n"]).is_empty());
528 let out = run(&mut delimiter, &["2026-09-23 20:24:25+08 host a: two\n"]);
529 assert_eq!(out.len(), 1);
530 assert_eq!(
531 out[0].body,
532 "2026-09-23 20:24:24+08 host a: one\n\tcont A\n\tcont B\n"
533 );
534 assert_eq!(out[0].lines, 3);
535 }
536
537 #[test]
538 fn a_collecting_start_drops_nothing_up_front() {
539 let mut delimiter = Delimiter::new(LIMITS, Start::Collect);
541 let mut records = Vec::new();
542 let mut offset = 0;
543 for text in ["first\n", "second\n"] {
544 if let Some(record) = delimiter.push(line(text, offset), blank_separated) {
545 records.push(record);
546 }
547 offset += text.len() as u64;
548 }
549 records.extend(delimiter.flush());
550 assert_eq!(delimiter.dropped_lines(), 0);
551 assert_eq!(records.len(), 1);
552 assert_eq!(records[0].body, "first\nsecond\n");
553 assert_eq!(records[0].start_offset, 0);
554 assert_eq!(records[0].end_offset, 13);
555 }
556
557 #[test]
558 fn a_start_signal_also_works_in_a_collecting_format() {
559 let mut delimiter = Delimiter::new(LIMITS, Start::Collect);
561 let signal = |line: &str| match anchor(line) {
562 Boundary::Starts => Boundary::Starts,
563 _ if line.trim().is_empty() => Boundary::Ends,
564 _ => Boundary::Neither,
565 };
566 let mut records = Vec::new();
567 let mut offset = 0;
568 for text in [
569 "junk without an anchor\n",
570 "2026-09-23 20:24:24+08 host a: one\n",
571 "\n",
572 "2026-09-23 20:24:25+08 host a: two\n",
573 ] {
574 if let Some(record) = delimiter.push(line(text, offset), signal) {
575 records.push(record);
576 }
577 offset += text.len() as u64;
578 }
579 records.extend(delimiter.flush());
580 assert_eq!(
581 records
582 .iter()
583 .map(|record| record.body.as_str())
584 .collect::<Vec<_>>(),
585 vec![
586 "junk without an anchor\n",
587 "2026-09-23 20:24:24+08 host a: one\n",
588 "2026-09-23 20:24:25+08 host a: two\n",
589 ]
590 );
591 assert_eq!(delimiter.dropped_lines(), 0);
592 }
593
594 #[test]
597 fn the_line_limit_is_inclusive_and_the_next_line_is_refused() {
598 let mut delimiter = Delimiter::new(Limits::new(3, 1 << 20), Start::WaitForStart);
600 let out = run(
601 &mut delimiter,
602 &[
603 "2026-09-23 20:24:24+08 host a: one\n",
604 "\tcont 1\n",
605 "\tcont 2\n",
606 "\tcont 3\n",
607 "2026-09-23 20:24:25+08 host a: two\n",
608 ],
609 );
610 assert_eq!(out.len(), 1);
611 assert_eq!(out[0].lines, 3, "正好等于上限要放行");
612 assert_eq!(out[0].completion, Completion::Oversized);
613 assert_eq!(delimiter.dropped_lines(), 1);
614 }
615
616 #[test]
617 fn the_byte_limit_is_inclusive() {
618 let one = "2026-09-23 20:24:24+08 host a: one\n";
619 let cont = "\tcont\n";
620 let mut delimiter = Delimiter::new(
622 Limits::new(1000, one.len() + cont.len()),
623 Start::WaitForStart,
624 );
625 assert!(delimiter.push(line(one, 0), anchor).is_none());
626 assert!(
627 delimiter
628 .push(line(cont, one.len() as u64), anchor)
629 .is_none(),
630 "正好等于上限要放行"
631 );
632 let sealed = delimiter
633 .push(line(cont, (one.len() + cont.len()) as u64), anchor)
634 .expect("第三行越限");
635 assert_eq!(sealed.completion, Completion::Oversized);
636 assert_eq!(sealed.lines, 2);
637 assert_eq!(sealed.body.len(), one.len() + cont.len());
638 }
639
640 #[test]
641 fn an_oversized_record_is_sealed_and_marked_and_the_rest_is_dropped() {
642 let mut delimiter = Delimiter::new(Limits::new(3, 4096), Start::WaitForStart);
644 let out = run(
645 &mut delimiter,
646 &[
647 "2026-09-23 20:24:24+08 host a: one\n",
648 "\tcont 1\n",
649 "\tcont 2\n",
650 "\tcont 3\n",
651 "\tcont 4\n",
652 "2026-09-23 20:24:25+08 host a: two\n",
653 ],
654 );
655 assert_eq!(out.len(), 1);
656 assert_eq!(out[0].completion, Completion::Oversized);
657 assert_eq!(
658 out[0].body,
659 "2026-09-23 20:24:24+08 host a: one\n\tcont 1\n\tcont 2\n"
660 );
661 assert_eq!(delimiter.dropped_lines(), 2);
663 assert_eq!(
664 delimiter.flush().expect("flush").body,
665 "2026-09-23 20:24:25+08 host a: two\n"
666 );
667 }
668
669 #[test]
670 fn a_single_line_over_the_limit_is_not_cut_in_half() {
671 let mut delimiter = Delimiter::new(Limits::new(10, 8), Start::WaitForStart);
674 let long = "2026-09-23 20:24:24+08 host a: a very long line\n";
675 assert!(
676 delimiter.push(line(long, 0), anchor).is_none(),
677 "单行本身就超限:先收下,不假装能截"
678 );
679 let sealed = delimiter
680 .push(line("\tcont\n", long.len() as u64), anchor)
681 .expect("第二条续行到限,把上一条封上");
682 assert_eq!(sealed.completion, Completion::Oversized);
683 assert_eq!(sealed.body, long, "单行原样保留,没被截");
684 assert_eq!(sealed.lines, 1);
685 assert_eq!(delimiter.dropped_lines(), 1);
686 }
687
688 #[test]
689 fn after_an_oversized_seal_a_collecting_format_restarts_on_the_next_line() {
690 let mut delimiter = Delimiter::new(Limits::new(2, 1 << 20), Start::Collect);
694 let mut records = Vec::new();
695 let mut offset = 0;
696 for text in ["r1 a\n", "r1 b\n", "r1 c\n", "r1 d\n", "\n", "r2\n"] {
697 if let Some(record) = delimiter.push(line(text, offset), blank_separated) {
698 records.push(record);
699 }
700 offset += text.len() as u64;
701 }
702 records.extend(delimiter.flush());
703 assert_eq!(records.len(), 3);
704 assert_eq!(records[0].completion, Completion::Oversized);
705 assert_eq!(records[0].body, "r1 a\nr1 b\n");
706 assert_eq!(records[1].body, "r1 d\n", "无头的那截");
707 assert_eq!(records[2].body, "r2\n");
708 assert_eq!(delimiter.dropped_lines(), 1);
709 }
710
711 #[test]
714 fn a_delimiter_survives_a_round_trip_through_json() {
715 let junk = "\tno head\n";
718 let one = "2026-09-23 20:24:24+08 host a: one\n";
719 let cont = "\tcont A\n";
720 let mut delimiter = anchored();
721 assert!(delimiter.push(line(junk, 0), anchor).is_none());
723 let base = junk.len() as u64;
725 assert!(delimiter.push(line(one, base), anchor).is_none());
726 assert!(
727 delimiter
728 .push(line(cont, base + one.len() as u64), anchor)
729 .is_none()
730 );
731 assert_eq!(delimiter.dropped_lines(), 1);
732
733 let json = serde_json::to_string(&delimiter).expect("serialize");
734 let mut resumed: Delimiter = serde_json::from_str(&json).expect("deserialize");
735 assert_eq!(
736 resumed.pending().expect("pending").body,
737 format!("{one}{cont}")
738 );
739 assert_eq!(resumed.dropped_lines(), 1);
740 assert_eq!(resumed.dropped_bytes(), junk.len());
741
742 let next = base + one.len() as u64 + cont.len() as u64;
744 let sealed = resumed
745 .push(line("2026-09-23 20:24:25+08 host a: two\n", next), anchor)
746 .expect("sealed");
747 assert_eq!(sealed.body, format!("{one}{cont}"));
748 assert_eq!(sealed.lines, 2);
749 assert_eq!(sealed.completion, Completion::Boundary);
750 assert_eq!(sealed.start_offset, base);
751 assert_eq!(sealed.end_offset, next);
752 assert_eq!(resumed.emitted(), 1);
753 }
754
755 #[test]
756 fn offsets_tile_the_input_without_gaps_or_overlap() {
757 let mut delimiter = anchored();
759 let lines = [
760 "junk\n",
761 "2026-09-23 20:24:24+08 host a: one\n",
762 "\tcont\n",
763 "2026-09-23 20:24:25+08 host a: two\n",
764 "2026-09-23 20:24:26+08 host a: three\n",
765 ];
766 let records = run_and_flush(&mut delimiter, &lines);
767
768 assert_eq!(records[0].start_offset, "junk\n".len() as u64);
770 assert_eq!(delimiter.dropped_bytes(), "junk\n".len());
771 for pair in records.windows(2) {
772 assert_eq!(
773 pair[0].end_offset, pair[1].start_offset,
774 "记录之间不许有洞或重叠:{pair:?}"
775 );
776 }
777 let total: u64 = lines.iter().map(|text| text.len() as u64).sum();
778 assert_eq!(records.last().expect("records").end_offset, total);
779 }
780
781 #[test]
782 fn records_come_out_in_order() {
783 let mut delimiter = anchored();
784 let records = run_and_flush(
785 &mut delimiter,
786 &[
787 "2026-09-23 20:24:24+08 host a: one\n",
788 "2026-09-23 20:24:25+08 host a: two\n",
789 "2026-09-23 20:24:26+08 host a: three\n",
790 ],
791 );
792 let starts: Vec<u64> = records.iter().map(|record| record.start_offset).collect();
794 let mut sorted = starts.clone();
795 sorted.sort_unstable();
796 assert_eq!(starts, sorted);
797 assert!(starts.windows(2).all(|pair| pair[0] < pair[1]));
798 assert_eq!(
799 records
800 .iter()
801 .map(|record| record.body.split_once(": ").expect("body").1)
802 .collect::<Vec<_>>(),
803 vec!["one\n", "two\n", "three\n"]
804 );
805 }
806
807 struct Rng(u64);
809
810 impl Rng {
811 fn next(&mut self) -> u64 {
812 self.0 ^= self.0 << 13;
813 self.0 ^= self.0 >> 7;
814 self.0 ^= self.0 << 17;
815 self.0
816 }
817
818 fn pick(&mut self, bound: usize) -> usize {
819 (self.next() % bound as u64) as usize
820 }
821 }
822
823 #[test]
824 fn no_byte_is_ever_silently_lost() {
825 let shapes = [
828 "2026-09-23 20:24:24+08 host a: anchor line\n",
829 "\tcontinuation\n",
830 " another continuation\n",
831 "plain line without an anchor\n",
832 ];
833 let mut rng = Rng(0x5eed_1234_5678_9abc);
834 let mut input = String::new();
835 let mut lines: Vec<(&str, u64)> = Vec::new();
836 for _ in 0..400 {
837 let text = shapes[rng.pick(shapes.len())];
838 lines.push((text, input.len() as u64));
839 input.push_str(text);
840 }
841 lines.push((shapes[0], input.len() as u64));
844 input.push_str(shapes[0]);
845
846 let mut delimiter = Delimiter::new(Limits::new(3, 64), Start::WaitForStart);
847 let mut records = Vec::new();
848 let mut index = 0;
849 while index < lines.len() {
850 let batch = 1 + rng.pick(5);
851 for (text, start) in &lines[index..(index + batch).min(lines.len())] {
852 if let Some(record) = delimiter.push(line(text, *start), anchor) {
853 records.push(record);
854 }
855 }
856 index += batch;
857 }
858 records.extend(delimiter.flush());
859
860 assert!(
861 records.len() > 10,
862 "流里应当确实产出记录:{}",
863 records.len()
864 );
865 assert!(delimiter.dropped_bytes() > 0, "应当走到过丢弃路径");
867 assert!(
868 records
869 .iter()
870 .any(|record| record.completion == Completion::Oversized),
871 "应当走到过超限路径"
872 );
873 assert!(
874 records
875 .iter()
876 .any(|record| record.completion == Completion::Deadline),
877 "应当走到过到期封口路径"
878 );
879 for record in &records {
880 assert_eq!(
881 record.body,
882 input[record.start_offset as usize..record.end_offset as usize],
883 "记录的正文必须与它的区间逐字节对得上"
884 );
885 assert!(record.start_offset < record.end_offset, "{record:?}");
886 }
887 for pair in records.windows(2) {
888 assert!(
889 pair[0].end_offset <= pair[1].start_offset,
890 "记录不许重叠或倒序:{pair:?}"
891 );
892 }
893 let covered: u64 = records
894 .iter()
895 .map(|record| record.end_offset - record.start_offset)
896 .sum();
897 assert_eq!(
898 covered + delimiter.dropped_bytes() as u64,
899 input.len() as u64,
900 "有字节去向不明(既不在记录里,也没被计入丢弃)"
901 );
902 }
903
904 #[test]
907 fn the_indented_reader_says_what_it_means() {
908 assert_eq!(indented("plain line\n"), Boundary::Starts);
909 assert_eq!(indented(" spaced\n"), Boundary::Neither);
910 assert_eq!(indented("\ttabbed\n"), Boundary::Neither);
911 assert_eq!(indented("\n"), Boundary::Starts);
913 }
914
915 #[test]
916 fn an_install_log_shaped_stream_folds_by_its_two_anchors() {
917 let mut delimiter = Delimiter::new(LIMITS, Start::WaitForStart);
920 let lines = [
921 "2026-09-23 20:24:24+08 MBP softwareupdated[565]: Setting up (\n",
922 "\t\"<SUOSUProduct: MSU>\",\n",
923 "\t)\n",
924 "Jul 17 12:00:47 MBP Installer Progress[66]: phases set to (\n",
925 "\t\"phase one\",\n",
926 "\t)\n",
927 "2026-09-23 20:24:25+08 MBP loginwindow[428]: policy = 0\n",
928 ];
929 let mut offset = 0;
930 let mut records = Vec::new();
931 for text in lines {
932 if let Some(record) = delimiter.push(line(text, offset), install_log_anchor) {
933 records.push(record);
934 }
935 offset += text.len() as u64;
936 }
937 records.extend(delimiter.flush());
938 assert_eq!(records.len(), 3);
939 assert_eq!(records[0].lines, 3);
940 assert_eq!(records[1].lines, 3, "BSD 锚也要认(它同样开一条新记录)");
941 assert_eq!(records[2].lines, 1);
942 assert_eq!(delimiter.dropped_lines(), 0);
943 assert_eq!(
944 records[0].body,
945 "2026-09-23 20:24:24+08 MBP softwareupdated[565]: Setting up (\n\t\"<SUOSUProduct: MSU>\",\n\t)\n"
946 );
947 }
948
949 #[test]
952 fn limits_new_clamps_both_fields_to_at_least_one() {
953 let clamped = Limits::new(0, 0);
954 assert_eq!(clamped.max_lines, 1);
955 assert_eq!(clamped.max_bytes, 1);
956 assert_eq!(
957 Limits::new(5, 9),
958 Limits {
959 max_lines: 5,
960 max_bytes: 9,
961 }
962 );
963 }
964
965 #[test]
966 fn pending_exposes_the_open_record_without_sealing_it() {
967 let one = "2026-09-23 20:24:24+08 host a: one\n";
968 let cont = "\tcont\n";
969 let mut delimiter = anchored();
970 assert!(delimiter.push(line(one, 0), anchor).is_none());
971 assert!(
972 delimiter
973 .push(line(cont, one.len() as u64), anchor)
974 .is_none()
975 );
976
977 let pending = delimiter.pending().expect("pending");
978 assert_eq!(pending.body, format!("{one}{cont}"));
979 assert_eq!(pending.lines, 2);
980 assert_eq!(pending.start_offset, 0);
981 assert_eq!(pending.end_offset, (one.len() + cont.len()) as u64);
982 assert_eq!(delimiter.emitted(), 0, "pending 不算已产出");
983 assert!(delimiter.is_accumulating());
984 }
985
986 #[test]
987 fn into_pending_hands_over_content_and_is_none_when_idle() {
988 assert!(anchored().into_pending().is_none());
989
990 let one = "2026-09-23 20:24:24+08 host a: one\n";
991 let mut delimiter = anchored();
992 assert!(delimiter.push(line(one, 0), anchor).is_none());
993 let taken = delimiter.into_pending().expect("pending");
994 assert_eq!(taken.body, one);
995 assert_eq!(taken.lines, 1);
996 }
997
998 #[test]
999 fn resume_takes_the_record_as_is_and_resets_counters() {
1000 let one = "2026-09-23 20:24:24+08 host a: one\n";
1001 let record = Record {
1002 body: one.to_string(),
1003 start_offset: 10,
1004 end_offset: 10 + one.len() as u64,
1005 lines: 1,
1006 completion: Completion::Oversized,
1008 };
1009 let mut delimiter = Delimiter::resume(Limits::new(10, 4096), Start::WaitForStart, record);
1010 assert_eq!(delimiter.dropped_lines(), 0);
1011 assert_eq!(delimiter.dropped_bytes(), 0);
1012 assert_eq!(delimiter.emitted(), 0);
1013 assert!(delimiter.is_accumulating());
1014
1015 let sealed = delimiter.flush().expect("flush");
1016 assert_eq!(sealed.completion, Completion::Deadline);
1017 assert_eq!(sealed.body, one);
1018 assert_eq!(sealed.start_offset, 10);
1019 assert_eq!(sealed.end_offset, 10 + one.len() as u64);
1020 assert_eq!(delimiter.emitted(), 1);
1021 }
1022
1023 #[test]
1024 fn completion_and_record_serde_contract() {
1025 assert_eq!(
1027 serde_json::to_string(&Completion::Boundary).expect("serialize"),
1028 "\"Boundary\""
1029 );
1030 assert_eq!(
1031 serde_json::to_string(&Completion::Deadline).expect("serialize"),
1032 "\"Deadline\""
1033 );
1034 assert_eq!(
1035 serde_json::to_string(&Completion::Oversized).expect("serialize"),
1036 "\"Oversized\""
1037 );
1038 assert_eq!(
1039 serde_json::from_str::<Completion>("\"Oversized\"").expect("deserialize"),
1040 Completion::Oversized
1041 );
1042
1043 let with_extra = r#"{"body":"x\n","start_offset":0,"end_offset":2,"lines":1,"completion":"Boundary","future":42}"#;
1045 let record: Record = serde_json::from_str(with_extra).expect("未知字段应被忽略");
1046 assert_eq!(record.body, "x\n");
1047 assert_eq!(record.completion, Completion::Boundary);
1048
1049 let missing = r#"{"body":"x\n","start_offset":0,"end_offset":2,"lines":1}"#;
1050 assert!(
1051 serde_json::from_str::<Record>(missing).is_err(),
1052 "缺 completion 应当报错"
1053 );
1054 }
1055
1056 #[test]
1059 fn end_signals_are_separators_so_their_bytes_are_accounted_for_separately() {
1060 let shapes = [
1061 "2026-09-23 20:24:24+08 host a: anchor line\n",
1062 "\tcontinuation\n",
1063 " another continuation\n",
1064 "\n", ];
1066 let signal = |line: &str| {
1067 if line.trim().is_empty() {
1068 Boundary::Ends
1069 } else if anchor(line) == Boundary::Starts {
1070 Boundary::Starts
1071 } else {
1072 Boundary::Neither
1073 }
1074 };
1075
1076 let mut rng = Rng(0x0bad_c0de_dead_beef);
1077 let mut input = String::new();
1078 let mut lines: Vec<(&str, u64)> = Vec::new();
1079 for _ in 0..400 {
1080 let text = shapes[rng.pick(shapes.len())];
1081 lines.push((text, input.len() as u64));
1082 input.push_str(text);
1083 }
1084 lines.push((shapes[3], input.len() as u64));
1086 input.push_str(shapes[3]);
1087
1088 let mut delimiter = Delimiter::new(Limits::new(3, 48), Start::WaitForStart);
1089 let mut records = Vec::new();
1090 let mut index = 0;
1091 while index < lines.len() {
1092 let batch = 1 + rng.pick(4);
1093 for (text, start) in &lines[index..(index + batch).min(lines.len())] {
1094 if let Some(record) = delimiter.push(line(text, *start), signal) {
1095 records.push(record);
1096 }
1097 }
1098 index += batch;
1099 }
1100
1101 assert!(
1103 records
1104 .iter()
1105 .any(|record| record.completion == Completion::Oversized),
1106 "应当走到过超限路径"
1107 );
1108 assert!(delimiter.dropped_bytes() > 0, "应当走到过丢弃路径");
1109
1110 let separators: u64 = lines
1111 .iter()
1112 .filter(|(text, _)| signal(text) == Boundary::Ends)
1113 .map(|(text, _)| text.len() as u64)
1114 .sum();
1115 assert!(separators > 0, "必须真的出现过结束信号");
1116
1117 for record in &records {
1118 assert_eq!(
1119 record.body,
1120 input[record.start_offset as usize..record.end_offset as usize],
1121 "记录的正文必须与它的区间逐字节对得上"
1122 );
1123 }
1124 let covered: u64 = records
1125 .iter()
1126 .map(|record| record.end_offset - record.start_offset)
1127 .sum();
1128 assert_eq!(
1129 covered + delimiter.dropped_bytes() as u64 + separators,
1130 input.len() as u64,
1131 "有字节去向不明(既不在记录里,也没计入丢弃或间隔)"
1132 );
1133
1134 let mut collecting = Delimiter::new(Limits::new(3, 48), Start::Collect);
1136 let mut records = Vec::new();
1137 for (text, start) in &lines {
1138 if let Some(record) = collecting.push(line(text, *start), signal) {
1139 records.push(record);
1140 }
1141 }
1142 records.extend(collecting.flush());
1143 let covered: u64 = records
1144 .iter()
1145 .map(|record| record.end_offset - record.start_offset)
1146 .sum();
1147 assert_eq!(
1148 covered + collecting.dropped_bytes() as u64 + separators,
1149 input.len() as u64
1150 );
1151 }
1152}