Skip to main content

datui_lib/
framed_records.rs

1//! Records that are not all one size, read from a memory map: length-prefixed records,
2//! variants a type field picks, records found by a sync marker, records in compressed
3//! blocks or in the payloads of a packet capture, and fixed records with fields that
4//! are read one at a time (varints, NUL-terminated text, counted groups, bit fields,
5//! deltas, checksums).
6//!
7//! A first pass walks the records and keeps where every 1024th one starts (and, for
8//! delta columns, the running sums there), so the index stays small and a window
9//! anywhere is read by walking from the nearest checkpoint. Fixed records with a
10//! known size need no pass at all: where a record starts is arithmetic.
11//!
12//! A compressed block is decompressed when it is first read, and a few are kept.
13//!
14//! The frame is decoded over a row index ([`crate::row_index`]), as
15//! [`crate::fixed_records`]' is, and a window deeper in the file is read through
16//! [`FramedRecords::window`].
17
18use crate::fixed_records::{Bytes, ColumnLayout, Physical};
19use crate::formats::{
20    self, Amount, Codec, Compression, Delta, Encoding, Field, Framing, HeaderValues, Spec, Type,
21};
22use polars::prelude::*;
23use std::collections::VecDeque;
24use std::ops::Range;
25use std::sync::{Arc, Mutex};
26
27/// Records between checkpoints of the index.
28const CHECKPOINT: u64 = 1024;
29
30/// The most bytes one block may decompress to.
31pub const MAX_BLOCK: usize = 256 << 20;
32
33/// Decompressed blocks kept for reading again.
34const CACHED_BLOCKS: usize = 8;
35
36/// Items one group may hold, and values one counted field may hold.
37const MAX_ITEMS: u64 = 1 << 20;
38
39/// The most rows a table holds: Polars counts rows in 32 bits.
40const MAX_ROWS: usize = IdxSize::MAX as usize;
41
42/// A size known now, or read from an earlier field of the same record.
43#[derive(Debug, Clone, Copy)]
44enum SizeRef {
45    Given(usize),
46    Slot { slot: usize, adjust: i64 },
47    Rest,
48}
49
50/// How an integer's raw value is read, for fields others refer to.
51#[derive(Debug, Clone, Copy)]
52struct IntRead {
53    signed: bool,
54    big: bool,
55}
56
57/// One field of a record (or a block header, a payload header, a group item).
58#[derive(Debug, Clone)]
59struct FieldPlan {
60    name: String,
61    slot: usize,
62    kind: Kind,
63    /// Values per cell: `None` for one; a size for an Array (given) or a List.
64    count: Option<SizeRef>,
65    /// The columns its values go to: one, or one per value when flattened.
66    outs: Vec<usize>,
67    /// Bit fields: column, lowest bit, width.
68    bits: Vec<(usize, u32, u32)>,
69    delta: Delta,
70    /// The index of its running sum, for a delta.
71    delta_index: usize,
72    /// The stored value that means null, for a delta (checked before summing).
73    sentinel: Option<i128>,
74}
75
76#[derive(Debug, Clone)]
77enum Kind {
78    /// A fixed-width value: an integer, a float, a bool, or text and bytes of a size
79    /// known before the record is read.
80    Fixed {
81        width: usize,
82        int: Option<IntRead>,
83    },
84    /// Text or bytes whose size is read from the record. `encoding` is `None` for bytes.
85    Sized {
86        size: SizeRef,
87        encoding: Option<Encoding>,
88    },
89    /// Text to its NUL, within `max` bytes when given (which it then always takes).
90    Strz {
91        max: Option<SizeRef>,
92        encoding: Encoding,
93    },
94    Var {
95        signed: bool,
96    },
97    Pad {
98        size: SizeRef,
99    },
100    /// An offset into a section, where NUL-terminated text is.
101    StringAt {
102        width: usize,
103        big: bool,
104        section: Range<usize>,
105    },
106    Group {
107        items: Vec<FieldPlan>,
108        slots: usize,
109    },
110}
111
112/// One variant of the records.
113#[derive(Debug, Clone)]
114struct VariantPlan {
115    name: Arc<str>,
116    ints: Vec<i128>,
117    texts: Vec<String>,
118    fields: Vec<FieldPlan>,
119    size: Option<SizeRef>,
120}
121
122/// What a column collects while records are read.
123#[derive(Debug, Clone)]
124enum Sink {
125    /// Fixed-width cells, decoded together by the reader of fixed records.
126    Packed {
127        layout: ColumnLayout,
128        cell: usize,
129        buf: Vec<u8>,
130        valid: Vec<bool>,
131    },
132    Text(Vec<Option<String>>),
133    Binary(Vec<Option<Vec<u8>>>),
134    Bits {
135        width: u32,
136        labels: Option<Arc<std::collections::BTreeMap<i64, String>>>,
137        values: Vec<Option<u64>>,
138    },
139    Flag(Vec<Option<bool>>),
140    Label(Vec<Option<Arc<str>>>),
141    Time(Vec<Option<i64>>),
142    List {
143        inner: Box<Sink>,
144        offsets: Vec<i64>,
145        valid: Vec<bool>,
146    },
147    Struct {
148        names: Vec<PlSmallStr>,
149        fields: Vec<Sink>,
150        len: usize,
151    },
152    /// A column the read does not want.
153    Skip,
154}
155
156impl Sink {
157    fn push_null(&mut self) {
158        match self {
159            Self::Packed {
160                cell, buf, valid, ..
161            } => {
162                buf.resize(buf.len() + *cell, 0);
163                valid.push(false);
164            }
165            Self::Text(v) => v.push(None),
166            Self::Label(v) => v.push(None),
167            Self::Binary(v) => v.push(None),
168            Self::Bits { values, .. } => values.push(None),
169            Self::Flag(v) => v.push(None),
170            Self::Time(v) => v.push(None),
171            Self::List { offsets, valid, .. } => {
172                offsets.push(*offsets.last().unwrap_or(&0));
173                valid.push(false);
174            }
175            Self::Struct { fields, len, .. } => {
176                for f in fields {
177                    f.push_null();
178                }
179                *len += 1;
180            }
181            Self::Skip => {}
182        }
183    }
184
185    fn push_bytes(&mut self, bytes: &[u8]) {
186        if let Self::Packed { buf, valid, .. } = self {
187            buf.extend_from_slice(bytes);
188            valid.push(true);
189        }
190    }
191
192    fn finish(self, name: PlSmallStr) -> PolarsResult<Series> {
193        Ok(match self {
194            Self::Packed {
195                layout,
196                cell,
197                buf,
198                valid,
199            } => {
200                let rows = valid.len();
201                // The cells sit side by side in `buf`, whatever the record's stride was:
202                // a delta's running sum is wider than the value stored.
203                let layout = ColumnLayout {
204                    start: 0,
205                    stride: cell,
206                    ..layout
207                };
208                let column = crate::fixed_records::decode(&buf, &layout, rows)?;
209                let series = column
210                    .as_materialized_series()
211                    .clone()
212                    .with_name(name.clone());
213                if valid.iter().all(|v| *v) {
214                    series
215                } else {
216                    let mask: BooleanChunked = valid.into_iter().collect();
217                    let nulls = Series::full_null(PlSmallStr::EMPTY, rows, series.dtype());
218                    series.zip_with(&mask, &nulls)?.with_name(name)
219                }
220            }
221            Self::Text(v) => StringChunked::from_iter_options(name, v.into_iter()).into_series(),
222            Self::Label(v) => {
223                // Few labels, many rows: each label is cast once and the rows gathered,
224                // where casting every row's text took most of a variant column's time.
225                let mut distinct: Vec<Arc<str>> = Vec::new();
226                let mut by_text = std::collections::HashMap::new();
227                let codes: IdxCa = v
228                    .iter()
229                    .map(|label| {
230                        let label = label.as_ref()?;
231                        // A variant's name is one shared string.
232                        if let Some(i) =
233                            distinct.iter().take(16).position(|d| Arc::ptr_eq(d, label))
234                        {
235                            return Some(i as IdxSize);
236                        }
237                        Some(*by_text.entry(label.clone()).or_insert_with(|| {
238                            distinct.push(label.clone());
239                            (distinct.len() - 1) as IdxSize
240                        }))
241                    })
242                    .collect();
243                StringChunked::from_iter_values(name, distinct.iter().map(|d| &**d))
244                    .into_series()
245                    .cast(&DataType::from_categories(Categories::global()))?
246                    .take(&codes)?
247            }
248            Self::Binary(v) => BinaryChunked::from_iter_options(name, v.into_iter()).into_series(),
249            Self::Bits {
250                width,
251                labels,
252                values,
253            } => match (labels, width) {
254                (Some(labels), _) => StringChunked::from_iter_options(
255                    name,
256                    values.into_iter().map(|v| {
257                        v.map(|v| {
258                            i64::try_from(v)
259                                .ok()
260                                .and_then(|k| labels.get(&k).cloned())
261                                .unwrap_or_else(|| v.to_string())
262                        })
263                    }),
264                )
265                .into_series(),
266                (None, 1) => BooleanChunked::from_iter_options(
267                    name,
268                    values.into_iter().map(|v| v.map(|v| v != 0)),
269                )
270                .into_series(),
271                (None, w) => {
272                    let wide =
273                        UInt64Chunked::from_iter_options(name, values.into_iter()).into_series();
274                    let dtype = match w {
275                        2..=8 => DataType::UInt8,
276                        9..=16 => DataType::UInt16,
277                        17..=32 => DataType::UInt32,
278                        _ => DataType::UInt64,
279                    };
280                    wide.strict_cast(&dtype)?
281                }
282            },
283            Self::Flag(v) => BooleanChunked::from_iter_options(name, v.into_iter()).into_series(),
284            Self::Time(v) => Int64Chunked::from_iter_options(name, v.into_iter())
285                .into_datetime(TimeUnit::Nanoseconds, None)
286                .into_series(),
287            Self::List {
288                inner,
289                offsets,
290                valid,
291            } => {
292                let values = inner.finish(PlSmallStr::from_static("item"))?.rechunk();
293                list_series(name, values, offsets, valid)?
294            }
295            Self::Struct { names, fields, len } => {
296                let series = names
297                    .into_iter()
298                    .zip(fields)
299                    .map(|(n, f)| f.finish(n))
300                    .collect::<PolarsResult<Vec<_>>>()?;
301                StructChunked::from_series(name, len, series.iter())?.into_series()
302            }
303            Self::Skip => Series::new_empty(name, &DataType::Null),
304        })
305    }
306}
307
308/// A List column of `values`, row `i` holding `offsets[i]..offsets[i + 1]`.
309fn list_series(
310    name: PlSmallStr,
311    values: Series,
312    offsets: Vec<i64>,
313    valid: Vec<bool>,
314) -> PolarsResult<Series> {
315    use polars_arrow::array::ListArray;
316    use polars_arrow::bitmap::Bitmap;
317    use polars_arrow::offset::OffsetsBuffer;
318    let inner_dtype = values.dtype().clone();
319    let array = values.to_arrow(0, CompatLevel::newest());
320    let mut all = Vec::with_capacity(offsets.len() + 1);
321    all.push(0i64);
322    all.extend(offsets);
323    let offsets = OffsetsBuffer::<i64>::try_from(all)?;
324    let validity = (!valid.iter().all(|v| *v)).then(|| Bitmap::from_iter(valid));
325    let dtype = ListArray::<i64>::default_datatype(array.dtype().clone());
326    let list = ListArray::<i64>::try_new(dtype, offsets, array, validity)?;
327    Series::from_arrow(name, Box::new(list))?.cast(&DataType::List(Box::new(inner_dtype)))
328}
329
330/// A column of the table and what collects it.
331#[derive(Debug, Clone)]
332struct OutColumn {
333    name: PlSmallStr,
334    dtype: DataType,
335    proto: Sink,
336}
337
338/// Where a run of records is.
339#[derive(Debug, Clone)]
340enum ChunkSource {
341    /// Bytes of the file.
342    Map(Range<usize>),
343    /// A block body of the file, compressed.
344    Block {
345        body: Range<usize>,
346        codec: Compression,
347        uncompressed: Option<usize>,
348    },
349}
350
351#[derive(Debug, Clone)]
352struct Chunk {
353    source: ChunkSource,
354    /// Records in it, when its header (or the index) says.
355    records: Option<u64>,
356    /// A packet's capture time.
357    time_ns: Option<i64>,
358}
359
360/// Where the reading of rows can start: the `row`th record, at `pos` in `chunk`.
361#[derive(Debug, Clone)]
362struct Checkpoint {
363    row: u64,
364    chunk: u32,
365    pos: u32,
366    /// Records of the chunk before `pos`, for a chunk that counts its records.
367    taken: u32,
368    acc: Box<[i128]>,
369}
370
371/// Where the records are.
372#[derive(Debug)]
373enum Index {
374    /// Records all `size` bytes from `start`; row `i` is record `(ring + i) % rows`.
375    Stride {
376        start: usize,
377        size: usize,
378        ring: usize,
379    },
380    Walk(Vec<Checkpoint>),
381}
382
383/// The variant tag of a row read field by field: one cut short, or of no variant.
384const WALK: u8 = u8::MAX;
385
386/// Where each row's record starts and which variant it is, kept from the walk that
387/// opens the file, so every column of a query reads from the same starts rather than
388/// walking the records again (#662). One run of the map only: a block's records are
389/// in its decompressed copy, and a packet's in its payload.
390#[derive(Debug)]
391struct RowTable {
392    /// Where each row's record starts, its sync marker included.
393    starts: crate::indexed::Offsets,
394    /// Each row's variant (0 without variants), or [`WALK`].
395    tags: Vec<u8>,
396}
397
398/// A file's walk, kept for the next open of it with the same spec: the walk is the one
399/// read of the file an open makes, and a file opened again (its variants, `H`) is not
400/// walked again. Kept by [`crate::indexed::keep`], which bounds what it keeps by the
401/// size of the files.
402struct KeptWalk {
403    spec: Spec,
404    data: Range<usize>,
405    index: Arc<Index>,
406    table: Option<Arc<RowTable>>,
407    rows: usize,
408    notes: Vec<String>,
409}
410
411/// Where one column's value is in a record of one variant.
412#[derive(Debug, Clone)]
413enum Source {
414    /// The variant has no such field.
415    Null,
416    /// At a fixed place from the record's start, so read without a walk.
417    At { offset: usize, field: FieldPlan },
418    /// The variant's name.
419    Label,
420    /// Behind a field whose size the record says: the record is walked.
421    Walk,
422    /// A running sum, which needs every record before it: walked from a checkpoint.
423    Summed,
424}
425
426/// The compiled spec: how to walk one record.
427#[derive(Debug)]
428struct Plan {
429    framing: Framing,
430    common: Vec<FieldPlan>,
431    /// The slot of the type field and whether it is text.
432    type_slot: Option<(usize, bool)>,
433    variants: Vec<VariantPlan>,
434    type_out: Option<usize>,
435    /// The record's size, when the spec gives one or a field holds it.
436    size: Option<SizeRef>,
437    /// The width and byte order of a length repeated after the record.
438    suffix: Option<(usize, bool)>,
439    align: usize,
440    sync: Vec<u8>,
441    checksum: Option<ChecksumPlan>,
442    /// Payload header fields, and the slot counting its records.
443    chunk_header: Vec<FieldPlan>,
444    chunk_count: Option<usize>,
445    chunk_slots: usize,
446    time_out: Option<usize>,
447    slots: usize,
448    deltas: Vec<Delta>,
449    columns: Vec<OutColumn>,
450    /// The variant read alone.
451    only: Option<usize>,
452    /// Per column, per variant (one entry without variants): where its value is.
453    sources: Vec<Vec<Source>>,
454}
455
456#[derive(Debug, Clone)]
457struct ChecksumPlan {
458    algo: formats::ChecksumAlgo,
459    field: usize,
460    from: Option<usize>,
461    to: usize,
462    out: usize,
463}
464
465/// A record walk's per-record state: where each field started and the integers read.
466struct Frame {
467    starts: Vec<Option<usize>>,
468    ends: Vec<Option<usize>>,
469    ints: Vec<Option<i128>>,
470}
471
472impl Frame {
473    fn new(slots: usize) -> Self {
474        Self {
475            starts: vec![None; slots],
476            ends: vec![None; slots],
477            ints: vec![None; slots],
478        }
479    }
480
481    /// Forget slots `range`.
482    fn clear(&mut self, range: Range<usize>) {
483        for s in range {
484            self.starts[s] = None;
485            self.ends[s] = None;
486            self.ints[s] = None;
487        }
488    }
489}
490
491/// Why a record could not be read.
492enum Stop {
493    /// The bytes ran out before the record did.
494    Truncated,
495    /// Anything else: said as it is.
496    Said(String),
497}
498
499/// Records read through a spec, by walking them.
500pub struct FramedRecords {
501    bytes: Arc<Bytes>,
502    plan: Arc<Plan>,
503    chunks: Vec<Chunk>,
504    index: Arc<Index>,
505    /// Built by the walk at open, when the records are one run of the map.
506    table: Option<Arc<RowTable>>,
507    rows: usize,
508    schema: SchemaRef,
509    cache: Mutex<VecDeque<(usize, Arc<Vec<u8>>)>>,
510}
511
512impl std::fmt::Debug for FramedRecords {
513    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
514        f.debug_struct("FramedRecords")
515            .field("rows", &self.rows)
516            .field("chunks", &self.chunks.len())
517            .finish()
518    }
519}
520
521/// Whether `spec` needs records walked rather than the fixed reader.
522pub fn needed(spec: &Spec) -> bool {
523    fn walked(field: &Field) -> bool {
524        !field.is_fixed_width()
525            || field.size.as_ref().is_some_and(|s| !s.is_fixed())
526            || field.count.as_ref().is_some_and(|s| !s.is_fixed())
527            || field.delta != Delta::None
528            || !field.bits.is_empty()
529            || field.string_at.is_some()
530    }
531    let r = &spec.records;
532    r.framing != Framing::Fixed
533        || !r.variants.is_empty()
534        || r.checksum.is_some()
535        || r.ring.is_some()
536        || spec.blocks.is_some()
537        || spec.capture.is_some()
538        || formats::all_fields(r).any(walked)
539}
540
541/// What compiling the spec builds: columns, slots and running sums.
542struct Compiler<'a> {
543    spec: &'a Spec,
544    header: &'a HeaderValues,
545    data: &'a [u8],
546    slots: usize,
547    deltas: Vec<Delta>,
548}
549
550/// Names of a scope's fields and their slots.
551#[derive(Default, Clone)]
552struct Scope(Vec<(String, usize)>);
553
554impl Scope {
555    fn slot(&self, name: &str) -> Option<usize> {
556        self.0
557            .iter()
558            .rev()
559            .find(|(n, _)| n == name)
560            .map(|(_, s)| *s)
561    }
562}
563
564impl Compiler<'_> {
565    fn size(&self, amount: &Amount, scope: &Scope, what: &str) -> Result<SizeRef, String> {
566        Ok(match amount {
567            Amount::Record { field, adjust } => SizeRef::Slot {
568                slot: scope
569                    .slot(field)
570                    .ok_or_else(|| format!("{what}: `{field}` is not an earlier field"))?,
571                adjust: *adjust,
572            },
573            Amount::Rest => SizeRef::Rest,
574            other => SizeRef::Given(
575                usize::try_from(self.header.resolve(other, what)?)
576                    .map_err(|_| format!("{what}: too large"))?,
577            ),
578        })
579    }
580
581    /// The fields' plans; their columns are added to `columns`.
582    fn fields(
583        &mut self,
584        fields: &[Field],
585        scope: &mut Scope,
586        columns: &mut Vec<OutColumn>,
587        shared: bool,
588    ) -> Result<Vec<FieldPlan>, String> {
589        let mut out = Vec::new();
590        for field in fields {
591            out.push(self.field(field, scope, columns, shared)?);
592        }
593        Ok(out)
594    }
595
596    /// The column named `name`, new or (for fields variants share) the one there is.
597    fn column(
598        columns: &mut Vec<OutColumn>,
599        name: &str,
600        dtype: DataType,
601        proto: Sink,
602        shared: bool,
603    ) -> usize {
604        if shared && let Some(i) = columns.iter().position(|c| c.name == name) {
605            return i;
606        }
607        columns.push(OutColumn {
608            name: name.into(),
609            dtype,
610            proto,
611        });
612        columns.len() - 1
613    }
614
615    fn field(
616        &mut self,
617        field: &Field,
618        scope: &mut Scope,
619        columns: &mut Vec<OutColumn>,
620        shared: bool,
621    ) -> Result<FieldPlan, String> {
622        let name = field.name.clone().unwrap_or_default();
623        let slot = self.slots;
624        self.slots += 1;
625        let big = field.endian.unwrap_or(self.spec.endian) == formats::Endian::Big;
626        let count = field
627            .count
628            .as_ref()
629            .map(|c| self.size(c, scope, "count"))
630            .transpose()?;
631        let given_count = match count {
632            Some(SizeRef::Given(n)) => {
633                if n as u64 > MAX_ITEMS {
634                    return Err(format!(
635                        "field `{name}`: {n} values is more than {MAX_ITEMS}"
636                    ));
637                }
638                Some(n)
639            }
640            _ => None,
641        };
642        let mut plan = FieldPlan {
643            name: name.clone(),
644            slot,
645            kind: Kind::Pad {
646                size: SizeRef::Given(0),
647            },
648            count,
649            outs: Vec::new(),
650            bits: Vec::new(),
651            delta: field.delta,
652            delta_index: 0,
653            sentinel: None,
654        };
655        let layout = |width: usize, cells: usize| -> Result<ColumnLayout, String> {
656            let mut layout = formats::layout_of(
657                self.spec,
658                field,
659                &name,
660                formats::Place {
661                    start: 0,
662                    stride: None,
663                    width,
664                    count: cells,
665                },
666                self.header,
667            )?;
668            layout.big_endian = big;
669            Ok(layout)
670        };
671        let packed = |layout: ColumnLayout| {
672            let cell = layout.width * layout.count;
673            Sink::Packed {
674                layout,
675                cell,
676                buf: Vec::new(),
677                valid: Vec::new(),
678            }
679        };
680        // A fixed value, possibly counted: an Array, flattened columns, or a List.
681        let mut fixed_out = |this: &mut Self,
682                             plan: &mut FieldPlan,
683                             mut layout: ColumnLayout|
684         -> Result<(), String> {
685            let _ = this;
686            match (count, given_count) {
687                (None, _) => {
688                    let dtype = layout.dtype();
689                    plan.outs = vec![Self::column(columns, &name, dtype, packed(layout), shared)];
690                }
691                (Some(_), Some(n)) if field.flatten => {
692                    layout.count = 1;
693                    for i in 0..n {
694                        let n_name = format!("{name}_{i}");
695                        let mut l = layout.clone();
696                        l.name = n_name.as_str().into();
697                        let dtype = l.dtype();
698                        plan.outs
699                            .push(Self::column(columns, &n_name, dtype, packed(l), shared));
700                    }
701                }
702                (Some(_), Some(n)) => {
703                    layout.count = n.max(1);
704                    let dtype = layout.dtype();
705                    plan.outs = vec![Self::column(columns, &name, dtype, packed(layout), shared)];
706                }
707                (Some(_), None) => {
708                    layout.count = 1;
709                    let item = layout.dtype();
710                    let sink = Sink::List {
711                        inner: Box::new(packed(layout)),
712                        offsets: Vec::new(),
713                        valid: Vec::new(),
714                    };
715                    plan.outs = vec![Self::column(
716                        columns,
717                        &name,
718                        DataType::List(Box::new(item)),
719                        sink,
720                        shared,
721                    )];
722                }
723            }
724            Ok(())
725        };
726        match field.ty {
727            Type::Pad => {
728                let size = field
729                    .size
730                    .as_ref()
731                    .map_or(Ok(SizeRef::Given(0)), |s| self.size(s, scope, "size"))?;
732                plan.kind = Kind::Pad { size };
733            }
734            Type::Group => {
735                let mut inner_scope = Scope::default();
736                let mut inner_columns = Vec::new();
737                let first = self.slots;
738                let items =
739                    self.fields(&field.group, &mut inner_scope, &mut inner_columns, false)?;
740                let slots = self.slots - first;
741                let names: Vec<PlSmallStr> = inner_columns.iter().map(|c| c.name.clone()).collect();
742                let struct_dtype = DataType::Struct(
743                    inner_columns
744                        .iter()
745                        .map(|c| polars::prelude::Field::new(c.name.clone(), c.dtype.clone()))
746                        .collect(),
747                );
748                let sink = Sink::List {
749                    inner: Box::new(Sink::Struct {
750                        names,
751                        fields: inner_columns.into_iter().map(|c| c.proto).collect(),
752                        len: 0,
753                    }),
754                    offsets: Vec::new(),
755                    valid: Vec::new(),
756                };
757                plan.outs = vec![Self::column(
758                    columns,
759                    &name,
760                    DataType::List(Box::new(struct_dtype)),
761                    sink,
762                    shared,
763                )];
764                plan.kind = Kind::Group { items, slots };
765            }
766            Type::VarU | Type::VarS => {
767                let signed = field.ty == Type::VarS;
768                let mut l = layout(8, 1)?;
769                if field.delta != Delta::None {
770                    l.physical = Physical::Signed(8);
771                    plan.sentinel = sentinel(field, 64, signed);
772                    l.null = None;
773                }
774                plan.kind = Kind::Var { signed };
775                fixed_out(self, &mut plan, l)?;
776            }
777            Type::Strz => {
778                let max = field
779                    .size
780                    .as_ref()
781                    .map(|s| self.size(s, scope, "size"))
782                    .transpose()?;
783                plan.kind = Kind::Strz {
784                    max,
785                    encoding: field.encoding,
786                };
787                plan.outs = vec![Self::column(
788                    columns,
789                    &name,
790                    DataType::String,
791                    Sink::Text(Vec::new()),
792                    shared,
793                )];
794            }
795            Type::Str | Type::Bytes if !field.size.as_ref().is_some_and(Amount::is_fixed) => {
796                let size = self.size(
797                    field.size.as_ref().expect("checked at parse"),
798                    scope,
799                    "size",
800                )?;
801                let text = field.ty == Type::Str;
802                plan.kind = Kind::Sized {
803                    size,
804                    encoding: text.then_some(field.encoding),
805                };
806                let (dtype, sink) = if text {
807                    (DataType::String, Sink::Text(Vec::new()))
808                } else {
809                    (DataType::Binary, Sink::Binary(Vec::new()))
810                };
811                plan.outs = vec![Self::column(columns, &name, dtype, sink, shared)];
812            }
813            _ => {
814                let width = match (field.ty.width(), &field.size) {
815                    (Some(w), _) => w as usize,
816                    (None, Some(amount)) => usize::try_from(self.header.resolve(amount, "size")?)
817                        .map_err(|_| "size: too large".to_string())?,
818                    (None, None) => 0,
819                };
820                if width == 0 {
821                    return Err(format!("field `{name}` takes no bytes"));
822                }
823                let int = match field.ty {
824                    Type::Unsigned(_) => Some(IntRead { signed: false, big }),
825                    Type::Signed(_) => Some(IntRead { signed: true, big }),
826                    _ => None,
827                };
828                if let Some(section) = &field.string_at {
829                    let section = self.section(section)?;
830                    plan.kind = Kind::StringAt {
831                        width,
832                        big,
833                        section,
834                    };
835                    plan.outs = vec![Self::column(
836                        columns,
837                        &name,
838                        DataType::String,
839                        Sink::Text(Vec::new()),
840                        shared,
841                    )];
842                } else {
843                    let mut l = layout(width, 1)?;
844                    if field.delta != Delta::None {
845                        l.physical = Physical::Signed(8);
846                        l.width = 8;
847                        plan.sentinel =
848                            sentinel(field, width as u32 * 8, int.is_some_and(|i| i.signed));
849                        l.null = None;
850                    }
851                    plan.kind = Kind::Fixed { width, int };
852                    fixed_out(self, &mut plan, l)?;
853                }
854            }
855        }
856        if field.delta != Delta::None {
857            plan.delta_index = self.deltas.len();
858            self.deltas.push(field.delta);
859        }
860        for bit in &field.bits {
861            let sink = Sink::Bits {
862                width: bit.width,
863                labels: bit.labels.clone(),
864                values: Vec::new(),
865            };
866            let dtype = match (&bit.labels, bit.width) {
867                (Some(_), _) => DataType::String,
868                (None, 1) => DataType::Boolean,
869                (None, 2..=8) => DataType::UInt8,
870                (None, 9..=16) => DataType::UInt16,
871                (None, 17..=32) => DataType::UInt32,
872                _ => DataType::UInt64,
873            };
874            let col = Self::column(columns, &bit.name, dtype, sink, shared);
875            plan.bits.push((col, bit.bit, bit.width));
876        }
877        if field.name.is_some() {
878            scope.0.push((name, slot));
879        }
880        Ok(plan)
881    }
882
883    /// The bytes of section `name`, bounded by the file.
884    fn section(&self, name: &str) -> Result<Range<usize>, String> {
885        let section = self
886            .spec
887            .sections
888            .iter()
889            .find(|s| s.name == name)
890            .ok_or_else(|| format!("no section named {name}"))?;
891        let offset = self.header.resolve_any(&section.offset, "section offset")?;
892        let size = self.header.resolve_any(&section.size, "section size")?;
893        let end = offset.saturating_add(size);
894        if end > self.data.len() as u64 {
895            return Err(format!(
896                "section {name} runs from byte {offset} to {end}, past the file's {} bytes",
897                self.data.len()
898            ));
899        }
900        Ok(offset as usize..end as usize)
901    }
902}
903
904/// The raw value a delta field's `null` stands for.
905fn sentinel(field: &Field, bits: u32, signed: bool) -> Option<i128> {
906    use crate::fixed_records::Null;
907    let bits = bits.clamp(1, 64);
908    match field.null? {
909        Null::Min if signed => Some(-(1i128 << (bits - 1))),
910        Null::Min => Some(0),
911        Null::Max if signed => Some((1i128 << (bits - 1)) - 1),
912        Null::Max => Some((1i128 << bits) - 1),
913        Null::Value(v) => Some(v),
914        Null::NaN => None,
915    }
916}
917
918/// An integer of `bytes`.
919fn int_of(bytes: &[u8], read: IntRead) -> i128 {
920    if read.signed {
921        i128::from(crate::fixed_records::read_signed(bytes, read.big))
922    } else {
923        i128::from(crate::fixed_records::read_unsigned(bytes, read.big))
924    }
925}
926
927/// A LEB128 value at the front of `bytes`, and the bytes it took.
928fn leb128(bytes: &[u8]) -> Option<(u64, usize)> {
929    let mut value = 0u64;
930    for (i, b) in bytes.iter().take(10).enumerate() {
931        let part = u64::from(b & 0x7f);
932        if i == 9 && part > 1 {
933            return None;
934        }
935        value |= part << (7 * i);
936        if b & 0x80 == 0 {
937            return Some((value, i + 1));
938        }
939    }
940    None
941}
942
943/// Text of `bytes` in `encoding`, trimmed as a fixed field's is.
944fn decode_text(bytes: &[u8], encoding: Encoding) -> String {
945    match encoding {
946        Encoding::Utf8 => crate::fixed_records::text(bytes),
947        Encoding::Latin1 => crate::fixed_records::latin1(bytes),
948        Encoding::Utf16Le => crate::fixed_records::utf16(bytes, false),
949        Encoding::Utf16Be => crate::fixed_records::utf16(bytes, true),
950    }
951}
952
953/// Where the NUL unit ends text at the front of `bytes`, if there is one.
954fn nul_at(bytes: &[u8], unit: usize) -> Option<usize> {
955    if unit == 1 {
956        memchr::memchr(0, bytes)
957    } else {
958        bytes
959            .chunks_exact(unit)
960            .position(|c| c.iter().all(|b| *b == 0))
961            .map(|i| i * unit)
962    }
963}
964
965/// What reading at a place found.
966#[derive(Debug, Clone, Copy, PartialEq, Eq)]
967enum Got {
968    /// A row.
969    Row,
970    /// A record of a variant not read.
971    Skipped,
972    /// No record left.
973    None,
974}
975
976/// What a walk pushes to: the sinks of one level's columns, and which it filled.
977struct Out<'s> {
978    sinks: &'s mut [Sink],
979    filled: &'s mut [bool],
980}
981
982impl Out<'_> {
983    /// Nulls in every column not filled, and a fresh start for the next row.
984    fn finish_row(&mut self) {
985        for (sink, filled) in self.sinks.iter_mut().zip(self.filled.iter_mut()) {
986            if !*filled {
987                sink.push_null();
988            }
989            *filled = false;
990        }
991    }
992}
993
994/// Reads records out of one run of bytes.
995struct Walker<'a> {
996    plan: &'a Plan,
997    /// The whole file, for sections.
998    file: &'a [u8],
999    frame: Frame,
1000    acc: Vec<i128>,
1001    /// The end of what the current record (or group item) may take.
1002    end: usize,
1003    /// Whether `end` is the record's own end, so a field past it is null.
1004    bounded: bool,
1005    /// A field did not fit, so the rest of the record is null.
1006    short: bool,
1007    record_start: usize,
1008    /// The field that holds the record's size, and what to add to it.
1009    size_slot: Option<(usize, i64)>,
1010    /// Bytes skipped looking for sync markers.
1011    skipped: u64,
1012    /// The variant of the record read last, if it had one.
1013    variant: Option<usize>,
1014}
1015
1016impl<'a> Walker<'a> {
1017    fn new(plan: &'a Plan, file: &'a [u8]) -> Self {
1018        Self {
1019            plan,
1020            file,
1021            frame: Frame::new(plan.slots),
1022            acc: vec![0; plan.deltas.len()],
1023            end: 0,
1024            bounded: false,
1025            short: false,
1026            record_start: 0,
1027            size_slot: None,
1028            skipped: 0,
1029            variant: None,
1030        }
1031    }
1032
1033    fn size_of(&self, s: SizeRef, pos: usize, what: &str) -> Result<usize, Stop> {
1034        match s {
1035            SizeRef::Given(n) => Ok(n),
1036            SizeRef::Rest => Ok(self.end.saturating_sub(pos)),
1037            SizeRef::Slot { slot, adjust } => {
1038                let v = self.frame.ints[slot].ok_or(Stop::Truncated)? + i128::from(adjust);
1039                usize::try_from(v)
1040                    .ok()
1041                    .filter(|n| *n as u64 <= formats::MAX_SIZE)
1042                    .ok_or_else(|| {
1043                        Stop::Said(format!(
1044                            "`{what}` at byte {pos} is {v}, outside 0 to {}",
1045                            formats::MAX_SIZE
1046                        ))
1047                    })
1048            }
1049        }
1050    }
1051
1052    fn take(&self, pos: &mut usize, n: usize) -> Result<Range<usize>, Stop> {
1053        let start = *pos;
1054        let stop = start.checked_add(n).ok_or(Stop::Truncated)?;
1055        if stop > self.end {
1056            return Err(Stop::Truncated);
1057        }
1058        *pos = stop;
1059        Ok(start..stop)
1060    }
1061
1062    /// Read `fields` from `data` at `pos`, pushing to `out`. A field that does not fit
1063    /// is null once the record's end is known, and truncates the record before then.
1064    fn walk(
1065        &mut self,
1066        fields: &[FieldPlan],
1067        data: &[u8],
1068        pos: &mut usize,
1069        mut out: Option<&mut Out<'_>>,
1070    ) -> Result<(), Stop> {
1071        for f in fields {
1072            let res = if self.short {
1073                Err(Stop::Truncated)
1074            } else {
1075                self.field(f, data, pos, out.as_deref_mut())
1076            };
1077            match res {
1078                Ok(()) => {}
1079                Err(Stop::Truncated) if self.bounded => {
1080                    self.short = true;
1081                    if let Some(out) = out.as_deref_mut() {
1082                        null_outs(f, out);
1083                    }
1084                }
1085                Err(e) => return Err(e),
1086            }
1087            if let Some((slot, adjust)) = self.size_slot
1088                && slot == f.slot
1089                && !self.bounded
1090            {
1091                let start = self.record_start;
1092                let len = self.frame.ints[slot].ok_or(Stop::Truncated)? + i128::from(adjust);
1093                let total = usize::try_from(len)
1094                    .ok()
1095                    .filter(|t| *t as u64 <= formats::MAX_SIZE && *t >= *pos - start)
1096                    .ok_or_else(|| {
1097                        Stop::Said(format!(
1098                            "the record at byte {start} gives its size as {len}, outside {} to {}",
1099                            *pos - start,
1100                            formats::MAX_SIZE
1101                        ))
1102                    })?;
1103                let rec_end = start + total;
1104                if rec_end > self.end {
1105                    return Err(Stop::Truncated);
1106                }
1107                self.end = rec_end;
1108                self.bounded = true;
1109            }
1110        }
1111        Ok(())
1112    }
1113
1114    fn field(
1115        &mut self,
1116        f: &FieldPlan,
1117        data: &[u8],
1118        pos: &mut usize,
1119        out: Option<&mut Out<'_>>,
1120    ) -> Result<(), Stop> {
1121        self.frame.starts[f.slot] = Some(*pos);
1122        let count = match f.count {
1123            None => None,
1124            Some(c) => {
1125                let n = self.size_of(c, *pos, &f.name)?;
1126                if n as u64 > MAX_ITEMS {
1127                    return Err(Stop::Said(format!(
1128                        "`{}` at byte {} counts {n} values, more than {MAX_ITEMS}",
1129                        f.name, *pos
1130                    )));
1131                }
1132                Some(n)
1133            }
1134        };
1135        match &f.kind {
1136            Kind::Pad { size } => {
1137                let n = self.size_of(*size, *pos, "pad")?;
1138                self.take(pos, n)?;
1139            }
1140            Kind::Fixed { width, int } => {
1141                let cells = count.unwrap_or(1);
1142                let range = self.take(pos, width.checked_mul(cells).ok_or(Stop::Truncated)?)?;
1143                let bytes = &data[range];
1144                let raw = int.and_then(|r| (cells > 0).then(|| int_of(&bytes[..*width], r)));
1145                if count.is_none() {
1146                    self.frame.ints[f.slot] = raw;
1147                }
1148                self.value(f, raw, Some((bytes, *width)), count, out);
1149            }
1150            Kind::Var { signed } => {
1151                let cells = count.unwrap_or(1);
1152                let mut bytes = Vec::with_capacity(cells.min(64) * 8);
1153                let mut first = None;
1154                for _ in 0..cells {
1155                    let (v, n) =
1156                        leb128(&data[(*pos).min(self.end)..self.end]).ok_or(Stop::Truncated)?;
1157                    *pos += n;
1158                    let v = if *signed {
1159                        i128::from(((v >> 1) as i64) ^ -((v & 1) as i64))
1160                    } else {
1161                        i128::from(v)
1162                    };
1163                    first.get_or_insert(v);
1164                    if *signed {
1165                        bytes.extend((v as i64).to_le_bytes());
1166                    } else {
1167                        bytes.extend((v as u64).to_le_bytes());
1168                    }
1169                }
1170                if count.is_none() {
1171                    self.frame.ints[f.slot] = first;
1172                }
1173                self.value(f, first, Some((&bytes, 8)), count, out);
1174            }
1175            Kind::Sized { size, encoding } => {
1176                let n = self.size_of(*size, *pos, &f.name)?;
1177                let range = self.take(pos, n)?;
1178                if let Some(out) = out {
1179                    let col = f.outs[0];
1180                    out.filled[col] = true;
1181                    match (&mut out.sinks[col], encoding) {
1182                        (Sink::Text(v), Some(enc)) => {
1183                            v.push(Some(decode_text(&data[range], *enc)));
1184                        }
1185                        (Sink::Binary(v), None) => v.push(Some(data[range].to_vec())),
1186                        _ => {}
1187                    }
1188                }
1189            }
1190            Kind::Strz { max, encoding } => {
1191                let unit = encoding.unit();
1192                let text = match max {
1193                    Some(m) => {
1194                        let n = self.size_of(*m, *pos, &f.name)?;
1195                        let range = self.take(pos, n)?;
1196                        let bytes = &data[range];
1197                        let stop = nul_at(bytes, unit).unwrap_or(bytes.len());
1198                        decode_text(&bytes[..stop], *encoding)
1199                    }
1200                    None => {
1201                        let rest = &data[(*pos).min(self.end)..self.end];
1202                        let stop = nul_at(rest, unit).ok_or(Stop::Truncated)?;
1203                        let text = decode_text(&rest[..stop], *encoding);
1204                        *pos += stop + unit;
1205                        text
1206                    }
1207                };
1208                if let Some(out) = out {
1209                    let col = f.outs[0];
1210                    out.filled[col] = true;
1211                    if let Sink::Text(v) = &mut out.sinks[col] {
1212                        v.push(Some(text));
1213                    }
1214                }
1215            }
1216            Kind::StringAt {
1217                width,
1218                big,
1219                section,
1220            } => {
1221                let range = self.take(pos, *width)?;
1222                let offset = crate::fixed_records::read_unsigned(&data[range], *big);
1223                self.frame.ints[f.slot] = Some(i128::from(offset));
1224                if let Some(out) = out {
1225                    let col = f.outs[0];
1226                    out.filled[col] = true;
1227                    let heap = &self.file[section.clone()];
1228                    let text = usize::try_from(offset)
1229                        .ok()
1230                        .filter(|o| *o < heap.len())
1231                        .map(|o| {
1232                            let rest = &heap[o..];
1233                            let stop = memchr::memchr(0, rest).unwrap_or(rest.len());
1234                            crate::fixed_records::text(&rest[..stop])
1235                        });
1236                    if let Sink::Text(v) = &mut out.sinks[col] {
1237                        v.push(text);
1238                    }
1239                }
1240            }
1241            Kind::Group { items, slots } => {
1242                let n = count.unwrap_or(0);
1243                let first = items.first().map_or(0, |i| i.slot);
1244                let saved = (
1245                    self.end,
1246                    self.bounded,
1247                    self.short,
1248                    self.record_start,
1249                    self.size_slot,
1250                );
1251                self.size_slot = None;
1252                // Walked once to see the whole group fits, then again into the sinks, so a
1253                // group cut short leaves no half-pushed items.
1254                let start = *pos;
1255                let mut probe = start;
1256                let mut fits = Ok(());
1257                for _ in 0..n {
1258                    self.frame.clear(first..first + slots);
1259                    self.bounded = true;
1260                    self.short = false;
1261                    self.record_start = probe;
1262                    if let Err(e) = self.walk(items, data, &mut probe, None) {
1263                        fits = Err(e);
1264                        break;
1265                    }
1266                    if self.short {
1267                        fits = Err(Stop::Truncated);
1268                        break;
1269                    }
1270                }
1271                let (end, bounded, short, record_start, size_slot) = saved;
1272                self.end = end;
1273                if let Err(e) = fits {
1274                    self.bounded = bounded;
1275                    self.short = short;
1276                    self.record_start = record_start;
1277                    self.size_slot = size_slot;
1278                    return Err(e);
1279                }
1280                *pos = probe;
1281                if let Some(out) = out {
1282                    let col = f.outs[0];
1283                    out.filled[col] = true;
1284                    if let Sink::List {
1285                        inner,
1286                        offsets,
1287                        valid,
1288                    } = &mut out.sinks[col]
1289                        && let Sink::Struct { fields, len, .. } = inner.as_mut()
1290                    {
1291                        let mut at = start;
1292                        let mut filled = vec![false; fields.len()];
1293                        for _ in 0..n {
1294                            self.frame.clear(first..first + slots);
1295                            self.bounded = true;
1296                            self.short = false;
1297                            self.record_start = at;
1298                            let mut item_out = Out {
1299                                sinks: fields,
1300                                filled: &mut filled,
1301                            };
1302                            // It fitted above, and reads the same again.
1303                            let _ = self.walk(items, data, &mut at, Some(&mut item_out));
1304                            item_out.finish_row();
1305                            *len += 1;
1306                        }
1307                        offsets.push(offsets.last().copied().unwrap_or(0) + n as i64);
1308                        valid.push(true);
1309                    }
1310                }
1311                self.end = end;
1312                self.bounded = bounded;
1313                self.short = short;
1314                self.record_start = record_start;
1315                self.size_slot = size_slot;
1316            }
1317        }
1318        self.frame.ends[f.slot] = Some(*pos);
1319        Ok(())
1320    }
1321
1322    /// Push a number (and its bits, and its running sum) to `out`, or only keep the sum.
1323    fn value(
1324        &mut self,
1325        f: &FieldPlan,
1326        raw: Option<i128>,
1327        bytes: Option<(&[u8], usize)>,
1328        count: Option<usize>,
1329        out: Option<&mut Out<'_>>,
1330    ) {
1331        let summed = if f.delta == Delta::None {
1332            None
1333        } else {
1334            let raw = raw.unwrap_or(0);
1335            if Some(raw) == f.sentinel {
1336                Some(None)
1337            } else {
1338                let sum = self.acc[f.delta_index].wrapping_add(raw);
1339                self.acc[f.delta_index] = sum;
1340                Some(i64::try_from(sum).ok())
1341            }
1342        };
1343        let Some(out) = out else { return };
1344        match (summed, bytes) {
1345            (Some(sum), _) => {
1346                let col = f.outs[0];
1347                out.filled[col] = true;
1348                match sum {
1349                    Some(v) => out.sinks[col].push_bytes(&v.to_le_bytes()),
1350                    None => out.sinks[col].push_null(),
1351                }
1352            }
1353            (None, Some((bytes, width))) => push_fixed(f, bytes, width, count, out),
1354            (None, None) => {}
1355        }
1356        push_bits(f, raw, out);
1357    }
1358
1359    /// Read one record at `pos`, before `chunk_end`; `pos` ends after it.
1360    fn record(
1361        &mut self,
1362        data: &[u8],
1363        pos: &mut usize,
1364        chunk_end: usize,
1365        out: Option<&mut Out<'_>>,
1366        time: Option<i64>,
1367    ) -> Result<Got, Stop> {
1368        let Some(only) = self.plan.only else {
1369            return Ok(match self.record_inner(data, pos, chunk_end, out, time)? {
1370                None => Got::None,
1371                Some(_) => Got::Row,
1372            });
1373        };
1374        // One variant read alone: the record is walked to learn its variant, and walked
1375        // again into the columns when it is the one.
1376        let (start, acc, skipped) = (*pos, self.acc.clone(), self.skipped);
1377        match self.record_inner(data, pos, chunk_end, None, time)? {
1378            None => Ok(Got::None),
1379            Some(Some(v)) if v == only => {
1380                if out.is_some() {
1381                    *pos = start;
1382                    self.acc = acc;
1383                    self.skipped = skipped;
1384                    self.record_inner(data, pos, chunk_end, out, time)?;
1385                }
1386                Ok(Got::Row)
1387            }
1388            Some(_) => Ok(Got::Skipped),
1389        }
1390    }
1391
1392    /// Read one record: `None` when there is none left (for sync framing, no marker),
1393    /// else the index of its variant, if it has one.
1394    fn record_inner(
1395        &mut self,
1396        data: &[u8],
1397        pos: &mut usize,
1398        chunk_end: usize,
1399        mut out: Option<&mut Out<'_>>,
1400        time: Option<i64>,
1401    ) -> Result<Option<Option<usize>>, Stop> {
1402        let plan = self.plan;
1403        if *pos >= chunk_end {
1404            return Ok(None);
1405        }
1406        if !plan.sync.is_empty() {
1407            let rest = &data[*pos..chunk_end];
1408            if !rest.starts_with(&plan.sync) {
1409                match memchr::memmem::find(rest, &plan.sync) {
1410                    Some(at) => {
1411                        self.skipped += at as u64;
1412                        *pos += at;
1413                    }
1414                    None => {
1415                        self.skipped += rest.len() as u64;
1416                        *pos = chunk_end;
1417                        return Ok(None);
1418                    }
1419                }
1420            }
1421            *pos += plan.sync.len();
1422        }
1423        let first = plan.chunk_slots;
1424        self.frame.clear(first..plan.slots);
1425        self.record_start = *pos;
1426        self.end = chunk_end;
1427        self.bounded = false;
1428        self.short = false;
1429        self.size_slot = None;
1430        match plan.size {
1431            Some(SizeRef::Given(n)) => {
1432                let end = pos.checked_add(n).ok_or(Stop::Truncated)?;
1433                if end > chunk_end {
1434                    return Err(Stop::Truncated);
1435                }
1436                self.end = end;
1437                self.bounded = true;
1438            }
1439            Some(SizeRef::Slot { slot, adjust }) => self.size_slot = Some((slot, adjust)),
1440            _ => {}
1441        }
1442        let start = *pos;
1443        let mut p = start;
1444        self.walk(&plan.common, data, &mut p, out.as_deref_mut())?;
1445        let mut label = None;
1446        let mut chosen = None;
1447        if let Some((slot, text)) = plan.type_slot {
1448            let int = self.frame.ints[slot];
1449            let shown = if text {
1450                match (self.frame.starts[slot], self.frame.ends[slot]) {
1451                    (Some(a), Some(b)) => Some(crate::fixed_records::text(&data[a..b])),
1452                    _ => None,
1453                }
1454            } else {
1455                None
1456            };
1457            let variant = plan
1458                .variants
1459                .iter()
1460                .enumerate()
1461                .find(|(_, v)| match (&shown, int) {
1462                    (Some(t), _) => v.texts.iter().any(|w| w == t),
1463                    (None, Some(i)) => v.ints.contains(&i),
1464                    _ => false,
1465                });
1466            match variant {
1467                Some((index, variant)) => {
1468                    chosen = Some(index);
1469                    if let Some(size) = variant.size
1470                        && !self.bounded
1471                    {
1472                        let n = self.size_of(size, p, "size")?;
1473                        let end = start.checked_add(n).ok_or(Stop::Truncated)?;
1474                        if end > chunk_end || end < p {
1475                            return Err(if end > chunk_end {
1476                                Stop::Truncated
1477                            } else {
1478                                Stop::Said(format!(
1479                                    "variant {} at byte {start} is {n} bytes, less than its common fields",
1480                                    variant.name
1481                                ))
1482                            });
1483                        }
1484                        self.end = end;
1485                        self.bounded = true;
1486                    }
1487                    self.walk(&variant.fields, data, &mut p, out.as_deref_mut())?;
1488                    label = Some(variant.name.clone());
1489                }
1490                None => {
1491                    let shown = shown.or_else(|| int.map(|i| i.to_string()));
1492                    if !self.bounded {
1493                        return Err(Stop::Said(format!(
1494                            "the record at byte {start} has type {}, which no variant names",
1495                            shown.as_deref().unwrap_or("(none)")
1496                        )));
1497                    }
1498                    label = shown.map(|s| Arc::from(format!("?{s}")));
1499                }
1500            }
1501        }
1502        self.variant = chosen;
1503        *pos = if self.bounded { self.end } else { p };
1504        if let Some((width, big)) = plan.suffix {
1505            let range = self.take_at(*pos, width, chunk_end)?;
1506            let again = crate::fixed_records::read_unsigned(&data[range], big);
1507            let (slot, _) = self.size_slot.unwrap_or((usize::MAX, 0));
1508            let first = self.frame.ints.get(slot).copied().flatten();
1509            if first != Some(i128::from(again)) {
1510                return Err(Stop::Said(format!(
1511                    "the record at byte {start} ends with length {again}, not the {} it starts with",
1512                    first.map_or_else(|| "?".to_string(), |v| v.to_string())
1513                )));
1514            }
1515            *pos += width;
1516        }
1517        if let Some(out) = out {
1518            if let (Some(col), Some(label)) = (plan.type_out, label)
1519                && let Sink::Label(v) = &mut out.sinks[col]
1520            {
1521                v.push(Some(label));
1522                out.filled[col] = true;
1523            }
1524            if let Some(check) = &plan.checksum {
1525                let from = check.from.map_or(Some(start), |s| self.frame.starts[s]);
1526                let to = self.frame.starts[check.to];
1527                let stored = self.frame.ints[check.field];
1528                let ok = match (from, to, stored) {
1529                    (Some(a), Some(b), Some(v)) if a <= b => {
1530                        Some(i128::from(check.algo.compute(&data[a..b])) == v)
1531                    }
1532                    _ => None,
1533                };
1534                if let Sink::Flag(v) = &mut out.sinks[check.out] {
1535                    v.push(ok);
1536                    out.filled[check.out] = true;
1537                }
1538            }
1539            if let Some(col) = plan.time_out
1540                && let Sink::Time(v) = &mut out.sinks[col]
1541            {
1542                v.push(time);
1543                out.filled[col] = true;
1544            }
1545            out.finish_row();
1546        }
1547        Ok(Some(chosen))
1548    }
1549
1550    fn take_at(&self, pos: usize, n: usize, end: usize) -> Result<Range<usize>, Stop> {
1551        let stop = pos.checked_add(n).ok_or(Stop::Truncated)?;
1552        if stop > end {
1553            return Err(Stop::Truncated);
1554        }
1555        Ok(pos..stop)
1556    }
1557}
1558
1559/// The bit fields of `f`'s value `raw` to their columns.
1560fn push_bits(f: &FieldPlan, raw: Option<i128>, out: &mut Out<'_>) {
1561    for (col, bit, width) in &f.bits {
1562        out.filled[*col] = true;
1563        match (raw, &mut out.sinks[*col]) {
1564            (Some(raw), Sink::Bits { values, .. }) => {
1565                let mask = if *width >= 64 {
1566                    u64::MAX
1567                } else {
1568                    (1u64 << width) - 1
1569                };
1570                values.push(Some(((raw as u64) >> bit) & mask));
1571            }
1572            (_, sink) => sink.push_null(),
1573        }
1574    }
1575}
1576
1577/// Nulls in every column of `f`.
1578fn null_outs(f: &FieldPlan, out: &mut Out<'_>) {
1579    for col in f.outs.iter().chain(f.bits.iter().map(|(c, _, _)| c)) {
1580        if !out.filled[*col] {
1581            out.sinks[*col].push_null();
1582            out.filled[*col] = true;
1583        }
1584    }
1585}
1586
1587fn push_fixed(f: &FieldPlan, bytes: &[u8], width: usize, count: Option<usize>, out: &mut Out<'_>) {
1588    if f.outs.len() == 1 {
1589        let col = f.outs[0];
1590        out.filled[col] = true;
1591        match (&mut out.sinks[col], count) {
1592            (
1593                Sink::List {
1594                    inner,
1595                    offsets,
1596                    valid,
1597                },
1598                Some(n),
1599            ) => {
1600                for i in 0..n {
1601                    inner.push_bytes(&bytes[i * width..(i + 1) * width]);
1602                }
1603                offsets.push(offsets.last().copied().unwrap_or(0) + n as i64);
1604                valid.push(true);
1605            }
1606            (sink, _) => sink.push_bytes(bytes),
1607        }
1608    } else {
1609        for (i, col) in f.outs.iter().enumerate() {
1610            out.filled[*col] = true;
1611            out.sinks[*col].push_bytes(&bytes[i * width..(i + 1) * width]);
1612        }
1613    }
1614}
1615
1616/// Bytes of one chunk: a range of the map, or a decompressed block.
1617enum ChunkData {
1618    Map(Range<usize>),
1619    Owned(Arc<Vec<u8>>),
1620}
1621
1622/// Where a read is: in which chunk, where, and how many of its records it has taken.
1623struct Cursor {
1624    pos: usize,
1625    data_start: usize,
1626    end: usize,
1627    taken: u64,
1628    limit: Option<u64>,
1629    time: Option<i64>,
1630}
1631
1632/// The rows a read wants, and how many it has.
1633struct Want {
1634    start: u64,
1635    len: usize,
1636    produced: usize,
1637}
1638
1639impl Want {
1640    /// A null row at `row`, if it is one wanted.
1641    fn null_row(&mut self, row: u64, out: &mut Out<'_>) {
1642        if row >= self.start && self.produced < self.len {
1643            out.finish_row();
1644            self.produced += 1;
1645        }
1646    }
1647}
1648
1649/// What opening found to say: trailing bytes, skipped bytes, records cut short.
1650#[derive(Default)]
1651struct Found {
1652    notes: Vec<String>,
1653    skipped: u64,
1654}
1655
1656impl FramedRecords {
1657    /// Read `spec`'s records from `bytes[data]`, with `header` read already.
1658    pub fn open(
1659        spec: &Spec,
1660        bytes: Arc<Bytes>,
1661        header: &HeaderValues,
1662        data: Range<usize>,
1663        named: &str,
1664        path: Option<&std::path::Path>,
1665    ) -> Result<(Self, Vec<String>), String> {
1666        let file = bytes.as_slice();
1667        let mut compiler = Compiler {
1668            spec,
1669            header,
1670            data: file,
1671            slots: 0,
1672            deltas: Vec::new(),
1673        };
1674        let mut columns = Vec::new();
1675        let mut chunk_scope = Scope::default();
1676        let mut discard = Vec::new();
1677        let chunk_header = match &spec.capture {
1678            Some(capture) => {
1679                compiler.fields(&capture.header, &mut chunk_scope, &mut discard, false)?
1680            }
1681            None => Vec::new(),
1682        };
1683        let chunk_count = spec
1684            .capture
1685            .as_ref()
1686            .and_then(|c| c.count.as_ref())
1687            .and_then(|n| chunk_scope.slot(n));
1688        let chunk_slots = compiler.slots;
1689        let time_out = spec
1690            .capture
1691            .as_ref()
1692            .and_then(|c| c.time.as_ref())
1693            .map(|name| {
1694                Compiler::column(
1695                    &mut columns,
1696                    name,
1697                    DataType::Datetime(TimeUnit::Nanoseconds, None),
1698                    Sink::Time(Vec::new()),
1699                    false,
1700                )
1701            });
1702        let records = &spec.records;
1703        let mut scope = Scope::default();
1704        let common = compiler.fields(&records.fields, &mut scope, &mut columns, false)?;
1705        let type_slot = records.type_field.as_ref().map(|name| {
1706            let field = records
1707                .fields
1708                .iter()
1709                .find(|f| f.name.as_deref() == Some(name.as_str()))
1710                .expect("checked at parse");
1711            (
1712                scope.slot(name).expect("a common field"),
1713                field.ty.is_text(),
1714            )
1715        });
1716        let only = match &spec.variant {
1717            Some(name) => Some(
1718                records
1719                    .variants
1720                    .iter()
1721                    .position(|v| &v.name == name)
1722                    .ok_or_else(|| format!("no variant named {name}"))?,
1723            ),
1724            None => None,
1725        };
1726        let type_out = type_slot.filter(|_| only.is_none()).map(|_| {
1727            Compiler::column(
1728                &mut columns,
1729                "type",
1730                DataType::from_categories(Categories::global()),
1731                Sink::Label(Vec::new()),
1732                false,
1733            )
1734        });
1735        let size = records
1736            .size
1737            .as_ref()
1738            .map(|s| compiler.size(s, &scope, "record size"))
1739            .transpose()?;
1740        let mut variants = Vec::new();
1741        let mut unread = Vec::new();
1742        for (i, variant) in records.variants.iter().enumerate() {
1743            let mut vscope = scope.clone();
1744            let into = if only.is_some_and(|o| o != i) {
1745                &mut unread
1746            } else {
1747                &mut columns
1748            };
1749            let fields = compiler.fields(&variant.fields, &mut vscope, into, true)?;
1750            let size = variant
1751                .size
1752                .as_ref()
1753                .map(|s| compiler.size(s, &vscope, "variant size"))
1754                .transpose()?;
1755            if matches!(size, Some(SizeRef::Slot { .. })) {
1756                return Err(format!(
1757                    "variant {}: its size comes from the header or is written in the spec",
1758                    variant.name
1759                ));
1760            }
1761            let mut ints = Vec::new();
1762            let mut texts = Vec::new();
1763            for when in &variant.when {
1764                match when {
1765                    formats::Expected::Int(v) => ints.push(*v),
1766                    // Text fields read with their padding trimmed: `"fmt "` is `fmt`.
1767                    formats::Expected::Text(t) => {
1768                        texts.push(t.trim_end_matches(['\0', ' ']).to_string())
1769                    }
1770                }
1771            }
1772            variants.push(VariantPlan {
1773                name: Arc::from(variant.name.as_str()),
1774                ints,
1775                texts,
1776                fields,
1777                size,
1778            });
1779        }
1780        let suffix = if records.length_suffix {
1781            let Some(Amount::Record { field, .. }) = &records.size else {
1782                return Err("length_suffix: needs size from a field".into());
1783            };
1784            let target = records
1785                .fields
1786                .iter()
1787                .find(|f| f.name.as_deref() == Some(field.as_str()))
1788                .ok_or_else(|| format!("length_suffix: `{field}` is not a common field"))?;
1789            let width = target
1790                .ty
1791                .width()
1792                .ok_or("length_suffix: the length field is fixed-width")?;
1793            let big = target.endian.unwrap_or(spec.endian) == formats::Endian::Big;
1794            Some((width as usize, big))
1795        } else {
1796            None
1797        };
1798        let checksum = records
1799            .checksum
1800            .as_ref()
1801            .map(|c| -> Result<ChecksumPlan, String> {
1802                // A field of a variant has a slot of its own in each scope; the checksum
1803                // names one by name, so look in every variant's.
1804                let slot_of = |name: &str| -> Option<usize> {
1805                    scope.slot(name).or_else(|| {
1806                        variants
1807                            .iter()
1808                            .find_map(|v| v.fields.iter().find(|f| f.name == name).map(|f| f.slot))
1809                    })
1810                };
1811                let field = slot_of(&c.field).ok_or("checksum: no such field")?;
1812                let to =
1813                    c.to.as_deref()
1814                        .map_or(Some(field), slot_of)
1815                        .ok_or("checksum: no such field")?;
1816                let from = c
1817                    .from
1818                    .as_deref()
1819                    .map(slot_of)
1820                    .map(|s| s.ok_or("checksum: no such field"))
1821                    .transpose()?;
1822                let out = Compiler::column(
1823                    &mut columns,
1824                    "checksum_ok",
1825                    DataType::Boolean,
1826                    Sink::Flag(Vec::new()),
1827                    false,
1828                );
1829                Ok(ChecksumPlan {
1830                    algo: c.algo,
1831                    field,
1832                    from,
1833                    to,
1834                    out,
1835                })
1836            })
1837            .transpose()?;
1838        if columns.is_empty() {
1839            return Err("the records have no named field".into());
1840        }
1841        // Record size known up front: written, or what fixed fields take.
1842        let fixed_size = match size {
1843            Some(SizeRef::Given(n)) => Some(n),
1844            None if variants.is_empty() => common.iter().try_fold(0usize, |sum, f| {
1845                let width = match (&f.kind, f.count) {
1846                    (Kind::Fixed { width, .. }, None) => *width,
1847                    (Kind::Fixed { width, .. }, Some(SizeRef::Given(n))) => width.checked_mul(n)?,
1848                    (
1849                        Kind::Pad {
1850                            size: SizeRef::Given(n),
1851                        },
1852                        None,
1853                    ) => *n,
1854                    _ => return None,
1855                };
1856                sum.checked_add(width)
1857            }),
1858            _ => None,
1859        };
1860        let size = match (size, fixed_size, records.framing) {
1861            (None, Some(n), Framing::Fixed)
1862                if !common.iter().any(|f| matches!(f.kind, Kind::Group { .. })) =>
1863            {
1864                Some(SizeRef::Given(n))
1865            }
1866            (size, _, _) => size,
1867        };
1868        if matches!(size, Some(SizeRef::Given(0))) && records.sync.is_empty() {
1869            return Err("a record takes no bytes".into());
1870        }
1871        let plan = Plan {
1872            framing: records.framing,
1873            common,
1874            type_slot,
1875            variants,
1876            type_out,
1877            size,
1878            suffix,
1879            align: records.align as usize,
1880            sync: records.sync.clone(),
1881            checksum,
1882            chunk_header,
1883            chunk_count,
1884            chunk_slots,
1885            time_out,
1886            slots: compiler.slots,
1887            deltas: compiler.deltas,
1888            columns,
1889            only,
1890            sources: Vec::new(),
1891        };
1892        let plan = Plan {
1893            sources: sources(&plan),
1894            ..plan
1895        };
1896        let mut found = Found::default();
1897        let chunks = if let Some(capture) = &spec.capture {
1898            let _ = capture;
1899            crate::framed_records::capture::packets(file, &mut found.notes)?
1900        } else if let Some(blocks) = &spec.blocks {
1901            list_blocks(spec, blocks, header, file, data.clone(), &mut found.notes)?
1902        } else {
1903            vec![Chunk {
1904                source: ChunkSource::Map(data.clone()),
1905                records: None,
1906                time_ns: None,
1907            }]
1908        };
1909        let schema: Schema = plan
1910            .columns
1911            .iter()
1912            .map(|c| polars::prelude::Field::new(c.name.clone(), c.dtype.clone()))
1913            .collect();
1914        let mut records_read = Self {
1915            bytes,
1916            plan: Arc::new(plan),
1917            chunks,
1918            index: Arc::new(Index::Walk(Vec::new())),
1919            table: None,
1920            rows: 0,
1921            schema: Arc::new(schema),
1922            cache: Mutex::new(VecDeque::new()),
1923        };
1924        let ring = records
1925            .ring
1926            .as_ref()
1927            .map(|r| header.resolve_any(r, "ring"))
1928            .transpose()?;
1929        let count = records
1930            .count
1931            .as_ref()
1932            .map(|c| header.resolve_any(c, "count"))
1933            .transpose()?;
1934        // The walk of this file with this spec, kept from an earlier open of it.
1935        let kept = path
1936            .and_then(crate::indexed::peek::<KeptWalk>)
1937            .filter(|k| k.spec == *spec && k.spec.variant == spec.variant && k.data == data);
1938        if let Some(kept) = kept {
1939            records_read.index = kept.index.clone();
1940            records_read.table = kept.table.clone();
1941            records_read.rows = kept.rows;
1942            found.notes.extend(kept.notes.iter().cloned());
1943            return Ok((records_read, found.notes));
1944        }
1945        let before = found.notes.len();
1946        records_read.build_index(named, ring, count, &mut found)?;
1947        if found.skipped > 0 {
1948            found.notes.push(format!(
1949                "{} {} skipped between records, to the next sync marker",
1950                found.skipped,
1951                if found.skipped == 1 { "byte" } else { "bytes" }
1952            ));
1953        }
1954        if let Some(path) = path
1955            && matches!(*records_read.index, Index::Walk(_))
1956        {
1957            crate::indexed::keep(
1958                path,
1959                Arc::new(KeptWalk {
1960                    spec: spec.clone(),
1961                    data,
1962                    index: records_read.index.clone(),
1963                    table: records_read.table.clone(),
1964                    rows: records_read.rows,
1965                    notes: found.notes[before..].to_vec(),
1966                }),
1967            );
1968        }
1969        Ok((records_read, found.notes))
1970    }
1971
1972    /// Whether a stride of `size` bytes finds every record: one run of fixed records.
1973    fn stride(&self) -> Option<usize> {
1974        let plan = &self.plan;
1975        match (plan.size, &self.chunks[..]) {
1976            (Some(SizeRef::Given(n)), [chunk])
1977                if plan.framing == Framing::Fixed
1978                    && plan.only.is_none()
1979                    && plan.sync.is_empty()
1980                    && plan.suffix.is_none()
1981                    // A running sum needs every record before the one read.
1982                    && plan.deltas.is_empty()
1983                    && matches!(chunk.source, ChunkSource::Map(_))
1984                    && plan.chunk_header.is_empty() =>
1985            {
1986                let align = plan.align.max(1);
1987                Some(n.div_ceil(align) * align)
1988            }
1989            _ => None,
1990        }
1991    }
1992
1993    fn build_index(
1994        &mut self,
1995        named: &str,
1996        ring: Option<u64>,
1997        limit: Option<u64>,
1998        found: &mut Found,
1999    ) -> Result<(), String> {
2000        let file = self.bytes.as_slice();
2001        if let Some(size) = self.stride() {
2002            let ChunkSource::Map(range) = self.chunks[0].source.clone() else {
2003                unreachable!("a stride is over the map")
2004            };
2005            let len = range.len();
2006            let mut rows = len / size;
2007            // The last record need not be padded out to the alignment.
2008            let unpadded = self.plan.size.map_or(size, |s| match s {
2009                SizeRef::Given(n) => n,
2010                _ => size,
2011            });
2012            if len % size >= unpadded {
2013                rows += 1;
2014            }
2015            if let Some(count) = limit {
2016                if count < rows as u64 {
2017                    rows = count as usize;
2018                } else if count > rows as u64 {
2019                    found.notes.push(format!(
2020                        "header says {count} records {} {rows} whole ones shown",
2021                        crate::glyphs::get().middot
2022                    ));
2023                }
2024            } else {
2025                let used = (rows * size).min(len);
2026                let trailing = len - used;
2027                if trailing > 0 && len % size < unpadded {
2028                    found
2029                        .notes
2030                        .push(trailing_note(named, &file[range.end - trailing..range.end]));
2031                }
2032            }
2033            let rows = rows.min(MAX_ROWS);
2034            let ring = match ring {
2035                Some(r) if rows > 0 => {
2036                    if r >= rows as u64 {
2037                        found.notes.push(format!(
2038                            "ring's oldest record {r} past the {rows} records {} read from the first",
2039                            crate::glyphs::get().middot
2040                        ));
2041                        0
2042                    } else {
2043                        r as usize
2044                    }
2045                }
2046                _ => 0,
2047            };
2048            self.index = Arc::new(Index::Stride {
2049                start: range.start,
2050                size,
2051                ring,
2052            });
2053            self.rows = rows;
2054            return Ok(());
2055        }
2056        let plan = self.plan.clone();
2057        let mut walker = Walker::new(&plan, file);
2058        let mut checkpoints = Vec::new();
2059        let mut rows: u64 = 0;
2060        let max_rows = MAX_ROWS as u64;
2061        let needs_walk_everything = plan.deltas.contains(&Delta::All);
2062        let mut table = (self.chunks.len() == 1
2063            && matches!(self.chunks[0].source, ChunkSource::Map(_))
2064            && plan.chunk_header.is_empty()
2065            && plan.variants.len() < usize::from(WALK))
2066        .then(|| RowTable {
2067            starts: crate::indexed::Offsets::for_file(file.len()),
2068            tags: Vec::new(),
2069        });
2070        'chunks: for ci in 0..self.chunks.len() {
2071            if limit.is_some_and(|l| rows >= l) || rows >= max_rows {
2072                break;
2073            }
2074            reset_block_sums(&plan, &mut walker.acc);
2075            if let Some(n) = self.chunks[ci].records
2076                && !needs_walk_everything
2077            {
2078                checkpoints.push(Checkpoint {
2079                    row: rows,
2080                    chunk: ci as u32,
2081                    pos: u32::MAX,
2082                    taken: 0,
2083                    acc: walker.acc.clone().into_boxed_slice(),
2084                });
2085                rows = rows
2086                    .saturating_add(n)
2087                    .min(limit.unwrap_or(u64::MAX))
2088                    .min(max_rows);
2089                continue;
2090            }
2091            let (data, mut cursor) = match self.enter(ci, &mut walker) {
2092                Ok(c) => c,
2093                Err(e) => {
2094                    found.notes.push(format!("block {ci}: {e}; left out"));
2095                    continue;
2096                }
2097            };
2098            let chunk_bytes = self.slice(&data, file);
2099            loop {
2100                if cursor.limit.is_some_and(|l| cursor.taken >= l)
2101                    || limit.is_some_and(|l| rows >= l)
2102                    || rows >= max_rows
2103                {
2104                    break;
2105                }
2106                if cursor.taken.is_multiple_of(CHECKPOINT) {
2107                    checkpoints.push(Checkpoint {
2108                        row: rows,
2109                        chunk: ci as u32,
2110                        pos: cursor.pos as u32,
2111                        taken: cursor.taken as u32,
2112                        acc: walker.acc.clone().into_boxed_slice(),
2113                    });
2114                }
2115                let before = cursor.pos;
2116                match walker.record(chunk_bytes, &mut cursor.pos, cursor.end, None, None) {
2117                    Ok(got @ (Got::Row | Got::Skipped)) => {
2118                        if got == Got::Row {
2119                            rows += 1;
2120                            if let Some(t) = table.as_mut() {
2121                                if t.tags.len() >= crate::indexed::MAX_RECORDS {
2122                                    table = None;
2123                                } else {
2124                                    let tag = match (walker.short, plan.type_slot) {
2125                                        (true, _) => WALK,
2126                                        (false, None) => 0,
2127                                        (false, Some(_)) => walker
2128                                            .variant
2129                                            .and_then(|v| u8::try_from(v).ok())
2130                                            .unwrap_or(WALK),
2131                                    };
2132                                    t.starts.push(walker.record_start - plan.sync.len());
2133                                    t.tags.push(tag);
2134                                }
2135                            }
2136                        }
2137                        cursor.taken += 1;
2138                        align(&mut cursor, plan.align);
2139                        if cursor.pos <= before {
2140                            found.notes.push(format!(
2141                                "zero-length record at byte {before} {} rest left out",
2142                                crate::glyphs::get().middot
2143                            ));
2144                            break 'chunks;
2145                        }
2146                    }
2147                    Ok(Got::None) => {
2148                        // A checkpoint pushed for a record that is not there.
2149                        if checkpoints.last().is_some_and(|c: &Checkpoint| {
2150                            c.row == rows && c.chunk == ci as u32 && c.pos == before as u32
2151                        }) {
2152                            checkpoints.pop();
2153                        }
2154                        break;
2155                    }
2156                    Err(stop) => {
2157                        if checkpoints.last().is_some_and(|c: &Checkpoint| {
2158                            c.row == rows && c.chunk == ci as u32 && c.pos == before as u32
2159                        }) {
2160                            checkpoints.pop();
2161                        }
2162                        let what = if self.chunks.len() > 1 {
2163                            format!("{named}, block {ci},")
2164                        } else {
2165                            named.to_string()
2166                        };
2167                        match stop {
2168                            Stop::Truncated => {
2169                                let rest = &chunk_bytes[before..cursor.end];
2170                                found.notes.push(trailing_note(&what, rest));
2171                            }
2172                            Stop::Said(said) => {
2173                                found.notes.push(format!(
2174                                    "{said} {} {} bytes from there left out",
2175                                    crate::glyphs::get().middot,
2176                                    cursor.end - before
2177                                ));
2178                            }
2179                        }
2180                        if self.chunks.len() == 1 {
2181                            break 'chunks;
2182                        }
2183                        break;
2184                    }
2185                }
2186                if cursor.pos >= cursor.end {
2187                    break;
2188                }
2189            }
2190            // A chunk that counts its records and holds fewer is filled with nulls.
2191            if let Some(l) = cursor.limit
2192                && cursor.taken < l
2193                && self.chunks[ci].records.is_some()
2194            {
2195                found.notes.push(format!(
2196                    "block {ci}: {l} records declared, {} found",
2197                    cursor.taken
2198                ));
2199            }
2200        }
2201        if let Some(l) = limit
2202            && rows < l
2203        {
2204            found.notes.push(format!(
2205                "header says {l} records {} {rows} shown",
2206                crate::glyphs::get().middot
2207            ));
2208        }
2209        found.skipped += walker.skipped;
2210        self.rows = rows as usize;
2211        self.index = Arc::new(Index::Walk(checkpoints));
2212        self.table = table.filter(|t| t.tags.len() == self.rows).map(|mut t| {
2213            t.starts.shrink();
2214            t.tags.shrink_to_fit();
2215            Arc::new(t)
2216        });
2217        Ok(())
2218    }
2219
2220    fn slice<'b>(&'b self, data: &'b ChunkData, file: &'b [u8]) -> &'b [u8] {
2221        match data {
2222            ChunkData::Map(_) => file,
2223            ChunkData::Owned(v) => v,
2224        }
2225    }
2226
2227    /// The bytes of chunk `i`: its range of the map, or its block decompressed.
2228    fn chunk_data(&self, i: usize) -> Result<ChunkData, String> {
2229        match &self.chunks[i].source {
2230            ChunkSource::Map(range) => Ok(ChunkData::Map(range.clone())),
2231            ChunkSource::Block {
2232                body,
2233                codec,
2234                uncompressed,
2235            } => {
2236                if let Ok(cache) = self.cache.lock()
2237                    && let Some((_, data)) = cache.iter().find(|(k, _)| *k == i)
2238                {
2239                    return Ok(ChunkData::Owned(data.clone()));
2240                }
2241                self.bytes.still_whole().map_err(|e| e.to_string())?;
2242                let raw = &self.bytes.as_slice()[body.clone()];
2243                let data = Arc::new(decompress(raw, *codec, *uncompressed)?);
2244                if let Ok(mut cache) = self.cache.lock() {
2245                    cache.push_front((i, data.clone()));
2246                    cache.truncate(CACHED_BLOCKS);
2247                }
2248                Ok(ChunkData::Owned(data))
2249            }
2250        }
2251    }
2252
2253    /// Chunk `i`'s bytes, and a cursor at its start, after its payload header.
2254    fn enter(&self, i: usize, walker: &mut Walker<'_>) -> Result<(ChunkData, Cursor), String> {
2255        let data = self.chunk_data(i)?;
2256        let (start, end) = match &data {
2257            ChunkData::Map(r) => (r.start, r.end),
2258            ChunkData::Owned(v) => (0, v.len()),
2259        };
2260        let mut cursor = Cursor {
2261            pos: start,
2262            data_start: start,
2263            end,
2264            taken: 0,
2265            limit: self.chunks[i].records,
2266            time: self.chunks[i].time_ns,
2267        };
2268        if !self.plan.chunk_header.is_empty() {
2269            let file = self.bytes.as_slice();
2270            let bytes = self.slice(&data, file);
2271            walker.frame.clear(0..self.plan.chunk_slots);
2272            walker.end = end;
2273            walker.bounded = false;
2274            walker.short = false;
2275            walker.size_slot = None;
2276            walker.record_start = start;
2277            let mut p = start;
2278            match walker.walk(&self.plan.chunk_header, bytes, &mut p, None) {
2279                Ok(()) => {
2280                    cursor.pos = p;
2281                    cursor.data_start = p;
2282                    if let Some(slot) = self.plan.chunk_count {
2283                        cursor.limit = walker.frame.ints[slot].map(|v| v.max(0) as u64);
2284                    }
2285                }
2286                Err(_) => {
2287                    // Too short for its header: nothing in it.
2288                    cursor.pos = end;
2289                    cursor.limit = Some(0);
2290                }
2291            }
2292        }
2293        Ok((data, cursor))
2294    }
2295
2296    /// Rows `[start, start + len)`, of the columns `wanted` names, decoded now.
2297    fn decode(&self, start: usize, len: usize, wanted: &[bool]) -> PolarsResult<DataFrame> {
2298        self.bytes.still_whole()?;
2299        let start = start.min(self.rows);
2300        let len = len.min(self.rows - start);
2301        let plan = &*self.plan;
2302        let mut sinks: Vec<Sink> = plan
2303            .columns
2304            .iter()
2305            .zip(wanted)
2306            .map(|(c, w)| if *w { c.proto.clone() } else { Sink::Skip })
2307            .collect();
2308        let mut filled = vec![false; sinks.len()];
2309        let file = self.bytes.as_slice();
2310        let mut walker = Walker::new(plan, file);
2311        match &*self.index {
2312            Index::Stride {
2313                start: base,
2314                size,
2315                ring,
2316            } => {
2317                let mut out = Out {
2318                    sinks: &mut sinks,
2319                    filled: &mut filled,
2320                };
2321                let ChunkSource::Map(range) = &self.chunks[0].source else {
2322                    unreachable!("a stride is over the map")
2323                };
2324                for row in start..start + len {
2325                    let k = (ring + row) % self.rows.max(1);
2326                    let mut pos = base + k * size;
2327                    let end = (pos + size).min(range.end);
2328                    if walker
2329                        .record(file, &mut pos, end, Some(&mut out), None)
2330                        .is_err()
2331                    {
2332                        out.finish_row();
2333                    }
2334                }
2335            }
2336            Index::Walk(checkpoints) => {
2337                let at = checkpoints.partition_point(|c| c.row <= start as u64);
2338                let mut want = Want {
2339                    start: start as u64,
2340                    len,
2341                    produced: 0,
2342                };
2343                let mut out = Out {
2344                    sinks: &mut sinks,
2345                    filled: &mut filled,
2346                };
2347                if let Some(cp) = at.checked_sub(1).map(|i| &checkpoints[i]) {
2348                    self.read_from(cp, &mut walker, &mut want, &mut out);
2349                }
2350                // Rows the index counted that a read now does not find: a file changed
2351                // under the map, or a block that no longer decompresses.
2352                for _ in want.produced..len {
2353                    out.finish_row();
2354                }
2355            }
2356        }
2357        let columns = sinks
2358            .into_iter()
2359            .zip(&plan.columns)
2360            .zip(wanted)
2361            .filter(|(_, w)| **w)
2362            .map(|((sink, c), _)| sink.finish(c.name.clone()).map(Column::from))
2363            .collect::<PolarsResult<Vec<_>>>()?;
2364        DataFrame::new(len, columns)
2365    }
2366
2367    /// From checkpoint `cp`, skip to the row `want` starts at and read its rows.
2368    fn read_from(
2369        &self,
2370        cp: &Checkpoint,
2371        walker: &mut Walker<'_>,
2372        want: &mut Want,
2373        out: &mut Out<'_>,
2374    ) {
2375        let file = self.bytes.as_slice();
2376        let plan = &*self.plan;
2377        walker.acc.copy_from_slice(&cp.acc);
2378        let mut row = cp.row;
2379        let mut ci = cp.chunk as usize;
2380        let mut first = true;
2381        while want.produced < want.len && ci < self.chunks.len() {
2382            if !first {
2383                reset_block_sums(plan, &mut walker.acc);
2384            }
2385            let (data, mut cursor) = match self.enter(ci, walker) {
2386                Ok(c) => c,
2387                Err(_) => {
2388                    // A block that will not decompress now: its promised rows are null.
2389                    let n = self.chunks[ci].records.unwrap_or(0);
2390                    for _ in 0..n {
2391                        want.null_row(row, out);
2392                        row += 1;
2393                    }
2394                    ci += 1;
2395                    first = false;
2396                    continue;
2397                }
2398            };
2399            if first && cp.pos != u32::MAX {
2400                cursor.pos = cp.pos as usize;
2401                cursor.taken = u64::from(cp.taken);
2402            }
2403            first = false;
2404            let bytes = self.slice(&data, file);
2405            while want.produced < want.len {
2406                if cursor.limit.is_some_and(|l| cursor.taken >= l) || cursor.pos >= cursor.end {
2407                    break;
2408                }
2409                let reading = row >= want.start;
2410                let got = walker.record(
2411                    bytes,
2412                    &mut cursor.pos,
2413                    cursor.end,
2414                    if reading { Some(&mut *out) } else { None },
2415                    cursor.time,
2416                );
2417                match got {
2418                    Ok(Got::Row) => {
2419                        if reading {
2420                            want.produced += 1;
2421                        }
2422                        row += 1;
2423                        cursor.taken += 1;
2424                        align(&mut cursor, plan.align);
2425                    }
2426                    Ok(Got::Skipped) => {
2427                        cursor.taken += 1;
2428                        align(&mut cursor, plan.align);
2429                    }
2430                    Ok(Got::None) | Err(_) => break,
2431                }
2432            }
2433            // A chunk that promised more records than it held: nulls for the rest.
2434            if let Some(l) = self.chunks[ci].records {
2435                while cursor.taken < l && want.produced < want.len {
2436                    want.null_row(row, out);
2437                    row += 1;
2438                    cursor.taken += 1;
2439                }
2440            }
2441            ci += 1;
2442        }
2443    }
2444
2445    /// Column `column` of `rows`, each read where the row table says its record starts.
2446    fn decode_from(
2447        &self,
2448        table: &RowTable,
2449        column: usize,
2450        rows: &[IdxSize],
2451    ) -> PolarsResult<Column> {
2452        self.bytes.still_whole()?;
2453        let plan = &*self.plan;
2454        let file = self.bytes.as_slice();
2455        let ChunkSource::Map(range) = &self.chunks[0].source else {
2456            unreachable!("a row table is over the map")
2457        };
2458        let mut sinks: Vec<Sink> = vec![Sink::Skip; plan.columns.len()];
2459        sinks[column] = plan.columns[column].proto.clone();
2460        let mut filled = vec![false; sinks.len()];
2461        let mut out = Out {
2462            sinks: &mut sinks,
2463            filled: &mut filled,
2464        };
2465        let mut walker = Walker::new(plan, file);
2466        if plan.type_out == Some(column) {
2467            // Each row's variant is its tag: the names are cast once and the rows
2468            // gathered by tag. A row read field by field says its own label.
2469            let mut labels: Vec<Option<Arc<str>>> =
2470                plan.variants.iter().map(|v| Some(v.name.clone())).collect();
2471            let codes: IdxCa = rows
2472                .iter()
2473                .map(|&row| {
2474                    let row = row as usize;
2475                    let tag = table.tags[row];
2476                    if tag != WALK {
2477                        return Some(IdxSize::from(tag));
2478                    }
2479                    let mut pos = table.starts.get(row);
2480                    let _ = walker.record(file, &mut pos, range.end, Some(&mut out), None);
2481                    let Sink::Label(said) = &mut out.sinks[column] else {
2482                        return None;
2483                    };
2484                    let label = said.pop().flatten()?;
2485                    labels.push(Some(label));
2486                    Some((labels.len() - 1) as IdxSize)
2487                })
2488                .collect();
2489            let name = plan.columns[column].name.clone();
2490            let labels = Sink::Label(labels).finish(name)?;
2491            return Ok(labels.take(&codes)?.into_column());
2492        }
2493        let sources = &plan.sources[column];
2494        for &row in rows {
2495            let row = row as usize;
2496            let start = table.starts.get(row);
2497            let tag = table.tags[row];
2498            let source = match tag {
2499                WALK => &Source::Walk,
2500                v => &sources[usize::from(v)],
2501            };
2502            match source {
2503                Source::Null => out.finish_row(),
2504                Source::Label => {
2505                    if let Sink::Label(v) = &mut out.sinks[column] {
2506                        v.push(plan.variants.get(usize::from(tag)).map(|v| v.name.clone()));
2507                    }
2508                    out.filled[column] = true;
2509                    out.finish_row();
2510                }
2511                Source::At { offset, field } => {
2512                    let Kind::Fixed { width, int } = field.kind else {
2513                        unreachable!("a value at a place is fixed")
2514                    };
2515                    let cells = match field.count {
2516                        Some(SizeRef::Given(n)) => Some(n),
2517                        _ => None,
2518                    };
2519                    let at = start + offset;
2520                    let bytes = at
2521                        .checked_add(width * cells.unwrap_or(1))
2522                        .filter(|end| *end <= range.end)
2523                        .map(|end| &file[at..end]);
2524                    if let Some(bytes) = bytes {
2525                        let raw = int
2526                            .and_then(|r| (cells != Some(0)).then(|| int_of(&bytes[..width], r)));
2527                        push_fixed(field, bytes, width, cells, &mut out);
2528                        push_bits(field, raw, &mut out);
2529                    }
2530                    out.finish_row();
2531                }
2532                Source::Walk | Source::Summed => {
2533                    let mut pos = start;
2534                    // A row the walk does not finish (the file changed under the map)
2535                    // is null rather than missing.
2536                    if !matches!(
2537                        walker.record(file, &mut pos, range.end, Some(&mut out), None),
2538                        Ok(Got::Row)
2539                    ) {
2540                        out.finish_row();
2541                    }
2542                }
2543            }
2544        }
2545        let name = plan.columns[column].name.clone();
2546        Ok(sinks.swap_remove(column).finish(name)?.into_column())
2547    }
2548
2549    pub fn rows(&self) -> usize {
2550        self.rows
2551    }
2552
2553    pub fn schema(&self) -> SchemaRef {
2554        self.schema.clone()
2555    }
2556
2557    pub fn sources(&self) -> &[Arc<Bytes>] {
2558        std::slice::from_ref(&self.bytes)
2559    }
2560}
2561
2562impl FramedRecords {
2563    /// The frame: decoded over a row index, only what a query reaches.
2564    pub fn lazy(self: &Arc<Self>) -> LazyFrame {
2565        crate::row_index::lazy(self)
2566    }
2567
2568    /// Rows `[start, start + len)` as a frame of their own, read from the nearest
2569    /// checkpoint before them.
2570    pub fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
2571        let all = vec![true; self.plan.columns.len()];
2572        Ok(self.decode(start, len, &all)?.lazy())
2573    }
2574
2575    /// The first `rows` rows of every column, decoded now.
2576    pub fn collect(&self, rows: usize) -> PolarsResult<DataFrame> {
2577        let all = vec![true; self.plan.columns.len()];
2578        self.decode(0, rows, &all)
2579    }
2580}
2581
2582impl crate::row_index::RowSource for FramedRecords {
2583    fn height(&self) -> usize {
2584        self.rows
2585    }
2586
2587    fn schema(&self) -> SchemaRef {
2588        self.schema.clone()
2589    }
2590
2591    // Each column is decoded on its own, so a query that names one column decodes only
2592    // that one; the streaming engine asks for a morsel at a time. From the row table,
2593    // a column reads each row where its record starts, as every other column does;
2594    // without one, each column walks the span its rows cover.
2595    fn decode(&self, column: usize, index: &IdxCa) -> PolarsResult<Column> {
2596        let rows = crate::row_index::checked(index, self.rows)?;
2597        if let Some(table) = &self.table
2598            && let Some(sources) = self.plan.sources.get(column)
2599            && !sources.iter().any(|s| matches!(s, Source::Summed))
2600        {
2601            return self.decode_from(table, column, &rows);
2602        }
2603        let mut wanted = vec![false; self.plan.columns.len()];
2604        *wanted
2605            .get_mut(column)
2606            .ok_or_else(|| polars_err!(OutOfBounds: "no column {column}"))? = true;
2607        let (Some(&lo), Some(&hi)) = (rows.iter().min(), rows.iter().max()) else {
2608            return Ok(self.decode(0, 0, &wanted)?.columns()[0].clone());
2609        };
2610        let (lo, span) = (lo as usize, (hi - lo) as usize + 1);
2611        let values = self.decode(lo, span, &wanted)?.columns()[0].clone();
2612        let contiguous = rows.len() == span && rows.windows(2).all(|w| w[1] == w[0] + 1);
2613        if contiguous {
2614            return Ok(values);
2615        }
2616        let at = IdxCa::from_vec(
2617            PlSmallStr::EMPTY,
2618            rows.iter().map(|&r| r - lo as IdxSize).collect(),
2619        );
2620        values.take(&at)
2621    }
2622}
2623
2624impl crate::pushdown::Windowed for FramedRecords {
2625    fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
2626        FramedRecords::window(self, start, len)
2627    }
2628}
2629
2630/// Put `cursor` at the next multiple of `align` from the chunk's start.
2631fn align(cursor: &mut Cursor, align: usize) {
2632    if align > 1 {
2633        let offset = cursor.pos - cursor.data_start;
2634        let rounded = offset.div_ceil(align).saturating_mul(align);
2635        cursor.pos = cursor.data_start.saturating_add(rounded).min(cursor.end);
2636    }
2637}
2638
2639/// Running sums that start again with each block.
2640fn reset_block_sums(plan: &Plan, acc: &mut [i128]) {
2641    for (sum, delta) in acc.iter_mut().zip(&plan.deltas) {
2642        if *delta == Delta::Block {
2643            *sum = 0;
2644        }
2645    }
2646}
2647
2648/// A warning naming the bytes at the end that are not a whole record.
2649fn trailing_note(what: &str, bytes: &[u8]) -> String {
2650    formats::trailing_note(what, bytes)
2651}
2652
2653impl Plan {
2654    /// A plan of `fields` alone, to read a block header or an index entry.
2655    fn bare(fields: Vec<FieldPlan>, slots: usize) -> Self {
2656        Self {
2657            framing: Framing::Fixed,
2658            common: fields,
2659            type_slot: None,
2660            variants: Vec::new(),
2661            type_out: None,
2662            size: None,
2663            suffix: None,
2664            align: 1,
2665            sync: Vec::new(),
2666            checksum: None,
2667            chunk_header: Vec::new(),
2668            chunk_count: None,
2669            chunk_slots: 0,
2670            time_out: None,
2671            slots,
2672            deltas: Vec::new(),
2673            columns: Vec::new(),
2674            only: None,
2675            sources: Vec::new(),
2676        }
2677    }
2678}
2679
2680/// Where each column's value is in a record of each variant: at a fixed place while
2681/// every field before it in the record is fixed in size, else found by a walk.
2682fn sources(plan: &Plan) -> Vec<Vec<Source>> {
2683    let variants = plan.variants.len().max(1);
2684    let mut out = vec![vec![Source::Null; variants]; plan.columns.len()];
2685    // Read alone, the other variants' fields fill no column of this table.
2686    for v in (0..variants).filter(|v| plan.only.is_none_or(|only| only == *v)) {
2687        let fields = plan
2688            .common
2689            .iter()
2690            .chain(plan.variants.get(v).into_iter().flat_map(|p| &p.fields));
2691        // A row's start is its sync marker's.
2692        let mut offset = Some(plan.sync.len());
2693        for f in fields {
2694            let width = match (&f.kind, f.count) {
2695                (Kind::Fixed { width, .. }, None) => Some(*width),
2696                (Kind::Fixed { width, .. }, Some(SizeRef::Given(n))) => width.checked_mul(n),
2697                (
2698                    Kind::Pad {
2699                        size: SizeRef::Given(n),
2700                    },
2701                    None,
2702                ) => Some(*n),
2703                _ => None,
2704            };
2705            let source = match (offset, &f.kind, width) {
2706                (Some(offset), Kind::Fixed { .. }, Some(_)) if f.delta == Delta::None => {
2707                    Source::At {
2708                        offset,
2709                        field: f.clone(),
2710                    }
2711                }
2712                _ => Source::Walk,
2713            };
2714            for &c in &f.outs {
2715                out[c][v] = if f.delta == Delta::None {
2716                    source.clone()
2717                } else {
2718                    Source::Summed
2719                };
2720            }
2721            for (c, _, _) in &f.bits {
2722                out[*c][v] = source.clone();
2723            }
2724            offset = offset.zip(width).and_then(|(o, w)| o.checked_add(w));
2725        }
2726    }
2727    if let Some(c) = plan.type_out {
2728        out[c].fill(Source::Label);
2729    }
2730    for c in plan.checksum.iter().map(|c| c.out).chain(plan.time_out) {
2731        out[c].fill(Source::Walk);
2732    }
2733    out
2734}
2735
2736/// Fields read once, for their values: a block header, an index entry.
2737struct Struct {
2738    plan: Plan,
2739    scope: Scope,
2740}
2741
2742impl Struct {
2743    fn compile(
2744        spec: &Spec,
2745        header: &HeaderValues,
2746        file: &[u8],
2747        fields: &[Field],
2748    ) -> Result<Self, String> {
2749        let mut compiler = Compiler {
2750            spec,
2751            header,
2752            data: file,
2753            slots: 0,
2754            deltas: Vec::new(),
2755        };
2756        let mut scope = Scope::default();
2757        let mut columns = Vec::new();
2758        let plans = compiler.fields(fields, &mut scope, &mut columns, false)?;
2759        Ok(Self {
2760            plan: Plan::bare(plans, compiler.slots),
2761            scope,
2762        })
2763    }
2764
2765    /// Read the fields at `pos`, before `end`: the frame, and where they end.
2766    fn read(&self, file: &[u8], pos: usize, end: usize) -> Result<(Frame, usize), Stop> {
2767        let mut walker = Walker::new(&self.plan, file);
2768        walker.end = end;
2769        walker.record_start = pos;
2770        let mut p = pos;
2771        walker.walk(&self.plan.common, file, &mut p, None)?;
2772        Ok((walker.frame, p))
2773    }
2774
2775    fn value(&self, frame: &Frame, name: &str) -> Option<i128> {
2776        frame.ints[self.scope.slot(name)?]
2777    }
2778}
2779
2780/// The blocks of `data`: from the index the file keeps, or by walking their headers.
2781fn list_blocks(
2782    spec: &Spec,
2783    blocks: &formats::Blocks,
2784    header: &HeaderValues,
2785    file: &[u8],
2786    data: Range<usize>,
2787    notes: &mut Vec<String>,
2788) -> Result<Vec<Chunk>, String> {
2789    let head = Struct::compile(spec, header, file, &blocks.header)?;
2790    let size = {
2791        let compiler = Compiler {
2792            spec,
2793            header,
2794            data: file,
2795            slots: 0,
2796            deltas: Vec::new(),
2797        };
2798        compiler.size(&blocks.size, &head.scope, "block size")?
2799    };
2800    let codec_of = |frame: &Frame| -> Result<Compression, String> {
2801        match &blocks.codec {
2802            Codec::Fixed(c) => Ok(*c),
2803            Codec::ByField { field, values } => {
2804                let code = head
2805                    .value(frame, field)
2806                    .ok_or("the block header has no codec")?;
2807                i64::try_from(code)
2808                    .ok()
2809                    .and_then(|c| values.get(&c).copied())
2810                    .ok_or_else(|| format!("codec {code} is not one compression names"))
2811            }
2812        }
2813    };
2814    let mut chunks = Vec::new();
2815    let read_block = |at: usize,
2816                      end: usize,
2817                      rows: Option<u64>,
2818                      chunks: &mut Vec<Chunk>,
2819                      notes: &mut Vec<String>|
2820     -> Result<Option<usize>, String> {
2821        let (frame, body_start) = match head.read(file, at, end) {
2822            Ok(x) => x,
2823            Err(_) => {
2824                notes.push(format!(
2825                    "the block header at byte {at} runs past the data; the rest is left out"
2826                ));
2827                return Ok(None);
2828            }
2829        };
2830        let body_len = match size {
2831            SizeRef::Given(n) => n,
2832            SizeRef::Slot { slot, adjust } => {
2833                let v = frame.ints[slot].unwrap_or(0) + i128::from(adjust);
2834                usize::try_from(v)
2835                    .ok()
2836                    .filter(|n| *n <= MAX_BLOCK)
2837                    .ok_or_else(|| {
2838                        format!(
2839                            "the block at byte {at} gives its size as {v}, outside 0 to {MAX_BLOCK}"
2840                        )
2841                    })?
2842            }
2843            SizeRef::Rest => end - body_start,
2844        };
2845        let body_end = body_start.checked_add(body_len).filter(|e| *e <= end);
2846        let Some(body_end) = body_end else {
2847            notes.push(format!(
2848                "the block at byte {at} is {body_len} bytes and runs past the data; the rest is left out"
2849            ));
2850            return Ok(None);
2851        };
2852        let codec = match codec_of(&frame) {
2853            Ok(c) => c,
2854            Err(e) => {
2855                notes.push(format!("the block at byte {at}: {e}; left out"));
2856                return Ok(Some(body_end));
2857            }
2858        };
2859        let records = rows.or_else(|| {
2860            blocks
2861                .records
2862                .as_ref()
2863                .and_then(|r| head.value(&frame, r))
2864                .map(|v| v.max(0) as u64)
2865        });
2866        let uncompressed = blocks
2867            .uncompressed
2868            .as_ref()
2869            .and_then(|u| head.value(&frame, u))
2870            .and_then(|v| usize::try_from(v).ok());
2871        chunks.push(Chunk {
2872            source: if codec == Compression::None {
2873                ChunkSource::Map(body_start..body_end)
2874            } else {
2875                ChunkSource::Block {
2876                    body: body_start..body_end,
2877                    codec,
2878                    uncompressed,
2879                }
2880            },
2881            records,
2882            time_ns: None,
2883        });
2884        Ok(Some(body_end))
2885    };
2886    match &blocks.index {
2887        Some(index) => {
2888            let entry = Struct::compile(spec, header, file, &index.fields)?;
2889            let at = usize::try_from(header.resolve_any(&index.at, "index at")?)
2890                .map_err(|_| "index: too far")?;
2891            let count = header.resolve_any(&index.count, "index count")?;
2892            let width = formats::fields_width(&index.fields).unwrap_or(1).max(1) as usize;
2893            let fits = file.len().saturating_sub(at) / width;
2894            if count > fits as u64 {
2895                return Err(format!(
2896                    "the block index at byte {at} says {count} entries; the file has room for {fits}"
2897                ));
2898            }
2899            let mut pos = at;
2900            for _ in 0..count {
2901                let (frame, next) = head_or(entry.read(file, pos, file.len()))?;
2902                pos = next;
2903                let offset = entry.value(&frame, "offset").unwrap_or(-1);
2904                let rows = entry.value(&frame, "rows").map(|v| v.max(0) as u64);
2905                let Some(offset) = usize::try_from(offset).ok().filter(|o| *o < file.len()) else {
2906                    notes.push(format!(
2907                        "an index entry points at byte {offset}, outside the file; left out"
2908                    ));
2909                    continue;
2910                };
2911                read_block(offset, file.len(), rows, &mut chunks, notes)?;
2912            }
2913        }
2914        None => {
2915            let mut pos = data.start;
2916            while pos < data.end {
2917                match read_block(pos, data.end, None, &mut chunks, notes)? {
2918                    Some(next) if next > pos => pos = next,
2919                    Some(_) => {
2920                        notes.push(format!(
2921                            "zero-length block at byte {pos} {} rest left out",
2922                            crate::glyphs::get().middot
2923                        ));
2924                        break;
2925                    }
2926                    None => break,
2927                }
2928            }
2929        }
2930    }
2931    Ok(chunks)
2932}
2933
2934fn head_or(read: Result<(Frame, usize), Stop>) -> Result<(Frame, usize), String> {
2935    read.map_err(|_| "an index entry runs past the end of the file".to_string())
2936}
2937
2938/// `raw` decompressed by `codec`, at most [`MAX_BLOCK`] bytes.
2939pub fn decompress(
2940    raw: &[u8],
2941    codec: Compression,
2942    uncompressed: Option<usize>,
2943) -> Result<Vec<u8>, String> {
2944    use std::io::Read;
2945    let read_all = |mut reader: Box<dyn Read + '_>| -> Result<Vec<u8>, String> {
2946        let mut out = Vec::new();
2947        reader
2948            .by_ref()
2949            .take(MAX_BLOCK as u64 + 1)
2950            .read_to_end(&mut out)
2951            .map_err(|e| format!("{} block: {e}", codec.name()))?;
2952        if out.len() > MAX_BLOCK {
2953            return Err(format!(
2954                "a block decompresses to more than {MAX_BLOCK} bytes"
2955            ));
2956        }
2957        Ok(out)
2958    };
2959    match codec {
2960        Compression::None => Ok(raw.to_vec()),
2961        Compression::Gzip => read_all(Box::new(flate2::read::MultiGzDecoder::new(raw))),
2962        Compression::Deflate => read_all(Box::new(flate2::read::DeflateDecoder::new(raw))),
2963        Compression::Zlib => read_all(Box::new(flate2::read::ZlibDecoder::new(raw))),
2964        Compression::Zstd => read_all(Box::new(
2965            zstd::Decoder::new(raw).map_err(|e| format!("zstd block: {e}"))?,
2966        )),
2967        Compression::Lz4 => read_all(Box::new(
2968            lz4::Decoder::new(raw).map_err(|e| format!("lz4 block: {e}"))?,
2969        )),
2970        Compression::Lz4Block => {
2971            let size = uncompressed.filter(|n| *n <= MAX_BLOCK).ok_or_else(|| {
2972                format!("an lz4 block needs its decompressed size, at most {MAX_BLOCK}")
2973            })?;
2974            lz4::block::decompress(raw, Some(size as i32)).map_err(|e| format!("lz4 block: {e}"))
2975        }
2976        Compression::Snappy => {
2977            let size = snap::raw::decompress_len(raw).map_err(|e| format!("snappy block: {e}"))?;
2978            if size > MAX_BLOCK {
2979                return Err(format!(
2980                    "a block decompresses to more than {MAX_BLOCK} bytes"
2981                ));
2982            }
2983            snap::raw::Decoder::new()
2984                .decompress_vec(raw)
2985                .map_err(|e| format!("snappy block: {e}"))
2986        }
2987        Compression::SnappyFramed => read_all(Box::new(snap::read::FrameDecoder::new(raw))),
2988        Compression::Brotli => read_all(Box::new(brotli::Decompressor::new(raw, 4096))),
2989        Compression::Bzip2 => read_all(Box::new(bzip2::read::BzDecoder::new(raw))),
2990        Compression::Xz => read_all(Box::new(xz2::read::XzDecoder::new(raw))),
2991    }
2992}
2993
2994/// Packet captures (pcap and pcapng): the UDP payload of each packet, with its time.
2995pub mod capture {
2996    use super::{Chunk, ChunkSource};
2997
2998    /// The most packets one capture is read for.
2999    const MAX_PACKETS: usize = 64 << 20;
3000
3001    fn u16_at(b: &[u8], at: usize, big: bool) -> Option<u16> {
3002        let raw: [u8; 2] = b.get(at..at + 2)?.try_into().ok()?;
3003        Some(if big {
3004            u16::from_be_bytes(raw)
3005        } else {
3006            u16::from_le_bytes(raw)
3007        })
3008    }
3009
3010    fn u32_at(b: &[u8], at: usize, big: bool) -> Option<u32> {
3011        let raw: [u8; 4] = b.get(at..at + 4)?.try_into().ok()?;
3012        Some(if big {
3013            u32::from_be_bytes(raw)
3014        } else {
3015            u32::from_le_bytes(raw)
3016        })
3017    }
3018
3019    /// Whether `head` starts a pcap or pcapng file.
3020    pub fn is_capture(head: &[u8]) -> bool {
3021        matches!(
3022            head.get(..4),
3023            Some(
3024                [0xd4, 0xc3, 0xb2, 0xa1]
3025                    | [0xa1, 0xb2, 0xc3, 0xd4]
3026                    | [0x4d, 0x3c, 0xb2, 0xa1]
3027                    | [0xa1, 0xb2, 0x3c, 0x4d]
3028                    | [0x0a, 0x0d, 0x0d, 0x0a]
3029            )
3030        )
3031    }
3032
3033    /// Where the UDP payload of a frame of `link` type is in `frame`, relative to it.
3034    pub fn udp_payload(frame: &[u8], link: u32) -> Option<std::ops::Range<usize>> {
3035        // The link layer: where the IP packet starts and which protocol it is.
3036        let (mut at, mut ethertype) = match link {
3037            // Ethernet.
3038            1 => (14, u16_at(frame, 12, true)?),
3039            // Raw IP.
3040            101 | 12 | 14 => (0, 0),
3041            228 => (0, 0x0800),
3042            229 => (0, 0x86dd),
3043            // Linux cooked captures.
3044            113 => (16, u16_at(frame, 14, true)?),
3045            276 => (20, u16_at(frame, 0, true)?),
3046            // BSD loopback: a four-byte address family in host order.
3047            0 | 108 => {
3048                let family = u32_at(frame, 0, false)?;
3049                let family = if family > 0xffff {
3050                    family.swap_bytes()
3051                } else {
3052                    family
3053                };
3054                (4, if family == 2 { 0x0800 } else { 0x86dd })
3055            }
3056            _ => return None,
3057        };
3058        // VLAN tags.
3059        while ethertype == 0x8100 || ethertype == 0x88a8 {
3060            ethertype = u16_at(frame, at + 2, true)?;
3061            at += 4;
3062        }
3063        if ethertype == 0 {
3064            ethertype = match frame.get(at)? >> 4 {
3065                4 => 0x0800,
3066                6 => 0x86dd,
3067                _ => return None,
3068            };
3069        }
3070        let (udp, ip_end) = match ethertype {
3071            0x0800 => {
3072                let ihl = usize::from(frame.get(at)? & 0x0f) * 4;
3073                let total = usize::from(u16_at(frame, at + 2, true)?);
3074                let flags = u16_at(frame, at + 6, true)?;
3075                // Fragments other than a whole datagram are not read.
3076                if ihl < 20 || *frame.get(at + 9)? != 17 || flags & 0x3fff != 0 {
3077                    return None;
3078                }
3079                (at + ihl, (at + total).min(frame.len()))
3080            }
3081            0x86dd => {
3082                if *frame.get(at + 6)? != 17 {
3083                    return None;
3084                }
3085                let payload = usize::from(u16_at(frame, at + 4, true)?);
3086                (at + 40, (at + 40 + payload).min(frame.len()))
3087            }
3088            _ => return None,
3089        };
3090        let len = usize::from(u16_at(frame, udp + 4, true)?);
3091        let start = udp + 8;
3092        let end = (udp + len).min(ip_end);
3093        (len >= 8 && start <= end).then_some(start..end)
3094    }
3095
3096    /// The UDP payloads of the capture in `file`, each with its time.
3097    pub(super) fn packets(file: &[u8], notes: &mut Vec<String>) -> Result<Vec<Chunk>, String> {
3098        let mut out = Vec::new();
3099        let mut other = 0u64;
3100        let magic = file
3101            .get(..4)
3102            .ok_or("the capture is shorter than its header")?;
3103        if magic == [0x0a, 0x0d, 0x0d, 0x0a] {
3104            pcapng(file, &mut out, &mut other)?;
3105        } else {
3106            pcap(file, &mut out, &mut other)?;
3107        }
3108        if other > 0 {
3109            notes.push(format!(
3110                "{other} {} in the capture {} not UDP, left out",
3111                if other == 1 { "packet" } else { "packets" },
3112                if other == 1 { "is" } else { "are" }
3113            ));
3114        }
3115        Ok(out)
3116    }
3117
3118    fn push(
3119        out: &mut Vec<Chunk>,
3120        base: usize,
3121        payload: std::ops::Range<usize>,
3122        time_ns: Option<i64>,
3123    ) {
3124        out.push(Chunk {
3125            source: ChunkSource::Map(base + payload.start..base + payload.end),
3126            records: None,
3127            time_ns,
3128        });
3129    }
3130
3131    fn pcap(file: &[u8], out: &mut Vec<Chunk>, other: &mut u64) -> Result<(), String> {
3132        let (big, nanos) = match file.get(..4) {
3133            Some([0xd4, 0xc3, 0xb2, 0xa1]) => (false, false),
3134            Some([0xa1, 0xb2, 0xc3, 0xd4]) => (true, false),
3135            Some([0x4d, 0x3c, 0xb2, 0xa1]) => (false, true),
3136            Some([0xa1, 0xb2, 0x3c, 0x4d]) => (true, true),
3137            _ => return Err("not a pcap or pcapng capture: its magic is not one".into()),
3138        };
3139        let link =
3140            u32_at(file, 20, big).ok_or("the capture is shorter than its header")? & 0x0fff_ffff;
3141        let mut at = 24usize;
3142        while at + 16 <= file.len() && out.len() < MAX_PACKETS {
3143            let secs = i64::from(u32_at(file, at, big).unwrap_or(0));
3144            let frac = i64::from(u32_at(file, at + 4, big).unwrap_or(0));
3145            let caplen = u32_at(file, at + 8, big).unwrap_or(0) as usize;
3146            let start = at + 16;
3147            let Some(end) = start.checked_add(caplen).filter(|e| *e <= file.len()) else {
3148                break;
3149            };
3150            let time = secs
3151                .checked_mul(1_000_000_000)
3152                .and_then(|s| s.checked_add(if nanos { frac } else { frac * 1000 }));
3153            match udp_payload(&file[start..end], link) {
3154                Some(payload) => push(out, start, payload, time),
3155                None => *other += 1,
3156            }
3157            at = end;
3158        }
3159        Ok(())
3160    }
3161
3162    fn pcapng(file: &[u8], out: &mut Vec<Chunk>, other: &mut u64) -> Result<(), String> {
3163        let mut at = 0usize;
3164        let mut big = false;
3165        // Each interface's link type and ticks per second.
3166        let mut interfaces: Vec<(u32, u64)> = Vec::new();
3167        while at + 12 <= file.len() && out.len() < MAX_PACKETS {
3168            let kind = u32_at(file, at, big).unwrap_or(0);
3169            if kind == 0x0a0d_0d0a {
3170                big = match file.get(at + 8..at + 12) {
3171                    Some([0x1a, 0x2b, 0x3c, 0x4d]) => true,
3172                    Some([0x4d, 0x3c, 0x2b, 0x1a]) => false,
3173                    _ => return Err("a pcapng section header has no byte-order magic".into()),
3174                };
3175                interfaces.clear();
3176            }
3177            let len = u32_at(file, at + 4, big).unwrap_or(0) as usize;
3178            if len < 12 || !len.is_multiple_of(4) || at + len > file.len() {
3179                break;
3180            }
3181            let body = &file[at + 8..at + len - 4];
3182            match kind {
3183                1 => {
3184                    let link = u32::from(u16_at(body, 0, big).unwrap_or(0));
3185                    let mut ticks = 1_000_000u64;
3186                    // Options: if_tsresol (9) sets the resolution.
3187                    let mut o = 8;
3188                    while o + 4 <= body.len() {
3189                        let code = u16_at(body, o, big).unwrap_or(0);
3190                        let olen = usize::from(u16_at(body, o + 2, big).unwrap_or(0));
3191                        if code == 0 {
3192                            break;
3193                        }
3194                        if code == 9 && olen >= 1 {
3195                            let r = body[o + 4];
3196                            ticks = if r & 0x80 != 0 {
3197                                1u64.checked_shl(u32::from(r & 0x7f)).unwrap_or(1_000_000)
3198                            } else {
3199                                10u64.checked_pow(u32::from(r)).unwrap_or(1_000_000)
3200                            };
3201                        }
3202                        o += 4 + olen.div_ceil(4) * 4;
3203                    }
3204                    interfaces.push((link, ticks.max(1)));
3205                }
3206                6 => {
3207                    let iface = u32_at(body, 0, big).unwrap_or(0) as usize;
3208                    let high = u64::from(u32_at(body, 4, big).unwrap_or(0));
3209                    let low = u64::from(u32_at(body, 8, big).unwrap_or(0));
3210                    let caplen = u32_at(body, 12, big).unwrap_or(0) as usize;
3211                    let Some(frame) = body.get(20..20 + caplen) else {
3212                        at += len;
3213                        continue;
3214                    };
3215                    let (link, ticks) = interfaces.get(iface).copied().unwrap_or((1, 1_000_000));
3216                    let stamp = (high << 32) | low;
3217                    let time =
3218                        i64::try_from(u128::from(stamp) * 1_000_000_000 / u128::from(ticks)).ok();
3219                    match udp_payload(frame, link) {
3220                        Some(payload) => push(out, at + 8 + 20, payload, time),
3221                        None => *other += 1,
3222                    }
3223                }
3224                3 => {
3225                    let (link, _) = interfaces.first().copied().unwrap_or((1, 1_000_000));
3226                    let frame = &body[4.min(body.len())..];
3227                    match udp_payload(frame, link) {
3228                        Some(payload) => push(out, at + 8 + 4, payload, None),
3229                        None => *other += 1,
3230                    }
3231                }
3232                _ => {}
3233            }
3234            at += len;
3235        }
3236        Ok(())
3237    }
3238}
3239
3240#[cfg(test)]
3241mod tests {
3242    use crate::fixed_records::Bytes;
3243    use crate::formats::{Opened, Spec};
3244    use polars::prelude::*;
3245    use std::sync::Arc;
3246
3247    fn open(spec: &str, bytes: Vec<u8>) -> Opened {
3248        let spec = Spec::parse(spec, None).unwrap_or_else(|e| panic!("{e}"));
3249        spec.open_rows(Arc::new(Bytes::Owned(bytes)), "f")
3250            .unwrap_or_else(|e| panic!("{e}"))
3251    }
3252
3253    fn all(opened: &Opened) -> DataFrame {
3254        opened
3255            .records
3256            .clone()
3257            .into_lazy()
3258            .unwrap()
3259            .collect()
3260            .unwrap()
3261    }
3262
3263    fn cell(df: &DataFrame, column: &str, row: usize) -> String {
3264        match df.column(column).unwrap().get(row).unwrap() {
3265            AnyValue::String(s) => s.to_string(),
3266            AnyValue::StringOwned(s) => s.to_string(),
3267            AnyValue::Categorical(..) | AnyValue::CategoricalOwned(..) => df
3268                .column(column)
3269                .unwrap()
3270                .cast(&DataType::String)
3271                .unwrap()
3272                .get(row)
3273                .unwrap()
3274                .to_string()
3275                .trim_matches('"')
3276                .to_string(),
3277            v => v.to_string(),
3278        }
3279    }
3280
3281    /// Every window of `opened` is the same rows as the whole read. A window walks
3282    /// the records; the whole read is decoded column by column, from the row table
3283    /// when there is one.
3284    fn windows_agree(opened: &Opened) {
3285        let whole = all(opened);
3286        let rows = opened.records.rows();
3287        assert_eq!(whole.height(), rows);
3288        let walked = opened.records.window(0, rows).unwrap().collect().unwrap();
3289        assert!(walked.equals_missing(&whole), "{walked}\n{whole}");
3290        for start in [
3291            0,
3292            1,
3293            rows / 2,
3294            rows.saturating_sub(3),
3295            1023,
3296            1024,
3297            1025,
3298            2049,
3299        ] {
3300            if start >= rows {
3301                continue;
3302            }
3303            let window = opened.records.window(start, 3).unwrap().collect().unwrap();
3304            let expected = whole.slice(start as i64, 3);
3305            assert!(
3306                window.equals_missing(&expected),
3307                "window at {start}:\n{window}\n{expected}"
3308            );
3309        }
3310    }
3311
3312    const ITCH: &str = r#"
3313name = "t.itch"
3314endian = "be"
3315[records]
3316framing = "length_prefixed"
3317size = "len"
3318size_adjust = 2
3319fields = [{ name = "len", type = "u2" }, { name = "kind", type = "str", size = 1 }]
3320type = "kind"
3321
3322[[variants]]
3323name = "add"
3324when = "A"
3325fields = [{ name = "ref", type = "u8" }, { name = "shares", type = "u4" }, { name = "stock", type = "str", size = 8 }, { name = "price", type = "u4", scale = 4 }]
3326
3327[[variants]]
3328name = "exec"
3329when = ["E", "C"]
3330fields = [{ name = "ref", type = "u8" }, { name = "shares", type = "u4" }]
3331"#;
3332
3333    fn itch_message(kind: u8, body: &[u8]) -> Vec<u8> {
3334        let mut out = ((body.len() + 1) as u16).to_be_bytes().to_vec();
3335        out.push(kind);
3336        out.extend(body);
3337        out
3338    }
3339
3340    fn add(r: u64, shares: u32, stock: &str, price: u32) -> Vec<u8> {
3341        let mut body = r.to_be_bytes().to_vec();
3342        body.extend(shares.to_be_bytes());
3343        let mut s = stock.as_bytes().to_vec();
3344        s.resize(8, b' ');
3345        body.extend(s);
3346        body.extend(price.to_be_bytes());
3347        itch_message(b'A', &body)
3348    }
3349
3350    fn exec(r: u64, shares: u32) -> Vec<u8> {
3351        let mut body = r.to_be_bytes().to_vec();
3352        body.extend(shares.to_be_bytes());
3353        itch_message(b'E', &body)
3354    }
3355
3356    #[test]
3357    fn length_prefixed_variants_make_one_table_with_a_type_column() {
3358        let mut bytes = Vec::new();
3359        for i in 0..3000u64 {
3360            if i % 3 == 0 {
3361                bytes.extend(add(i, 100, "AAPL", 1_234_500));
3362            } else {
3363                bytes.extend(exec(i, 7));
3364            }
3365        }
3366        // A type no variant names still has a length: its row shows the type.
3367        bytes.extend(itch_message(b'Z', &[1, 2, 3]));
3368        let opened = open(ITCH, bytes);
3369        assert!(opened.notes.is_empty(), "{:?}", opened.notes);
3370        let df = all(&opened);
3371        assert_eq!(df.height(), 3001);
3372        assert_eq!(
3373            df.get_column_names(),
3374            ["len", "kind", "type", "ref", "shares", "stock", "price"]
3375        );
3376        assert_eq!(cell(&df, "type", 0), "add");
3377        assert_eq!(cell(&df, "type", 1), "exec");
3378        assert_eq!(cell(&df, "stock", 0), "AAPL");
3379        assert_eq!(cell(&df, "price", 0), "123.4500");
3380        assert_eq!(cell(&df, "stock", 1), "null");
3381        assert_eq!(cell(&df, "shares", 1), "7");
3382        assert_eq!(cell(&df, "type", 3000), "?Z");
3383        assert_eq!(cell(&df, "ref", 3000), "null");
3384        windows_agree(&opened);
3385        // A projection decodes the one column.
3386        let one = opened
3387            .records
3388            .clone()
3389            .into_lazy()
3390            .unwrap()
3391            .select([col("ref")])
3392            .collect()
3393            .unwrap();
3394        assert_eq!(one.width(), 1);
3395        assert_eq!(cell(&one, "ref", 2999), "2999");
3396    }
3397
3398    #[test]
3399    fn one_variant_reads_alone() {
3400        let mut bytes = Vec::new();
3401        for i in 0..2500u64 {
3402            if i % 5 == 0 {
3403                bytes.extend(add(i, 100, "MSFT", 1));
3404            } else {
3405                bytes.extend(exec(i, 7));
3406            }
3407        }
3408        let spec = Spec::parse(ITCH, None)
3409            .unwrap()
3410            .with_variant("add")
3411            .unwrap();
3412        let opened = spec.open_rows(Arc::new(Bytes::Owned(bytes)), "f").unwrap();
3413        let df = all(&opened);
3414        assert_eq!(df.height(), 500);
3415        assert_eq!(
3416            df.get_column_names(),
3417            ["len", "kind", "ref", "shares", "stock", "price"]
3418        );
3419        assert_eq!(cell(&df, "ref", 499), "2495");
3420        windows_agree(&opened);
3421        assert!(
3422            Spec::parse(ITCH, None)
3423                .unwrap()
3424                .with_variant("nope")
3425                .is_err()
3426        );
3427    }
3428
3429    #[test]
3430    fn a_variant_sizes_its_record_and_an_unknown_type_stops_the_read() {
3431        let spec = r#"
3432name = "t.var"
3433[records]
3434framing = "variant"
3435type = { field = "kind", type = "u1" }
3436[[variants]]
3437name = "a"
3438when = 1
3439fields = [{ name = "x", type = "u2" }]
3440[[variants]]
3441name = "b"
3442when = 2
3443fields = [{ name = "y", type = "u4" }, { name = "z", type = "u1" }]
3444"#;
3445        let bytes = vec![1, 5, 0, 2, 9, 0, 0, 0, 3, 1, 6, 0, 7, 0xaa];
3446        let opened = open(spec, bytes);
3447        let df = all(&opened);
3448        assert_eq!(df.height(), 3);
3449        assert_eq!(cell(&df, "x", 2), "6");
3450        assert_eq!(cell(&df, "y", 1), "9");
3451        assert_eq!(cell(&df, "z", 1), "3");
3452        assert!(opened.notes[0].contains("type 7"), "{:?}", opened.notes);
3453        windows_agree(&opened);
3454    }
3455
3456    #[test]
3457    fn sync_markers_find_frames_through_garbage() {
3458        let spec = r#"
3459name = "t.sync"
3460[records]
3461framing = "sync"
3462sync = "1ACFFC1D"
3463fields = [{ name = "seq", type = "u2" }, { name = "v", type = "u1" }]
3464"#;
3465        let mut bytes = vec![0xff, 0xee];
3466        for i in 0..3u16 {
3467            bytes.extend([0x1a, 0xcf, 0xfc, 0x1d]);
3468            bytes.extend(i.to_le_bytes());
3469            bytes.push(i as u8 * 10);
3470            if i == 1 {
3471                bytes.extend([1, 2, 3]);
3472            }
3473        }
3474        let opened = open(spec, bytes);
3475        let df = all(&opened);
3476        assert_eq!(df.height(), 3);
3477        assert_eq!(cell(&df, "v", 2), "20");
3478        windows_agree(&opened);
3479        assert!(
3480            opened.notes.iter().any(|n| n.starts_with("5 bytes")),
3481            "{:?}",
3482            opened.notes
3483        );
3484    }
3485
3486    #[test]
3487    fn checksums_bits_and_groups() {
3488        let spec = r#"
3489name = "t.ck"
3490[records]
3491framing = "length_prefixed"
3492size = "len"
3493fields = [
3494  { name = "len", type = "u1" },
3495  { name = "status", type = "u2", bits = [{ name = "valid", bit = 0 }, { name = "mode", bit = 4, width = 3, enum = { 0 = "IDLE", 1 = "RUN" } }] },
3496  { name = "n", type = "u1" },
3497  { name = "levels", group = { count = "n", fields = [{ name = "px", type = "s2" }, { name = "qty", type = "u1" }] } },
3498  { name = "crc", type = "u2" },
3499]
3500checksum = { algo = "crc16-ccitt", field = "crc", from = "status" }
3501"#;
3502        let record = |status: u16, levels: &[(i16, u8)], corrupt: bool| {
3503            let mut body = status.to_le_bytes().to_vec();
3504            body.push(levels.len() as u8);
3505            for (px, q) in levels {
3506                body.extend(px.to_le_bytes());
3507                body.push(*q);
3508            }
3509            let mut crc = crate::formats::ChecksumAlgo::Crc16Ccitt.compute(&body) as u16;
3510            if corrupt {
3511                crc ^= 1;
3512            }
3513            let mut out = vec![(1 + body.len() + 2) as u8];
3514            out.extend(body);
3515            out.extend(crc.to_le_bytes());
3516            out
3517        };
3518        let mut bytes = record(0x0011, &[(-5, 1), (7, 2)], false);
3519        bytes.extend(record(0x0000, &[], true));
3520        let opened = open(spec, bytes);
3521        let df = all(&opened);
3522        assert_eq!(df.height(), 2);
3523        assert_eq!(cell(&df, "valid", 0), "true");
3524        assert_eq!(cell(&df, "mode", 0), "RUN");
3525        assert_eq!(cell(&df, "mode", 1), "IDLE");
3526        assert_eq!(cell(&df, "checksum_ok", 0), "true");
3527        assert_eq!(cell(&df, "checksum_ok", 1), "false");
3528        let levels = df.column("levels").unwrap();
3529        assert!(
3530            matches!(levels.dtype(), DataType::List(inner) if matches!(**inner, DataType::Struct(_)))
3531        );
3532        let first = levels.get(0).unwrap().to_string();
3533        assert!(first.contains("-5") && first.contains('7'), "{first}");
3534        assert_eq!(levels.list().unwrap().lst_lengths().get(1), Some(0));
3535        windows_agree(&opened);
3536    }
3537
3538    #[test]
3539    fn fortran_records_have_their_length_on_both_ends() {
3540        let spec = r#"
3541name = "t.fortran"
3542[records]
3543framing = "length_prefixed"
3544size = "n"
3545size_adjust = 4
3546length_suffix = true
3547fields = [{ name = "n", type = "u4" }, { name = "v", type = "f8", count = "n_values" }]
3548"#;
3549        assert!(Spec::parse(spec, None).is_err());
3550        let spec = r#"
3551name = "t.fortran"
3552[records]
3553framing = "length_prefixed"
3554size = "n"
3555size_adjust = 4
3556length_suffix = true
3557fields = [{ name = "n", type = "u4" }, { name = "data", type = "bytes", size = "rest" }]
3558"#;
3559        let mut bytes = Vec::new();
3560        for payload in [&b"abc"[..], b"hello"] {
3561            bytes.extend((payload.len() as u32).to_le_bytes());
3562            bytes.extend(payload);
3563            bytes.extend((payload.len() as u32).to_le_bytes());
3564        }
3565        let opened = open(spec, bytes.clone());
3566        let df = all(&opened);
3567        assert_eq!(df.height(), 2);
3568        assert_eq!(
3569            df.column("data").unwrap().binary().unwrap().get(1),
3570            Some(&b"hello"[..])
3571        );
3572        // A suffix that disagrees stops the read there.
3573        let n = bytes.len();
3574        bytes[n - 1] = 9;
3575        let opened = open(spec, bytes);
3576        assert_eq!(opened.records.rows(), 1);
3577        assert!(
3578            opened.notes[0].contains("ends with length"),
3579            "{:?}",
3580            opened.notes
3581        );
3582    }
3583
3584    #[test]
3585    fn tagged_chunks_with_text_types_and_even_alignment() {
3586        let spec = r#"
3587name = "t.riff"
3588[records]
3589framing = "length_prefixed"
3590size = "size"
3591size_adjust = 8
3592align = 2
3593fields = [{ name = "id", type = "str", size = 4 }, { name = "size", type = "u4" }]
3594type = "id"
3595[[variants]]
3596name = "fmt"
3597when = "fmt "
3598fields = [{ name = "channels", type = "u2" }]
3599[[variants]]
3600name = "data"
3601when = "data"
3602fields = [{ name = "payload", type = "bytes", size = "rest" }]
3603"#;
3604        let mut bytes = b"fmt ".to_vec();
3605        bytes.extend(2u32.to_le_bytes());
3606        bytes.extend(2u16.to_le_bytes());
3607        bytes.extend(b"data");
3608        bytes.extend(3u32.to_le_bytes());
3609        bytes.extend([1, 2, 3, 0]);
3610        bytes.extend(b"LIST");
3611        bytes.extend(0u32.to_le_bytes());
3612        let opened = open(spec, bytes);
3613        let df = all(&opened);
3614        assert_eq!(df.height(), 3, "{:?}", opened.notes);
3615        assert_eq!(cell(&df, "channels", 0), "2");
3616        assert_eq!(
3617            df.column("payload").unwrap().binary().unwrap().get(1),
3618            Some(&[1u8, 2, 3][..])
3619        );
3620        assert_eq!(cell(&df, "type", 2), "?LIST");
3621        windows_agree(&opened);
3622    }
3623
3624    #[test]
3625    fn strz_heap_strings_varints_and_deltas() {
3626        let spec = r#"
3627name = "t.mixed"
3628[header]
3629fields = [{ name = "str_off", type = "u4" }, { name = "str_size", type = "u4" }, { name = "n", type = "u4" }]
3630[sections.strings]
3631offset = "header.str_off"
3632size = "header.str_size"
3633[records]
3634count = "header.n"
3635fields = [
3636  { name = "name", type = "strz" },
3637  { name = "label", type = "u2", string_at = "strings" },
3638  { name = "ts", type = "vu", delta = true, time = "ms" },
3639  { name = "dv", type = "vs" },
3640]
3641"#;
3642        let heap = b"alpha\0beta\0";
3643        let n = 2500u32;
3644        let mut records = Vec::new();
3645        for i in 0..n {
3646            records.extend(format!("r{i}\0").as_bytes());
3647            records.extend(if i % 2 == 0 { 0u16 } else { 6u16 }.to_le_bytes());
3648            // A delta of 1000 ms, as LEB128.
3649            records.extend([0xe8, 0x07]);
3650            // Zigzag -3 is 5.
3651            records.push(5);
3652        }
3653        let off = 12 + records.len() as u32;
3654        let mut bytes = off.to_le_bytes().to_vec();
3655        bytes.extend((heap.len() as u32).to_le_bytes());
3656        bytes.extend(n.to_le_bytes());
3657        bytes.extend(&records);
3658        bytes.extend(heap);
3659        let opened = open(spec, bytes);
3660        let df = all(&opened);
3661        assert_eq!(df.height(), 2500);
3662        assert_eq!(cell(&df, "name", 7), "r7");
3663        assert_eq!(cell(&df, "label", 0), "alpha");
3664        assert_eq!(cell(&df, "label", 1), "beta");
3665        assert_eq!(cell(&df, "dv", 0), "-3");
3666        assert_eq!(
3667            df.column("ts").unwrap().dtype(),
3668            &DataType::Datetime(TimeUnit::Milliseconds, None)
3669        );
3670        let ts = df.column("ts").unwrap().cast(&DataType::Int64).unwrap();
3671        assert_eq!(ts.get(2048).unwrap(), AnyValue::Int64(2_049_000));
3672        windows_agree(&opened);
3673    }
3674
3675    #[test]
3676    fn the_byte_order_comes_from_the_magic() {
3677        let spec = r#"
3678name = "t.auto"
3679endian = "auto"
3680match = { magic = [0xd4, 0xc3, 0xb2, 0xa1] }
3681[header]
3682fields = [{ type = "pad", size = 4 }]
3683[records]
3684fields = [{ name = "v", type = "u4" }]
3685"#;
3686        let le = [vec![0xd4, 0xc3, 0xb2, 0xa1], 7u32.to_le_bytes().to_vec()].concat();
3687        let be = [vec![0xa1, 0xb2, 0xc3, 0xd4], 7u32.to_be_bytes().to_vec()].concat();
3688        let parsed = Spec::parse(spec, None).unwrap();
3689        assert!(parsed.magic_matches(&be));
3690        for bytes in [le, be] {
3691            let df = all(&open(spec, bytes));
3692            assert_eq!(cell(&df, "v", 0), "7");
3693        }
3694    }
3695
3696    #[test]
3697    fn a_footer_counts_the_records_and_checks_the_file() {
3698        let spec = r#"
3699name = "t.foot"
3700[records]
3701count = "footer.n"
3702fields = [{ name = "v", type = "u2" }]
3703[footer]
3704fields = [{ name = "n", type = "u4" }, { name = "crc", type = "u4" }]
3705checksum = { algo = "crc32", field = "crc" }
3706"#;
3707        let mut bytes: Vec<u8> = (0..5u16).flat_map(|v| v.to_le_bytes()).collect();
3708        let crc = crate::formats::ChecksumAlgo::Crc32.compute(&bytes) as u32;
3709        bytes.extend(4u32.to_le_bytes());
3710        let mut good = bytes.clone();
3711        good.extend(crc.to_le_bytes());
3712        let opened = open(spec, good);
3713        assert_eq!(opened.records.rows(), 4);
3714        // A checksum that matches says nothing.
3715        assert!(
3716            !opened.notes.iter().any(|n| n.contains("checksum")),
3717            "{:?}",
3718            opened.notes
3719        );
3720        bytes.extend((crc ^ 1).to_le_bytes());
3721        let opened = open(spec, bytes);
3722        assert!(
3723            opened.notes.iter().any(|n| n.contains("the file's is")),
3724            "{:?}",
3725            opened.notes
3726        );
3727    }
3728
3729    fn compress(codec: &str, raw: &[u8]) -> Vec<u8> {
3730        use std::io::Write;
3731        match codec {
3732            "zstd" => zstd::encode_all(raw, 3).unwrap(),
3733            "gzip" => {
3734                let mut e = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::fast());
3735                e.write_all(raw).unwrap();
3736                e.finish().unwrap()
3737            }
3738            "zlib" => {
3739                let mut e =
3740                    flate2::write::ZlibEncoder::new(Vec::new(), flate2::Compression::fast());
3741                e.write_all(raw).unwrap();
3742                e.finish().unwrap()
3743            }
3744            "deflate" => {
3745                let mut e =
3746                    flate2::write::DeflateEncoder::new(Vec::new(), flate2::Compression::fast());
3747                e.write_all(raw).unwrap();
3748                e.finish().unwrap()
3749            }
3750            "lz4" => {
3751                let mut e = lz4::EncoderBuilder::new().build(Vec::new()).unwrap();
3752                e.write_all(raw).unwrap();
3753                let (out, r) = e.finish();
3754                r.unwrap();
3755                out
3756            }
3757            "lz4_block" => lz4::block::compress(raw, None, false).unwrap(),
3758            "snappy" => snap::raw::Encoder::new().compress_vec(raw).unwrap(),
3759            "snappy_framed" => {
3760                let mut e = snap::write::FrameEncoder::new(Vec::new());
3761                e.write_all(raw).unwrap();
3762                e.into_inner().unwrap()
3763            }
3764            "brotli" => {
3765                let mut out = Vec::new();
3766                {
3767                    let mut e = brotli::CompressorWriter::new(&mut out, 4096, 5, 22);
3768                    e.write_all(raw).unwrap();
3769                }
3770                out
3771            }
3772            "bzip2" => {
3773                let mut e = bzip2::write::BzEncoder::new(Vec::new(), bzip2::Compression::fast());
3774                e.write_all(raw).unwrap();
3775                e.finish().unwrap()
3776            }
3777            "xz" => {
3778                let mut e = xz2::write::XzEncoder::new(Vec::new(), 1);
3779                e.write_all(raw).unwrap();
3780                e.finish().unwrap()
3781            }
3782            _ => raw.to_vec(),
3783        }
3784    }
3785
3786    #[test]
3787    fn a_running_sum_in_one_plain_block_reads_the_same_in_any_window() {
3788        let spec = r#"
3789name = "t.one"
3790[blocks]
3791header = [{ name = "clen", type = "u4" }]
3792size = "clen"
3793[records]
3794fields = [{ name = "v", type = "u4", delta = "block" }]
3795"#;
3796        let mut bytes = 40u32.to_le_bytes().to_vec();
3797        bytes.extend((0..10u32).flat_map(|_| 1u32.to_le_bytes()));
3798        let opened = open(spec, bytes);
3799        let df = all(&opened);
3800        assert_eq!(cell(&df, "v", 9), "10");
3801        windows_agree(&opened);
3802    }
3803
3804    #[test]
3805    fn every_codec_reads_a_block() {
3806        for codec in [
3807            "none",
3808            "gzip",
3809            "deflate",
3810            "zlib",
3811            "zstd",
3812            "lz4",
3813            "lz4_block",
3814            "snappy",
3815            "snappy_framed",
3816            "brotli",
3817            "bzip2",
3818            "xz",
3819        ] {
3820            let spec = format!(
3821                r#"
3822name = "t.blocks"
3823[blocks]
3824header = [{{ name = "clen", type = "u4" }}, {{ name = "rawlen", type = "u4" }}]
3825size = "clen"
3826compression = "{codec}"
3827uncompressed = "rawlen"
3828[records]
3829fields = [{{ name = "v", type = "u4", delta = "block" }}]
3830"#
3831            );
3832            let mut bytes = Vec::new();
3833            for block in 0..3u32 {
3834                let raw: Vec<u8> = (0..1500u32)
3835                    .flat_map(|_| (block + 1).to_le_bytes())
3836                    .collect();
3837                let packed = compress(codec, &raw);
3838                bytes.extend((packed.len() as u32).to_le_bytes());
3839                bytes.extend((raw.len() as u32).to_le_bytes());
3840                bytes.extend(packed);
3841            }
3842            let opened = open(&spec, bytes);
3843            assert!(opened.notes.is_empty(), "{codec}: {:?}", opened.notes);
3844            let df = all(&opened);
3845            assert_eq!(df.height(), 4500, "{codec}");
3846            // The running sum starts again in each block.
3847            assert_eq!(cell(&df, "v", 1499), "1500", "{codec}");
3848            assert_eq!(cell(&df, "v", 1500), "2", "{codec}");
3849            assert_eq!(cell(&df, "v", 4499), "4500", "{codec}");
3850            windows_agree(&opened);
3851        }
3852    }
3853
3854    #[test]
3855    fn a_block_index_in_the_footer_and_a_codec_per_block() {
3856        let spec = r#"
3857name = "t.indexed"
3858[blocks]
3859header = [{ name = "codec", type = "u1" }, { name = "clen", type = "u4" }]
3860size = "clen"
3861compression = { field = "codec", values = { 0 = "none", 1 = "zstd" } }
3862index = { at = "footer.index_off", count = "footer.n_blocks", fields = [{ name = "offset", type = "u8" }, { name = "rows", type = "u4" }] }
3863[records]
3864fields = [{ name = "v", type = "u2" }]
3865[footer]
3866fields = [{ name = "index_off", type = "u8" }, { name = "n_blocks", type = "u4" }]
3867"#;
3868        let mut bytes = Vec::new();
3869        let mut index = Vec::new();
3870        for block in 0..4u16 {
3871            let raw: Vec<u8> = (0..100u16)
3872                .flat_map(|i| (block * 100 + i).to_le_bytes())
3873                .collect();
3874            let (codec, body) = if block % 2 == 0 {
3875                (0u8, raw.clone())
3876            } else {
3877                (1u8, zstd::encode_all(&raw[..], 1).unwrap())
3878            };
3879            index.extend((bytes.len() as u64).to_le_bytes());
3880            index.extend(100u32.to_le_bytes());
3881            bytes.push(codec);
3882            bytes.extend((body.len() as u32).to_le_bytes());
3883            bytes.extend(body);
3884        }
3885        let index_off = bytes.len() as u64;
3886        bytes.extend(index);
3887        bytes.extend(index_off.to_le_bytes());
3888        bytes.extend(4u32.to_le_bytes());
3889        let opened = open(spec, bytes);
3890        assert!(opened.notes.is_empty(), "{:?}", opened.notes);
3891        let df = all(&opened);
3892        assert_eq!(df.height(), 400);
3893        assert_eq!(cell(&df, "v", 399), "399");
3894        windows_agree(&opened);
3895    }
3896
3897    #[test]
3898    fn columns_at_offsets_in_one_file() {
3899        let spec = r#"
3900name = "t.cols"
3901layout = "columns"
3902[header]
3903fields = [{ name = "n", type = "u4" }, { name = "a_off", type = "u4" }, { name = "b_off", type = "u4" }]
3904[records]
3905count = "header.n"
3906fields = [{ name = "a", type = "u2", offset = "header.a_off" }, { name = "b", type = "f8", offset = "header.b_off" }]
3907"#;
3908        let mut bytes = Vec::new();
3909        bytes.extend(3u32.to_le_bytes());
3910        bytes.extend(12u32.to_le_bytes());
3911        bytes.extend(18u32.to_le_bytes());
3912        for v in [1u16, 2, 3] {
3913            bytes.extend(v.to_le_bytes());
3914        }
3915        for v in [0.5f64, 1.5, 2.5] {
3916            bytes.extend(v.to_le_bytes());
3917        }
3918        let df = all(&open(spec, bytes));
3919        assert_eq!(df.height(), 3);
3920        assert_eq!(cell(&df, "a", 2), "3");
3921        assert_eq!(cell(&df, "b", 1), "1.5");
3922    }
3923
3924    #[test]
3925    fn a_ring_buffer_starts_at_its_oldest_record() {
3926        let spec = r#"
3927name = "t.ring"
3928[header]
3929fields = [{ name = "head", type = "u4" }]
3930[records]
3931ring = "header.head"
3932fields = [{ name = "v", type = "u1" }]
3933"#;
3934        let mut bytes = 2u32.to_le_bytes().to_vec();
3935        bytes.extend([30, 40, 10, 20]);
3936        let df = all(&open(spec, bytes));
3937        let values: Vec<String> = (0..4).map(|i| cell(&df, "v", i)).collect();
3938        assert_eq!(values, ["10", "20", "30", "40"]);
3939    }
3940
3941    #[test]
3942    fn half_floats_and_text_encodings() {
3943        let spec = r#"
3944name = "t.enc"
3945[records]
3946fields = [
3947  { name = "h", type = "f2" }, { name = "b", type = "bf2" },
3948  { name = "latin", type = "str", size = 4, encoding = "latin1" },
3949  { name = "wide", type = "str", size = 6, encoding = "utf16le" },
3950]
3951"#;
3952        let mut bytes = half::f16::from_f32(1.5).to_bits().to_le_bytes().to_vec();
3953        bytes.extend(half::bf16::from_f32(-2.0).to_bits().to_le_bytes());
3954        bytes.extend([b'c', 0xe9, b' ', b' ']);
3955        bytes.extend([b'h', 0, b'i', 0, 0, 0]);
3956        let df = all(&open(spec, bytes));
3957        assert_eq!(cell(&df, "h", 0), "1.5");
3958        assert_eq!(cell(&df, "b", 0), "-2.0");
3959        assert_eq!(cell(&df, "latin", 0), "c\u{e9}");
3960        assert_eq!(cell(&df, "wide", 0), "hi");
3961    }
3962
3963    /// A UDP packet in Ethernet and IPv4, carrying `payload`.
3964    fn udp_frame(payload: &[u8]) -> Vec<u8> {
3965        let mut f = vec![0u8; 12];
3966        f.extend([0x08, 0x00]);
3967        let total = (20 + 8 + payload.len()) as u16;
3968        f.extend([0x45, 0]);
3969        f.extend(total.to_be_bytes());
3970        f.extend([0, 0, 0x40, 0, 64, 17, 0, 0]);
3971        f.extend([10, 0, 0, 1, 10, 0, 0, 2]);
3972        f.extend([0x30, 0x39, 0x30, 0x39]);
3973        f.extend(((8 + payload.len()) as u16).to_be_bytes());
3974        f.extend([0, 0]);
3975        f.extend(payload);
3976        f
3977    }
3978
3979    #[test]
3980    fn a_capture_s_udp_payloads_hold_the_records() {
3981        let spec = r#"
3982name = "t.mold"
3983endian = "be"
3984[capture]
3985header = [{ name = "session", type = "str", size = 10 }, { name = "seq", type = "u8" }, { name = "count", type = "u2" }]
3986count = "count"
3987time = "captured"
3988[records]
3989framing = "length_prefixed"
3990size = "len"
3991size_adjust = 2
3992fields = [{ name = "len", type = "u2" }, { name = "msg", type = "str", size = "rest" }]
3993"#;
3994        let mut pcap = vec![0xd4, 0xc3, 0xb2, 0xa1, 2, 0, 4, 0];
3995        pcap.extend([0u8; 8]);
3996        pcap.extend(65535u32.to_le_bytes());
3997        pcap.extend(1u32.to_le_bytes());
3998        for (i, msgs) in [vec!["hi", "there"], vec!["x"]].into_iter().enumerate() {
3999            let mut payload = b"SESSION001".to_vec();
4000            payload.extend((i as u64).to_be_bytes());
4001            payload.extend((msgs.len() as u16).to_be_bytes());
4002            for m in &msgs {
4003                payload.extend((m.len() as u16).to_be_bytes());
4004                payload.extend(m.as_bytes());
4005            }
4006            let frame = udp_frame(&payload);
4007            pcap.extend((1_700_000_000 + i as u32).to_le_bytes());
4008            pcap.extend(5u32.to_le_bytes());
4009            pcap.extend((frame.len() as u32).to_le_bytes());
4010            pcap.extend((frame.len() as u32).to_le_bytes());
4011            pcap.extend(frame);
4012        }
4013        let opened = open(spec, pcap);
4014        let df = all(&opened);
4015        assert_eq!(df.height(), 3, "{:?}", opened.notes);
4016        assert_eq!(cell(&df, "msg", 1), "there");
4017        assert_eq!(cell(&df, "msg", 2), "x");
4018        assert_eq!(cell(&df, "captured", 2), "2023-11-14 22:13:21.000005");
4019    }
4020
4021    #[test]
4022    fn a_variant_s_fields_past_its_record_s_end_are_null() {
4023        // The second add says it is 15 bytes: its stock and price are past its end.
4024        let mut short = add(2, 2, "B", 2);
4025        short[..2].copy_from_slice(&13u16.to_be_bytes());
4026        short.truncate(15);
4027        let bytes = [add(1, 1, "A", 1), short, exec(3, 3), add(4, 4, "D", 4)].concat();
4028        let opened = open(ITCH, bytes);
4029        let df = all(&opened);
4030        assert_eq!(df.height(), 4, "{:?}", opened.notes);
4031        assert_eq!(cell(&df, "shares", 1), "2");
4032        assert_eq!(cell(&df, "stock", 1), "null");
4033        assert_eq!(cell(&df, "stock", 3), "D");
4034        windows_agree(&opened);
4035    }
4036
4037    /// A file's walk is kept: opened again with the same spec it is not walked again,
4038    /// and with another variant it is walked and kept in its place.
4039    #[test]
4040    fn a_files_walk_is_kept_for_its_next_open() {
4041        let dir = tempfile::tempdir().unwrap();
4042        let path = dir.path().join("kept.itch");
4043        let bytes: Vec<u8> = (0..3000u64)
4044            .flat_map(|i| [add(i, i as u32, "S", 1), exec(i, 2)].concat())
4045            .collect();
4046        std::fs::write(&path, bytes).unwrap();
4047        let spec = Spec::parse(ITCH, None).unwrap();
4048        let first = spec.open(&path, "kept.itch").unwrap();
4049        let kept = crate::indexed::peek::<super::KeptWalk>(&path).expect("kept");
4050        assert_eq!(kept.rows, 6000);
4051        let again = spec.open(&path, "kept.itch").unwrap();
4052        let same = crate::indexed::peek::<super::KeptWalk>(&path).unwrap();
4053        assert!(Arc::ptr_eq(&kept, &same), "not walked again");
4054        assert!(all(&first).equals_missing(&all(&again)));
4055        windows_agree(&again);
4056        let exec_only = spec
4057            .with_variant("exec")
4058            .unwrap()
4059            .open(&path, "kept.itch")
4060            .unwrap();
4061        assert_eq!(exec_only.records.rows(), 3000);
4062        let replaced = crate::indexed::peek::<super::KeptWalk>(&path).unwrap();
4063        assert_eq!(replaced.rows, 3000);
4064        windows_agree(&exec_only);
4065    }
4066
4067    #[test]
4068    fn a_record_cut_short_is_left_out_and_said() {
4069        let opened = open(
4070            ITCH,
4071            [add(1, 1, "A", 1), add(2, 2, "B", 2)[..10].to_vec()].concat(),
4072        );
4073        assert_eq!(opened.records.rows(), 1);
4074        assert!(
4075            opened.notes[0].contains("not a whole record"),
4076            "{:?}",
4077            opened.notes
4078        );
4079    }
4080}