Skip to main content

wist_shared/
records.rs

1//! 记录边界界定:把**行流**切成**记录**。
2//!
3//! 一行 ≠ 一条记录。同一行文本,可能是一条记录的开头,也可能是上一条的续行。所以
4//! 处理多行日志的本质不是"识别多行",而是**先定「记录边界」** —— 边界不定,后面一切
5//! (解析、去重、计数)都建在流沙上。
6//!
7//! ## 判据是「边界信号」,不是「形态」
8//!
9//! 逐行只需回答一句话:**这一行开不开新记录?**可靠依据是这行**自身**有没有边界信号:
10//!
11//! - [`Boundary::Starts`](**开始型**):这行开一条新记录 —— 适用于边界信号落在**记录里面**
12//!   的格式(如行首时间戳,它是首行的一部分);
13//! - [`Boundary::Ends`](**结束型**):上一条到此结束 —— 适用于边界信号落在**记录之间**的格式
14//!   (如空行,它不属于任何一条);
15//! - [`Boundary::Neither`]:续行。
16//!
17//! 缩进、空行这些**形态**只是边界信号的**代理**。代理在某些文件上成立(续行恰好都缩进),
18//! 换个文件就错(空行分隔的记录、顶格闭合的 `}`)。用形态当判据,等于把"格式知识"降级成
19//! "排版巧合"。
20//!
21//! 本模块**不猜**边界信号长什么样 —— 那是格式知识,只有策展知道。调用方把它作为函数传进来
22//! ([`indented`] 是常见读法的现成实现)。
23//!
24//! ## 算法的四条要点
25//!
26//! 1. **边界是倒推出来的**:一条记录要等到**下一个边界信号**到达才确定结束 —— 所以最后一条
27//!    永远悬着,必须有 [`Delimiter::flush`](空闲 / 来源变化 / 停机)强制封口;
28//! 2. **起点不许发半条**:从文件中间落地时,那截内容没有头。宁可丢半条,不可发半条 ——
29//!    半条记录长得像完整记录,会一路骗过下游的字段解析,而且查不出来;
30//! 3. **有状态就有边界**:累积不能无限增长([`Limits`]),到限**封口并标记**,不静默截断;
31//! 4. **丢弃要可数**:被丢的是"没有头"的部分。计数是**累计**的 —— 要判"连续丢了很久"
32//!    (即边界信号可能配错),由调用方在每次起头后取快照做差。
33//!
34//! ## 无 IO、无时间、无配置
35//!
36//! 时间是调用方的事(该 [`Delimiter::flush`] 时调),所以它既能用在采集侧(从文件尾读),
37//! 也能用在从字节流任意位置读的地方。同一条记录**不会超过** [`Limits`]。
38
39/// 一行在**记录边界**上的角色。
40#[derive(Debug, Clone, Copy, PartialEq, Eq)]
41pub enum Boundary {
42    /// 开始信号:这一行开一条新记录。
43    Starts,
44    /// 结束信号:上一条到此结束,之后重新开始在下一行。这一行本身**不是内容**。
45    Ends,
46    /// 都不是 —— 这一行是上一条的续行。
47    Neither,
48}
49
50/// 常见读法:行首缩进/制表符是上一条的续行([`Boundary`] 的现成实现)。
51///
52/// 这是"边界信号 ≈ 缩进"的**代理读法**:适用于"续行都缩进、非续行都带锚"的文件
53/// (实测 `/var/log/install.log` 上误判 6/338592)。
54///
55/// 语义严格、无例外:**只有**行首是空格或制表符才算续行,其余(含空行)一律算开始信号。
56/// 行内的空行会因此被当成一条空记录 —— 那不是本函数的意外,而是代理读法本身的局限:
57/// 想要别的行为,自己写一个 [`Boundary`] 判定函数即可。
58pub 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/// 送进界定器的一行。
66///
67/// `text` 原样带上(是否含行尾换行由调用方决定),`start_offset`/`end_offset` 是它在
68/// 来源里的字节区间 —— 界定器只搬运、不改写,所以调用方拿到的记录区间能与来源对账。
69#[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/// 一条记录**怎么被封口的** —— 下游据此判它可不可信。
77#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
78pub enum Completion {
79    /// 下一个边界信号到来 —— 记录的边界是**确定**的,内容完整。
80    Boundary,
81    /// 到期封口(空闲 / 来源变化 / 停机)—— 内容完整,但边界是推断的。
82    Deadline,
83    /// 到上限被截 —— **内容可能不全**,下游要当心(解析失败不该归罪于格式)。
84    Oversized,
85}
86
87/// 界定出来的一条记录。
88#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
89pub struct Record {
90    /// 记录正文:由各行的 `text` 原样拼接(不插入、不删改任何字符)。
91    pub body: String,
92    /// 本记录在来源里的字节区间:首行的 `start_offset` 到末行的 `end_offset`。
93    pub start_offset: u64,
94    pub end_offset: u64,
95    /// 拼成这条记录的行数。
96    pub lines: usize,
97    /// 只有从 [`Delimiter::push`]/[`Delimiter::flush`] 交出来的记录,这个值才有意义。
98    pub completion: Completion,
99}
100
101/// 累积上限:先到者胜。
102///
103/// 同一条记录**不会超过**这两个数(唯一例外见 [`Delimiter::push`]:单行本身超限时,
104/// 那条单行不被切开 —— 单行的上限归读取器管)。上限是"某个块永不结束"的兜底:
105/// 一个没有边界信号的缩进块会一直涨,不给上限就等于把内存交给日志内容决定。
106#[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    /// 两个值都被抬到至少 1:0 会表示"连一行都不许",那不是一条能用的规则。
114    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/// 起点状态:取决于边界信号是**开始型**还是**结束型**。
123#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
124pub enum Start {
125    /// 等第一个**开始信号**才起头 —— 配**开始型**格式(如行首时间戳)。
126    ///
127    /// 这样从记录中间落地时,那一截不会被当成记录发出(不发半条)。
128    WaitForStart,
129    /// 没在累积时,下一行就起头 —— 配**结束型**格式(如空行分隔)。
130    ///
131    /// 这种格式没有可辨的开始信号,所以无法识别"落地在记录中间":第一条可能无头。
132    /// 这是格式本身的代价,不是界定器能补救的。
133    Collect,
134}
135
136/// 记录边界界定器:把行流切成记录。
137///
138/// # 序列化
139///
140/// 可以整体存盘、跨进程接着算(`agentd` 的 checkpoint 里就存着它)。
141/// 存的是**完整状态含策略**([`Limits`] / [`Start`] 一起落)—— 因为它们描述的是
142/// "这条未封口的记录是按什么规则在攒",属于状态而不只是配置,且会随着这条记录结束而自然失效。
143/// 想改用当前配置接续也可以:新建一个,再把 [`Delimiter::pending`] 拿走的那条搬过去。
144#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
145pub struct Delimiter {
146    limits: Limits,
147    start: Start,
148    /// 正在累积的那条;`None` = 手里没有未封口的记录。
149    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    /// 用一份未封口的记录恢复(跨进程接着同一条)。
168    ///
169    /// \(record.completion\) 被忽略 —— 它只在封口那一刻才有意义。计数器从零起:
170    /// 它们是"本次运行"的诊断,不跟着那条记录走。
171    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    /// 送一行进来。若这一行**封上了**上一条记录,返回它。
183    ///
184    /// 一次至多封上一条:一条记录只在**下一个边界信号**处结束。
185    ///
186    /// 到上限时:把已累积的那条以 [`Completion::Oversized`] 封口,**并丢掉这一行**
187    /// (它此刻已经没有头了,与"从记录中间落地"同一处置)。所以产出的记录不会超过上限,
188    /// 唯一例外是本就是**单行**的记录 —— 单行是读取器该管的边界,界定器不替它切半条。
189    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            // 结束信号是**间隔**不是内容:封上当前这条即可。它不算"丢弃" ——
197            // 它是一条正常的边界,不该污染"信号可能配错"的诊断计数。
198            Boundary::Ends => self.seal(Completion::Boundary),
199            Boundary::Neither => {
200                // 先算结论再动状态:把"借用 buffer 判断"与"改 buffer / 封口"分开,
201                // 否则封口与追加不能同时出现。
202                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                    // 落在记录中间、又没有头:丢。宁可丢半条,不可发半条。
224                    None => {
225                        self.drop_line(line);
226                        None
227                    }
228                }
229            }
230        }
231    }
232
233    /// 到期封口:这一段不会再有续行了(空闲到点 / 来源变化 / 停机)。
234    ///
235    /// 不封口的话,最后一条记录会一直悬着 —— 低频文件上就是"日志采了但看不到"。
236    pub fn flush(&mut self) -> Option<Record> {
237        self.seal(Completion::Deadline)
238    }
239
240    /// 手里有没有还没封口的记录。
241    pub fn is_accumulating(&self) -> bool {
242        self.buffer.is_some()
243    }
244
245    /// 未封口那条记录的只读视图(`None` = 手里没有)。
246    ///
247    /// 它的 `completion` 还没定,别当真 —— 要的是"这条攒到哪了"(正文、区间、行数)。
248    pub fn pending(&self) -> Option<&Record> {
249        self.buffer.as_ref()
250    }
251
252    /// 交出未封口的那条记录(`None` = 手里没有),并清空界定器。
253    ///
254    /// 给"把状态存下来、下次接着算"的调用方用 —— 与 [`Delimiter::pending`] 相比
255    /// 它**拿走**了内容,省掉一次正文本的拷贝。
256    pub fn into_pending(mut self) -> Option<Record> {
257        self.buffer.take()
258    }
259
260    /// 累计丢弃的**行数**:**没有头**的续行(落在记录中间落地,或被上限挡在门外)。
261    ///
262    /// 它一直涨却从不产出记录,说明边界信号可能配错了 —— 那不该表现成"这个文件没日志"。
263    pub fn dropped_lines(&self) -> usize {
264        self.dropped_lines
265    }
266
267    /// 累计丢弃的**正文**字节数(按各行的 `text` 长度算,不是来源区间宽度)。
268    ///
269    /// 计数是累计的:要判"**连续**丢了很久",由调用方在每次起头后取快照做差。
270    pub fn dropped_bytes(&self) -> usize {
271        self.dropped_bytes
272    }
273
274    /// 已产出的记录数。
275    pub fn emitted(&self) -> usize {
276        self.emitted
277    }
278
279    /// 起一条新记录。`completion` 只是占位 —— 所有对外交出的记录都在 [`Delimiter::seal`]
280    /// 里被改写,这个值不会流出去。
281    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
304/// 再加这一行会不会越过上限(先到者胜)。**含端点**:正好等于上限是允许的。
305fn 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    /// 带行尾换行的行(与文件读取器的口径一致)。
323    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    /// 开始型信号:记录开头的锚是 `2026-09-23 20:24:24+08 ...`。
332    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    /// `install.log` 的真实形状:两种时间戳锚 + 缩进续行。
341    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    /// 空行是间隔(结束型)的信号。
352    fn blank_separated(line: &str) -> Boundary {
353        if line.trim().is_empty() {
354            Boundary::Ends
355        } else {
356            Boundary::Neither
357        }
358    }
359
360    /// 逐行喂进去,返回所有记录。偏移按输入累加,便于断言区间。
361    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    /// 把行喂进去再到期封口,返回全部记录。
374    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    // ── 基本形状 ──────────────────────────────────────────────────────────────
381
382    #[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        // 只有第一条被封上:第二条要等**下一个**锚。
394        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        // 那条悬着的靠到期封口收尾。
403        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        // 相邻两条锚之间没有续行:各自就是一条一行的记录。
412        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        // 相邻单行记录首尾相接,没有洞。
431        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        // 两种信号混用:行首锚开记录,空行也可作终止符。
437        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        // 空行是间隔不是内容:既不进正文,也不算丢弃。
455        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        // 连着两个空行(或文件以空行开头):没有任何未封口的记录可断。
461        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        // `indented` 语义严格:空行不是续行,它算开始信号 —— 结果是一条正文为空的记录。
472        // 这不是意外,是代理读法本身的局限(要别的行为就自己写判定函数)。
473        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    // ── 起点与"不发半条" ──────────────────────────────────────────────────────
482
483    #[test]
484    fn an_indented_first_line_lands_mid_record_and_is_dropped() {
485        // 从记录中间落地:那一截没有头。宁可丢半条,不可发半条。
486        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        // 什么都没吃到(或只吃到没有头的碎片)就封口:不该凭空产出记录。
510        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        // 一次读到的行不构成完整记录:状态跨批次活着,记录只产出一次(不产生半条)。
524        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        // 结束型格式从第一行就起头:开头那截不是"落在记录中间",不该被丢。
540        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        // 结束型起点不排斥开始型信号:两者可以混用(锚开一条、空行也可断一条)。
560        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    // ── 上限 ─────────────────────────────────────────────────────────────────
595
596    #[test]
597    fn the_line_limit_is_inclusive_and_the_next_line_is_refused() {
598        // 上限**含端点**:max_lines = 3 就该允许 3 行,第 4 行才越限。
599        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        // 正好等于两行之和:第二行要放行;第三行才越限。
621        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        // 一个永不结束的块:到限就封口 + 标记;越限的行不再粘进任何记录。
643        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        // 越限的两行都没有头了 → 丢弃,且不粘进下一条。
662        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        // 上限约束的是**累积**,不是单行。单行的上限归读取器管 ——
672        // 界定器不替它把一行切成半条(半条记录会一路骗过下游解析)。
673        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        // 结束型格式没有锚可以重新对齐:越限之后剩下的续行只能起一条**无头**记录。
691        // 这是 `Start::Collect` 已经声明的代价(它同样无法识别"落地在记录中间"),
692        // 不是界定器能补救的 —— 钉住它,免得日后被当成新 bug 改坏。
693        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    // ── 不变量 ────────────────────────────────────────────────────────────────
712
713    #[test]
714    fn a_delimiter_survives_a_round_trip_through_json() {
715        // 跨进程接着算:整个界定器存盘再读回来,未封口的记录不许变形,计数也不该被清零
716        // (否则重启会把"丢了多少"抹掉,而那正是判断信号配错的依据)。
717        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        // 头一段没有锚 → 丢弃(记数)。
722        assert!(delimiter.push(line(junk, 0), anchor).is_none());
723        // 再开一条未封口的,等着跨进程继续。
724        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        // 接着喂:封口的结果与"从没断过"完全一样。
743        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        // 记录边界必须完整覆盖输入:不能丢内容,也不能让同一段字节落在两条里。
758        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        // 被丢弃的是开头那截(已知),其余必须首尾相接。
769        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        // 保序:偏移严格递增。
793        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    /// 确定性伪随机(xorshift):不引新依赖,但每次跑的是同一串序列。
808    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        // 核心不变量:输入字节流要么落在某条记录里,要么被计入丢弃 —— 没有第三种去处。
826        // 混三种形状、随机切批次、用**紧**上限,把跨批次 / 超限 / 丢弃 / 起头几条路一起压。
827        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        // 以一条锚收尾:保证流尾手里确实有一条未封口的记录,把"到期封口"也压上
842        // (否则流尾刚好吃到越限、被丢空,那条路就白测了)。
843        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        // 自证:这份输入必须真的走到过丢弃/超限/到期这几条路,否则这个性质测试是空转的。
866        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    // ── 现成读法与真实形状 ────────────────────────────────────────────────────
905
906    #[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        // 语义严格、无例外:空行也算开始信号(它不是续行)。
912        assert_eq!(indented("\n"), Boundary::Starts);
913    }
914
915    #[test]
916    fn an_install_log_shaped_stream_folds_by_its_two_anchors() {
917        // `/var/log/install.log` 的真实形状:两种时间戳锚(带年/时区 与 BSD)+ 缩进续行。
918        // 缩进读法在这个文件上能蒙对,但把锚写准才是不依赖运气的做法。
919        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    // ── 契约、状态交接与补充不变量 ──────────────────────────────────────────
950
951    #[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 只是占位:resume / flush 会按封口方式改写它。
1007            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        // 枚举用外部标签,名字就是契约(checkpoint 存盘依赖它)。
1026        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        // 未知字段被忽略;缺字段报错(Record 没有默认值)。
1044        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    /// 混合开始型 / 结束型信号时的守恒律:每个输入字节要么在某条记录里,
1057    /// 要么被计入丢弃,要么属于结束信号本身(它是间隔,不算内容)。
1058    #[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", // 结束信号
1065        ];
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        // 以结束信号收尾,逼出"到期封口"。
1085        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        // 自证走到了各条路径,否则这个性质测试是空转的。
1102        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        // 换成结束型起点再验一次守恒:此时除了超限外不应再有"没有头"的丢弃。
1135        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}