kumo_jsonl/
decompress.rs

1use std::collections::VecDeque;
2use std::io::{BufRead, BufReader};
3use thiserror::Error;
4use zstd_safe::{DCtx, InBuffer, OutBuffer};
5
6#[derive(Error, Debug)]
7#[error("{}", zstd_safe::get_error_name(*.0))]
8pub struct ZStdError(pub usize);
9
10/// Default limit, in bytes, on how large a decompressed record may be. A record
11/// that reaches it without a newline is discarded and the following records are
12/// still read. Large enough that legitimate records are not expected to
13/// approach it, while still bounding how much is buffered for one line.
14pub const DEFAULT_MAX_LINE_SIZE: usize = 128 * 1024 * 1024;
15
16/// A line extracted from the decompressed stream, along with its
17/// byte offset in the decompressed data.
18pub struct DecompressedLine {
19    pub text: String,
20    pub byte_offset: u64,
21}
22
23impl std::fmt::Debug for DecompressedLine {
24    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
25        // A line may be as large as the configured maximum (128 MiB by
26        // default). Truncate to a leading fragment to keep debug output, and
27        // anything that logs it, bounded regardless of line length.
28        const MAX: usize = 64;
29        let mut dbg = f.debug_struct("DecompressedLine");
30        dbg.field("byte_offset", &self.byte_offset)
31            .field("len", &self.text.len());
32        if self.text.len() <= MAX {
33            dbg.field("text", &self.text);
34        } else {
35            dbg.field(
36                "text_prefix",
37                &&self.text[..self.text.floor_char_boundary(MAX)],
38            );
39        }
40        dbg.finish()
41    }
42}
43
44/// The outcome of a call to [`FileDecompressor::next_line`].
45#[derive(Debug)]
46pub enum NextLine {
47    /// A complete line was extracted.
48    Line(DecompressedLine),
49    /// A record longer than the configured maximum was discarded and the
50    /// stream continues with the following record. Reports where the discarded
51    /// record began in the decompressed stream and how many bytes were dropped
52    /// (excluding the terminating newline).
53    Skipped { byte_offset: u64, bytes: u64 },
54    /// No line is available right now. The caller should check whether the file
55    /// is done or retry later.
56    None,
57}
58
59/// Bookkeeping for a record being discarded because it exceeds the maximum line
60/// size. Scanning continues across decompress calls until the record's
61/// terminating newline is found. Discarding only consumes bytes from the
62/// output buffer; the zstd stream itself is never put in an error state, so
63/// the segment's later records decode normally once the newline is reached.
64struct Skip {
65    /// Byte offset in the decompressed stream where the discarded record began.
66    byte_offset: u64,
67    /// Bytes of the record discarded so far, excluding the terminating newline.
68    discarded: u64,
69    /// Whether the record is at or after `skip_before` and should be reported
70    /// to the caller. A record before the checkpoint was consumed on an earlier
71    /// run and is dropped silently.
72    surface: bool,
73}
74
75/// State for incremental zstd decompression and line extraction from a single file.
76pub struct FileDecompressor {
77    file: BufReader<std::fs::File>,
78    context: DCtx<'static>,
79    out_buffer: Vec<u8>,
80    /// Steady-state size of `out_buffer`. The buffer grows past this to hold an
81    /// oversized record and is shrunk back to it once the record is emitted.
82    base_out_buffer: usize,
83    /// The maximum size `out_buffer` may grow to. A record that fills it
84    /// without a newline is discarded (reported as `NextLine::Skipped`) rather
85    /// than buffered further.
86    max_out_buffer: usize,
87    /// Start of the next unprocessed line in `out_buffer`.
88    line_start: usize,
89    /// Number of valid bytes in `out_buffer`.
90    out_pos: usize,
91    /// When discarding a record that exceeds `max_out_buffer`, the state of
92    /// that in-progress skip. `None` during normal line extraction.
93    skipping: Option<Skip>,
94    /// Total number of lines decompressed so far.
95    lines_decompressed: usize,
96    /// Global line index up to which lines have been consumed or skipped.
97    /// This is the value that should be used for checkpointing.
98    /// Equals skip_before + number of lines actually returned to caller.
99    pub lines_consumed: usize,
100    /// Buffered lines that have been extracted but not yet consumed.
101    pending_lines: VecDeque<DecompressedLine>,
102    /// Whether we've seen EOF on the compressed input.
103    saw_eof: bool,
104    /// Cumulative byte offset in the decompressed stream.
105    /// Tracks the position of `line_start` relative to the start
106    /// of the decompressed output.
107    decompressed_offset: u64,
108}
109
110impl FileDecompressor {
111    /// Open a file and prepare for incremental zstd decompression.
112    pub fn open(path: &std::path::Path) -> anyhow::Result<Self> {
113        Self::open_with_max_line_size(path, DEFAULT_MAX_LINE_SIZE)
114    }
115
116    /// Open a file, capping the output buffer at `max_out_buffer` bytes. A
117    /// record that fills the cap without a newline is discarded and reading
118    /// continues with the next record (see [`NextLine::Skipped`]).
119    pub fn open_with_max_line_size(
120        path: &std::path::Path,
121        max_out_buffer: usize,
122    ) -> anyhow::Result<Self> {
123        let file = BufReader::new(
124            std::fs::File::open(path)
125                .map_err(|e| anyhow::anyhow!("opening {} for read: {e}", path.display()))?,
126        );
127        let mut context = DCtx::create();
128        context
129            .init()
130            .map_err(ZStdError)
131            .map_err(|e| anyhow::anyhow!("initialize zstd decompression context: {e}"))?;
132        context
133            .load_dictionary(&[])
134            .map_err(ZStdError)
135            .map_err(|e| anyhow::anyhow!("load empty dictionary: {e}"))?;
136
137        // Keep at least one byte so the growth check has a non-empty buffer to
138        // fill. A misconfigured tiny cap simply discards records as oversized.
139        let max_out_buffer = max_out_buffer.max(1);
140        let base_out_buffer = DCtx::out_size().min(max_out_buffer);
141        Ok(Self {
142            file,
143            context,
144            out_buffer: vec![0u8; base_out_buffer],
145            base_out_buffer,
146            max_out_buffer,
147            line_start: 0,
148            out_pos: 0,
149            skipping: None,
150            lines_decompressed: 0,
151            lines_consumed: 0,
152            pending_lines: VecDeque::new(),
153            saw_eof: false,
154            decompressed_offset: 0,
155        })
156    }
157
158    /// Returns the next line from this file.
159    ///
160    /// `skip_before`: lines with index < skip_before are discarded.
161    ///
162    /// Returns:
163    /// - `Ok(NextLine::Line(line))` -- a complete line was extracted.
164    /// - `Ok(NextLine::Skipped { .. })` -- a record exceeding the maximum line
165    ///   size was discarded. Reading continues with the next record.
166    /// - `Ok(NextLine::None)` -- there isn't any more data available right now.
167    ///   The caller should check if the file is done or retry later.
168    pub fn next_line(&mut self, skip_before: usize) -> anyhow::Result<NextLine> {
169        // Return a buffered line if available
170        if let Some(line) = self.pending_lines.pop_front() {
171            self.lines_consumed += 1;
172            return Ok(NextLine::Line(line));
173        }
174
175        // If we previously saw EOF and have no buffered lines, signal EOF
176        if self.saw_eof {
177            return Ok(NextLine::None);
178        }
179
180        // Account for skipped lines in lines_consumed
181        if self.lines_consumed < skip_before {
182            self.lines_consumed = skip_before;
183        }
184
185        // Read and decompress more data
186        loop {
187            let in_buffer = self.file.fill_buf()?;
188            if in_buffer.is_empty() {
189                self.saw_eof = true;
190                // A skip in progress is left intact. If the file is still
191                // being written, the caller resets EOF and we resume discarding
192                // when more data arrives. If the file is done, has_partial_data
193                // returns true while a skip is in progress, which is how the
194                // caller recognizes that the trailing partial record should be
195                // dropped.
196                if let Some(line) = self.pending_lines.pop_front() {
197                    self.lines_consumed += 1;
198                    return Ok(NextLine::Line(line));
199                }
200                return Ok(NextLine::None);
201            }
202
203            let mut src = InBuffer::around(in_buffer);
204            let mut dest = OutBuffer::around_pos(&mut self.out_buffer, self.out_pos);
205
206            self.context
207                .decompress_stream(&mut dest, &mut src)
208                .map_err(ZStdError)
209                .map_err(|e| anyhow::anyhow!("zstd decompress: {e}"))?;
210
211            let bytes_read = {
212                let pos = src.pos();
213                drop(src);
214                pos
215            };
216            self.file.consume(bytes_read);
217            self.out_pos = dest.pos();
218
219            // Set when this iteration finishes discarding an oversized record
220            // that should be reported. Extraction below queues any lines found
221            // after the skip into pending_lines first. The code further down
222            // deliberately checks for a pending just_skipped report ahead of
223            // pending_lines, reversing the detection order: the skip is
224            // reported to the caller before the lines that follow it, even
225            // though those lines were extracted first.
226            let mut just_skipped: Option<(u64, u64)> = None;
227
228            // While discarding an oversized record, consume its bytes up to and
229            // including its terminating newline before resuming extraction.
230            if self.skipping.is_some() {
231                match memchr::memchr(b'\n', &self.out_buffer[..self.out_pos]) {
232                    None => {
233                        // The whole buffer is more of the record. Drop it and
234                        // read more. `self.skipping.is_some()` guards this
235                        // whole match, and the code that enlarges `out_buffer`
236                        // runs only when `self.skipping` is `None`. The
237                        // conditions cannot both hold, so entering this branch
238                        // guarantees the resize code does not run this
239                        // iteration. `out_buffer` holds at `base_out_buffer`
240                        // bytes for every iteration of a skip, whatever the
241                        // size of the discarded record.
242                        let skip = self.skipping.as_mut().expect("checked skipping");
243                        skip.discarded += self.out_pos as u64;
244                        self.decompressed_offset += self.out_pos as u64;
245                        self.out_pos = 0;
246                        self.line_start = 0;
247                        continue;
248                    }
249                    Some(idx) => {
250                        // The record ends at the newline. Discard through it
251                        // and resume normal extraction on whatever follows.
252                        let mut skip = self.skipping.take().expect("checked skipping");
253                        skip.discarded += idx as u64;
254                        let consumed = idx + 1;
255                        self.decompressed_offset += consumed as u64;
256                        self.lines_decompressed += 1;
257                        self.out_buffer.copy_within(consumed..self.out_pos, 0);
258                        self.out_pos -= consumed;
259                        self.line_start = 0;
260                        if skip.surface {
261                            just_skipped = Some((skip.byte_offset, skip.discarded));
262                        }
263                        // Fall through to extract the remainder. A surfaced
264                        // skip is returned below, ahead of those lines.
265                    }
266                }
267            }
268
269            // Extract complete lines
270            while let Some(idx) =
271                memchr::memchr(b'\n', &self.out_buffer[self.line_start..self.out_pos])
272            {
273                let line_byte_offset = self.decompressed_offset;
274                if self.lines_decompressed >= skip_before {
275                    let this_line = &self.out_buffer[self.line_start..self.line_start + idx];
276                    let line = String::from_utf8_lossy(this_line).into_owned();
277                    self.pending_lines.push_back(DecompressedLine {
278                        text: line,
279                        byte_offset: line_byte_offset,
280                    });
281                }
282                // Advance past the line content + newline
283                let consumed = idx + 1;
284                self.decompressed_offset += consumed as u64;
285                self.line_start += consumed;
286                self.lines_decompressed += 1;
287            }
288
289            // Compact the output buffer
290            if self.line_start == self.out_pos {
291                self.out_pos = 0;
292                self.line_start = 0;
293            } else if self.line_start > 0 {
294                self.out_buffer
295                    .copy_within(self.line_start..self.out_pos, 0);
296                self.out_pos -= self.line_start;
297                self.line_start = 0;
298            }
299
300            // Release memory borrowed to hold an oversized record (see
301            // base_out_buffer) once the remaining live data fits in the
302            // steady-state size again.
303            if self.out_buffer.len() > self.base_out_buffer && self.out_pos <= self.base_out_buffer
304            {
305                self.out_buffer.truncate(self.base_out_buffer);
306                self.out_buffer.shrink_to_fit();
307            }
308
309            // Report a just-finished oversized-record skip ahead of any lines
310            // that followed it in the same buffer (already queued above).
311            if let Some((byte_offset, bytes)) = just_skipped {
312                self.lines_consumed += 1;
313                return Ok(NextLine::Skipped { byte_offset, bytes });
314            }
315
316            // If we extracted any lines, return the first one
317            if let Some(line) = self.pending_lines.pop_front() {
318                self.lines_consumed += 1;
319                return Ok(NextLine::Line(line));
320            }
321
322            // A full buffer that doesn't hold any complete lines (compaction
323            // always leaves line_start at 0) means the current record is larger
324            // than the buffer. Grow it (up to max_out_buffer) to let the next
325            // decompress call make progress; otherwise zstd stalls with "no
326            // progress ... output buffer full".
327            if self.out_pos == self.out_buffer.len() {
328                if self.out_buffer.len() >= self.max_out_buffer {
329                    // If a record exceeds the buffer cap, we discard just that
330                    // record rather than failing the whole segment, since the
331                    // zstd stream is intact (this is our own cap, not a decode
332                    // error). We scan forward for its terminating newline and
333                    // resume with the next record. The current buffer is
334                    // dropped and the buffer shrinks back to base, since while
335                    // skipping we only need room to scan.
336                    let surface = self.lines_decompressed >= skip_before;
337                    self.skipping = Some(Skip {
338                        byte_offset: self.decompressed_offset,
339                        discarded: self.out_pos as u64,
340                        surface,
341                    });
342                    self.decompressed_offset += self.out_pos as u64;
343                    self.out_pos = 0;
344                    self.line_start = 0;
345                    if self.out_buffer.len() > self.base_out_buffer {
346                        self.out_buffer.truncate(self.base_out_buffer);
347                        self.out_buffer.shrink_to_fit();
348                    }
349                    continue;
350                }
351                let new_len = self
352                    .out_buffer
353                    .len()
354                    .saturating_mul(2)
355                    .min(self.max_out_buffer);
356                // Reserve fallibly: max_out_buffer comes straight from operator
357                // config and can be arbitrarily large. A huge cap is reported
358                // as an error here (skipping this segment) instead of aborting
359                // the process on a failed allocation.
360                let additional = new_len - self.out_buffer.len();
361                self.out_buffer.try_reserve(additional).map_err(|e| {
362                    anyhow::anyhow!("allocating {new_len} bytes for the decompression buffer: {e}")
363                })?;
364                self.out_buffer.resize(new_len, 0);
365            }
366
367            // No complete lines yet; read more data
368        }
369    }
370
371    /// Reset the EOF flag so we can try reading more data
372    /// (useful when tailing a file that is still being written to).
373    pub fn reset_eof(&mut self) {
374        self.saw_eof = false;
375    }
376
377    /// Returns true if there is partial (incomplete line) data remaining: either
378    /// buffered bytes with no terminating newline yet, or a record still being
379    /// discarded for exceeding the maximum line size. When a done file ends in
380    /// this state the trailing record is incomplete and is dropped.
381    pub fn has_partial_data(&self) -> bool {
382        self.out_pos > 0 || self.skipping.is_some()
383    }
384
385    /// Returns true if the trailing partial data is a record being discarded
386    /// for exceeding the maximum line size, as opposed to an unterminated line
387    /// left behind by a writer that exited without flushing.
388    pub fn is_discarding_oversized_record(&self) -> bool {
389        self.skipping.is_some()
390    }
391}
392
393#[cfg(test)]
394mod test {
395    use super::*;
396    use std::io::Write;
397    use zstd::stream::write::Encoder;
398
399    /// Compress the supplied lines into a zstd JSONL segment on disk, matching
400    /// the framing produced by the log writer, and return the path.
401    fn write_segment(lines: &[String]) -> tempfile::NamedTempFile {
402        let file = tempfile::NamedTempFile::new().unwrap();
403        let mut encoder = Encoder::new(file.reopen().unwrap(), 3).unwrap();
404        for line in lines {
405            encoder.write_all(line.as_bytes()).unwrap();
406            encoder.write_all(b"\n").unwrap();
407        }
408        encoder.finish().unwrap();
409        file
410    }
411
412    fn read_all(path: &std::path::Path) -> anyhow::Result<Vec<String>> {
413        let mut decompressor = FileDecompressor::open(path)?;
414        let mut lines = vec![];
415        loop {
416            match decompressor.next_line(0)? {
417                NextLine::Line(line) => lines.push(line.text),
418                NextLine::Skipped { .. } => panic!("unexpected skip"),
419                NextLine::None => break,
420            }
421        }
422        Ok(lines)
423    }
424
425    fn expect_line(decompressor: &mut FileDecompressor) -> DecompressedLine {
426        match decompressor.next_line(0).unwrap() {
427            NextLine::Line(line) => line,
428            other => panic!("expected a line, got {other:?}"),
429        }
430    }
431
432    #[test]
433    fn small_records_roundtrip() {
434        let input: Vec<String> = (0..1000).map(|i| format!("{{\"n\":{i}}}")).collect();
435        let seg = write_segment(&input);
436        let lines = read_all(seg.path()).unwrap();
437        k9::assert_equal!(lines, input);
438    }
439
440    /// Verifies that a record whose decompressed size exceeds the fixed output
441    /// buffer (`DCtx::out_size()`) is still recovered, along with the ordinary
442    /// records that share its segment.
443    #[test]
444    fn record_larger_than_output_buffer() {
445        let big_value = "x".repeat(DCtx::out_size() * 2);
446        let input = vec![
447            r#"{"n":"before"}"#.to_string(),
448            format!(r#"{{"big":"{big_value}"}}"#),
449            r#"{"n":"after"}"#.to_string(),
450        ];
451        let seg = write_segment(&input);
452        let lines = read_all(seg.path()).unwrap();
453        k9::assert_equal!(lines, input);
454    }
455
456    /// After an oversized record forces the buffer to grow, it shrinks back
457    /// down to the steady-state size once that record has been emitted.
458    #[test]
459    fn buffer_shrinks_back_after_large_record() {
460        let big_value = "x".repeat(DCtx::out_size() * 4);
461        let input = vec![
462            format!(r#"{{"big":"{big_value}"}}"#),
463            r#"{"n":"after"}"#.to_string(),
464        ];
465        let seg = write_segment(&input);
466        let mut decompressor = FileDecompressor::open(seg.path()).unwrap();
467
468        // Recovering the oversized record forces the buffer past its base size
469        // (proven by record_larger_than_output_buffer). By the time it is
470        // emitted the buffer has already been reclaimed.
471        let first = expect_line(&mut decompressor);
472        k9::assert_equal!(first.text, input[0]);
473        k9::assert_equal!(decompressor.out_buffer.len(), decompressor.base_out_buffer);
474
475        let second = expect_line(&mut decompressor);
476        k9::assert_equal!(second.text, input[1]);
477        k9::assert_equal!(decompressor.out_buffer.len(), decompressor.base_out_buffer);
478    }
479
480    /// Discards a record exceeding the cap and reports it as skipped, with its
481    /// start offset and byte count, rather than failing the whole segment.
482    #[test]
483    fn record_exceeding_cap_is_skipped() {
484        let cap = 64 * 1024;
485        let content = "x".repeat(cap * 2);
486        let seg = write_segment(&[content.clone()]);
487        let mut decompressor = FileDecompressor::open_with_max_line_size(seg.path(), cap).unwrap();
488        match decompressor.next_line(0).unwrap() {
489            NextLine::Skipped { byte_offset, bytes } => {
490                k9::assert_equal!(byte_offset, 0);
491                k9::assert_equal!(bytes, content.len() as u64);
492            }
493            other => panic!("expected the oversized record to be skipped, got {other:?}"),
494        }
495        assert!(matches!(decompressor.next_line(0).unwrap(), NextLine::None));
496    }
497
498    /// Reads a record one byte below the cap: its content plus the newline
499    /// separator fit within the buffer.
500    #[test]
501    fn record_just_below_cap_is_read() {
502        let cap = 64 * 1024;
503        let content = "x".repeat(cap - 1);
504        let seg = write_segment(&[content.clone()]);
505        let mut decompressor = FileDecompressor::open_with_max_line_size(seg.path(), cap).unwrap();
506        k9::assert_equal!(expect_line(&mut decompressor).text, content);
507    }
508
509    /// Skips a record whose content is exactly `cap` bytes: one byte too long
510    /// for its trailing newline to also fit in the buffer.
511    #[test]
512    fn record_at_cap_is_skipped() {
513        let cap = 64 * 1024;
514        let content = "y".repeat(cap);
515        let seg = write_segment(&[content]);
516        let mut decompressor = FileDecompressor::open_with_max_line_size(seg.path(), cap).unwrap();
517        assert!(matches!(
518            decompressor.next_line(0).unwrap(),
519            NextLine::Skipped { .. }
520        ));
521    }
522
523    /// Discards an oversized record between two ordinary records while the
524    /// records around it are still read, and the checkpoint line count keeps
525    /// advancing across the skip.
526    #[test]
527    fn oversized_record_skipped_rest_recovered() {
528        let cap = 64 * 1024;
529        let input = vec![
530            "before".to_string(),
531            "y".repeat(cap * 2),
532            "after".to_string(),
533        ];
534        let seg = write_segment(&input);
535        let mut decompressor = FileDecompressor::open_with_max_line_size(seg.path(), cap).unwrap();
536
537        let before = expect_line(&mut decompressor);
538        k9::assert_equal!(before.text, "before".to_string());
539
540        match decompressor.next_line(0).unwrap() {
541            NextLine::Skipped { bytes, .. } => {
542                k9::assert_equal!(bytes, (cap * 2) as u64);
543            }
544            other => panic!("expected the oversized record to be skipped, got {other:?}"),
545        }
546
547        let after = expect_line(&mut decompressor);
548        k9::assert_equal!(after.text, "after".to_string());
549        // before + skipped record + after = three lines consumed.
550        k9::assert_equal!(decompressor.lines_consumed, 3);
551
552        assert!(matches!(decompressor.next_line(0).unwrap(), NextLine::None));
553    }
554
555    /// A record that exceeds the cap but sits before `skip_before` (already
556    /// consumed on an earlier run) is dropped silently. Decompression always
557    /// restarts from the head of a segment. This keeps an oversized record
558    /// before the checkpoint from being reported again on every restart.
559    #[test]
560    fn oversized_record_before_checkpoint_dropped_silently() {
561        let cap = 64 * 1024;
562        let input = vec!["y".repeat(cap * 2), "after".to_string()];
563        let seg = write_segment(&input);
564        let mut decompressor = FileDecompressor::open_with_max_line_size(seg.path(), cap).unwrap();
565
566        // skip_before = 1 marks the oversized first record as already consumed;
567        // it must not surface as a skip, and the next call yields "after".
568        match decompressor.next_line(1).unwrap() {
569            NextLine::Line(line) => {
570                k9::assert_equal!(line.text, "after".to_string());
571            }
572            other => panic!("expected the record after the skipped one, got {other:?}"),
573        }
574        assert!(matches!(decompressor.next_line(1).unwrap(), NextLine::None));
575    }
576
577    /// Verifies that a skip reaching EOF before the newline of the record
578    /// arrives (the record is still being written) is preserved: after more
579    /// data is appended and `reset_eof` is called, the skip completes with the
580    /// full byte count and the following record is read with its correct
581    /// offset, rather than the appended suffix being mistaken for a new record.
582    #[test]
583    fn skip_resumes_across_eof_when_tailing() {
584        let cap = 64 * 1024;
585        let content = "y".repeat(cap * 2);
586
587        // Compress incrementally: first the content of the oversized record
588        // with no newline yet (flushed to decode on its own), then its newline
589        // and a following record.
590        let mut encoder = Encoder::new(Vec::new(), 3).unwrap();
591        encoder.write_all(content.as_bytes()).unwrap();
592        encoder.flush().unwrap();
593        let prefix_len = encoder.get_ref().len();
594        encoder.write_all(b"\nafter\n").unwrap();
595        let all = encoder.finish().unwrap();
596        let (prefix, suffix) = all.split_at(prefix_len);
597
598        let seg = tempfile::NamedTempFile::new().unwrap();
599        std::fs::write(seg.path(), prefix).unwrap();
600
601        let mut d = FileDecompressor::open_with_max_line_size(seg.path(), cap).unwrap();
602        // The record is missing its terminating newline: reading scans for one
603        // and finds none before the input runs out, then stops at EOF mid-skip.
604        assert!(matches!(d.next_line(0).unwrap(), NextLine::None));
605        assert!(d.has_partial_data());
606
607        // The rest of the record and a following record arrive.
608        let mut f = std::fs::OpenOptions::new()
609            .append(true)
610            .open(seg.path())
611            .unwrap();
612        f.write_all(suffix).unwrap();
613        f.flush().unwrap();
614        drop(f);
615        d.reset_eof();
616
617        match d.next_line(0).unwrap() {
618            NextLine::Skipped { byte_offset, bytes } => {
619                k9::assert_equal!(byte_offset, 0);
620                k9::assert_equal!(bytes, content.len() as u64);
621            }
622            other => panic!("expected the oversized record to be skipped, got {other:?}"),
623        }
624        let after = expect_line(&mut d);
625        k9::assert_equal!(after.text, "after".to_string());
626        k9::assert_equal!(after.byte_offset, content.len() as u64 + 1);
627        k9::assert_equal!(d.lines_consumed, 2);
628        assert!(matches!(d.next_line(0).unwrap(), NextLine::None));
629    }
630}