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}