Skip to main content

datui_lib/formats/
framed_records.rs

1//! Records not all one size, read from a memory map: length-prefixed records, variants
2//! picked by a type field, records found by a sync marker, records in compressed blocks
3//! or packet-capture payloads, and fixed records with fields read one at a time
4//! (varints, NUL-terminated text, counted groups, bit fields, deltas, checksums).
5//!
6//! A first pass keeps where every 1024th record starts (and delta columns' running sums
7//! there), so a window anywhere walks from the nearest checkpoint. Fixed-size records
8//! need no pass. Compressed blocks are decompressed on first read, a few kept. The frame
9//! decodes over a row index ([`crate::formats::row_index`]), as [`crate::formats::fixed_records`] does;
10//! deeper windows go through [`FramedRecords::window`].
11
12use crate::formats::fixed_records::{Bytes, ColumnLayout, Physical};
13use crate::formats::{
14    self, Amount, Codec, Compression, Delta, Encoding, Field, Framing, HeaderValues, Spec, Type,
15};
16use polars::prelude::*;
17use std::collections::VecDeque;
18use std::ops::Range;
19use std::sync::{Arc, Mutex};
20
21/// Records between checkpoints of the index.
22const CHECKPOINT: u64 = 1024;
23
24/// The most bytes one block may decompress to.
25pub const MAX_BLOCK: usize = 256 << 20;
26
27/// Decompressed blocks kept for reading again.
28const CACHED_BLOCKS: usize = 8;
29
30/// Items one group may hold, and values one counted field may hold.
31const MAX_ITEMS: u64 = 1 << 20;
32
33/// The most rows a table holds: Polars counts rows in 32 bits.
34const MAX_ROWS: usize = IdxSize::MAX as usize;
35
36/// A size known now, or read from an earlier field of the same record.
37#[derive(Debug, Clone, Copy)]
38enum SizeRef {
39    Given(usize),
40    Slot { slot: usize, adjust: i64 },
41    Rest,
42}
43
44/// How an integer's raw value is read, for fields others refer to.
45#[derive(Debug, Clone, Copy)]
46struct IntRead {
47    signed: bool,
48    big: bool,
49}
50
51/// One field of a record (or a block header, a payload header, a group item).
52#[derive(Debug, Clone)]
53struct FieldPlan {
54    name: String,
55    slot: usize,
56    kind: Kind,
57    /// Values per cell: `None` for one; a size for an Array (given) or a List.
58    count: Option<SizeRef>,
59    /// The columns its values go to: one, or one per value when flattened.
60    outs: Vec<usize>,
61    /// Bit fields: column, lowest bit, width.
62    bits: Vec<(usize, u32, u32)>,
63    delta: Delta,
64    /// The index of its running sum, for a delta.
65    delta_index: usize,
66    /// The stored value that means null, for a delta (checked before summing).
67    sentinel: Option<i128>,
68}
69
70#[derive(Debug, Clone)]
71enum Kind {
72    /// A fixed-width value: an integer, a float, a bool, or text and bytes of a size
73    /// known before the record is read.
74    Fixed {
75        width: usize,
76        int: Option<IntRead>,
77    },
78    /// Text or bytes whose size is read from the record. `encoding` is `None` for bytes.
79    Sized {
80        size: SizeRef,
81        encoding: Option<Encoding>,
82    },
83    /// Text to its NUL, within `max` bytes when given (which it then always takes).
84    Strz {
85        max: Option<SizeRef>,
86        encoding: Encoding,
87    },
88    Var {
89        signed: bool,
90    },
91    Pad {
92        size: SizeRef,
93    },
94    /// An offset into a section, where NUL-terminated text is.
95    StringAt {
96        width: usize,
97        big: bool,
98        section: Range<usize>,
99    },
100    Group {
101        items: Vec<FieldPlan>,
102        slots: usize,
103    },
104}
105
106/// One variant of the records.
107#[derive(Debug, Clone)]
108struct VariantPlan {
109    name: Arc<str>,
110    ints: Vec<i128>,
111    texts: Vec<String>,
112    fields: Vec<FieldPlan>,
113    size: Option<SizeRef>,
114}
115
116/// What a column collects while records are read.
117#[derive(Debug, Clone)]
118enum Sink {
119    /// Fixed-width cells, decoded together by the reader of fixed records.
120    Packed {
121        layout: ColumnLayout,
122        cell: usize,
123        buf: Vec<u8>,
124        valid: Vec<bool>,
125    },
126    Text(Vec<Option<String>>),
127    Binary(Vec<Option<Vec<u8>>>),
128    Bits {
129        width: u32,
130        labels: Option<Arc<std::collections::BTreeMap<i64, String>>>,
131        values: Vec<Option<u64>>,
132    },
133    Flag(Vec<Option<bool>>),
134    Label(Vec<Option<Arc<str>>>),
135    Time(Vec<Option<i64>>),
136    List {
137        inner: Box<Sink>,
138        offsets: Vec<i64>,
139        valid: Vec<bool>,
140    },
141    Struct {
142        names: Vec<PlSmallStr>,
143        fields: Vec<Sink>,
144        len: usize,
145    },
146    /// A column the read does not want.
147    Skip,
148}
149
150impl Sink {
151    fn push_null(&mut self) {
152        match self {
153            Self::Packed {
154                cell, buf, valid, ..
155            } => {
156                buf.resize(buf.len() + *cell, 0);
157                valid.push(false);
158            }
159            Self::Text(v) => v.push(None),
160            Self::Label(v) => v.push(None),
161            Self::Binary(v) => v.push(None),
162            Self::Bits { values, .. } => values.push(None),
163            Self::Flag(v) => v.push(None),
164            Self::Time(v) => v.push(None),
165            Self::List { offsets, valid, .. } => {
166                offsets.push(*offsets.last().unwrap_or(&0));
167                valid.push(false);
168            }
169            Self::Struct { fields, len, .. } => {
170                for f in fields {
171                    f.push_null();
172                }
173                *len += 1;
174            }
175            Self::Skip => {}
176        }
177    }
178
179    fn push_bytes(&mut self, bytes: &[u8]) {
180        if let Self::Packed { buf, valid, .. } = self {
181            buf.extend_from_slice(bytes);
182            valid.push(true);
183        }
184    }
185
186    fn finish(self, name: PlSmallStr) -> PolarsResult<Series> {
187        Ok(match self {
188            Self::Packed {
189                layout,
190                cell,
191                buf,
192                valid,
193            } => {
194                let rows = valid.len();
195                // The cells sit side by side in `buf`, whatever the record's stride was:
196                // a delta's running sum is wider than the value stored.
197                let layout = ColumnLayout {
198                    start: 0,
199                    stride: cell,
200                    ..layout
201                };
202                let column = crate::formats::fixed_records::decode(&buf, &layout, rows)?;
203                let series = column
204                    .as_materialized_series()
205                    .clone()
206                    .with_name(name.clone());
207                if valid.iter().all(|v| *v) {
208                    series
209                } else {
210                    let mask: BooleanChunked = valid.into_iter().collect();
211                    let nulls = Series::full_null(PlSmallStr::EMPTY, rows, series.dtype());
212                    series.zip_with(&mask, &nulls)?.with_name(name)
213                }
214            }
215            Self::Text(v) => StringChunked::from_iter_options(name, v.into_iter()).into_series(),
216            Self::Label(v) => {
217                // Few labels, many rows: each label is cast once and the rows gathered,
218                // where casting every row's text took most of a variant column's time.
219                let mut distinct: Vec<Arc<str>> = Vec::new();
220                let mut by_text = std::collections::HashMap::new();
221                let codes: IdxCa = v
222                    .iter()
223                    .map(|label| {
224                        let label = label.as_ref()?;
225                        // A variant's name is one shared string.
226                        if let Some(i) =
227                            distinct.iter().take(16).position(|d| Arc::ptr_eq(d, label))
228                        {
229                            return Some(i as IdxSize);
230                        }
231                        Some(*by_text.entry(label.clone()).or_insert_with(|| {
232                            distinct.push(label.clone());
233                            (distinct.len() - 1) as IdxSize
234                        }))
235                    })
236                    .collect();
237                StringChunked::from_iter_values(name, distinct.iter().map(|d| &**d))
238                    .into_series()
239                    .cast(&DataType::from_categories(Categories::global()))?
240                    .take(&codes)?
241            }
242            Self::Binary(v) => BinaryChunked::from_iter_options(name, v.into_iter()).into_series(),
243            Self::Bits {
244                width,
245                labels,
246                values,
247            } => match (labels, width) {
248                (Some(labels), _) => StringChunked::from_iter_options(
249                    name,
250                    values.into_iter().map(|v| {
251                        v.map(|v| {
252                            i64::try_from(v)
253                                .ok()
254                                .and_then(|k| labels.get(&k).cloned())
255                                .unwrap_or_else(|| v.to_string())
256                        })
257                    }),
258                )
259                .into_series(),
260                (None, 1) => BooleanChunked::from_iter_options(
261                    name,
262                    values.into_iter().map(|v| v.map(|v| v != 0)),
263                )
264                .into_series(),
265                (None, w) => {
266                    let wide =
267                        UInt64Chunked::from_iter_options(name, values.into_iter()).into_series();
268                    let dtype = match w {
269                        2..=8 => DataType::UInt8,
270                        9..=16 => DataType::UInt16,
271                        17..=32 => DataType::UInt32,
272                        _ => DataType::UInt64,
273                    };
274                    wide.strict_cast(&dtype)?
275                }
276            },
277            Self::Flag(v) => BooleanChunked::from_iter_options(name, v.into_iter()).into_series(),
278            Self::Time(v) => Int64Chunked::from_iter_options(name, v.into_iter())
279                .into_datetime(TimeUnit::Nanoseconds, None)
280                .into_series(),
281            Self::List {
282                inner,
283                offsets,
284                valid,
285            } => {
286                let values = inner.finish(PlSmallStr::from_static("item"))?.rechunk();
287                list_series(name, values, offsets, valid)?
288            }
289            Self::Struct { names, fields, len } => {
290                let series = names
291                    .into_iter()
292                    .zip(fields)
293                    .map(|(n, f)| f.finish(n))
294                    .collect::<PolarsResult<Vec<_>>>()?;
295                StructChunked::from_series(name, len, series.iter())?.into_series()
296            }
297            Self::Skip => Series::new_empty(name, &DataType::Null),
298        })
299    }
300}
301
302/// A List column of `values`, row `i` holding `offsets[i]..offsets[i + 1]`.
303fn list_series(
304    name: PlSmallStr,
305    values: Series,
306    offsets: Vec<i64>,
307    valid: Vec<bool>,
308) -> PolarsResult<Series> {
309    use polars_arrow::array::ListArray;
310    use polars_arrow::bitmap::Bitmap;
311    use polars_arrow::offset::OffsetsBuffer;
312    let inner_dtype = values.dtype().clone();
313    let array = values.to_arrow(0, CompatLevel::newest());
314    let mut all = Vec::with_capacity(offsets.len() + 1);
315    all.push(0i64);
316    all.extend(offsets);
317    let offsets = OffsetsBuffer::<i64>::try_from(all)?;
318    let validity = (!valid.iter().all(|v| *v)).then(|| Bitmap::from_iter(valid));
319    let dtype = ListArray::<i64>::default_datatype(array.dtype().clone());
320    let list = ListArray::<i64>::try_new(dtype, offsets, array, validity)?;
321    Series::from_arrow(name, Box::new(list))?.cast(&DataType::List(Box::new(inner_dtype)))
322}
323
324/// A column of the table and what collects it.
325#[derive(Debug, Clone)]
326struct OutColumn {
327    name: PlSmallStr,
328    dtype: DataType,
329    proto: Sink,
330}
331
332/// Where a run of records is.
333#[derive(Debug, Clone)]
334enum ChunkSource {
335    /// Bytes of the file.
336    Map(Range<usize>),
337    /// A block body of the file, compressed.
338    Block {
339        body: Range<usize>,
340        codec: Compression,
341        uncompressed: Option<usize>,
342    },
343}
344
345#[derive(Debug, Clone)]
346struct Chunk {
347    source: ChunkSource,
348    /// Records in it, when its header (or the index) says.
349    records: Option<u64>,
350    /// A packet's capture time.
351    time_ns: Option<i64>,
352}
353
354/// Where the reading of rows can start: the `row`th record, at `pos` in `chunk`.
355#[derive(Debug, Clone)]
356struct Checkpoint {
357    row: u64,
358    chunk: u32,
359    pos: u32,
360    /// Records of the chunk before `pos`, for a chunk that counts its records.
361    taken: u32,
362    acc: Box<[i128]>,
363}
364
365/// Where the records are.
366#[derive(Debug)]
367enum Index {
368    /// Records all `size` bytes from `start`; row `i` is record `(ring + i) % rows`.
369    Stride {
370        start: usize,
371        size: usize,
372        ring: usize,
373    },
374    Walk(Vec<Checkpoint>),
375}
376
377/// The variant tag of a row read field by field: one cut short, or of no variant.
378const WALK: u8 = u8::MAX;
379
380/// Each row's record start and variant, kept from the opening walk so every column of a
381/// query reads the same starts without walking again. One map run only: blocks' and
382/// packets' records live in their decompressed copies and payloads.
383#[derive(Debug)]
384struct RowTable {
385    /// Where each row's record starts, its sync marker included.
386    starts: crate::formats::indexed::Offsets,
387    /// Each row's variant (0 without variants), or [`WALK`].
388    tags: Vec<u8>,
389}
390
391/// A file's walk, kept for its next open with the same spec (variants, `H`), so it is
392/// walked once. Kept by [`crate::formats::indexed::keep`], bounded by file sizes.
393struct KeptWalk {
394    spec: Spec,
395    data: Range<usize>,
396    index: Arc<Index>,
397    table: Option<Arc<RowTable>>,
398    rows: usize,
399    notes: Vec<String>,
400}
401
402/// Where one column's value is in a record of one variant.
403#[derive(Debug, Clone)]
404enum Source {
405    /// The variant has no such field.
406    Null,
407    /// At a fixed place from the record's start, so read without a walk.
408    At { offset: usize, field: FieldPlan },
409    /// The variant's name.
410    Label,
411    /// Behind a field whose size the record says: the record is walked.
412    Walk,
413    /// A running sum, which needs every record before it: walked from a checkpoint.
414    Summed,
415}
416
417/// The compiled spec: how to walk one record.
418#[derive(Debug)]
419struct Plan {
420    framing: Framing,
421    common: Vec<FieldPlan>,
422    /// The slot of the type field and whether it is text.
423    type_slot: Option<(usize, bool)>,
424    variants: Vec<VariantPlan>,
425    type_out: Option<usize>,
426    /// The record's size, when the spec gives one or a field holds it.
427    size: Option<SizeRef>,
428    /// The width and byte order of a length repeated after the record.
429    suffix: Option<(usize, bool)>,
430    align: usize,
431    sync: Vec<u8>,
432    checksum: Option<ChecksumPlan>,
433    /// Payload header fields, and the slot counting its records.
434    chunk_header: Vec<FieldPlan>,
435    chunk_count: Option<usize>,
436    chunk_slots: usize,
437    time_out: Option<usize>,
438    slots: usize,
439    deltas: Vec<Delta>,
440    columns: Vec<OutColumn>,
441    /// The variant read alone.
442    only: Option<usize>,
443    /// Per column, per variant (one entry without variants): where its value is.
444    sources: Vec<Vec<Source>>,
445}
446
447#[derive(Debug, Clone)]
448struct ChecksumPlan {
449    algo: formats::ChecksumAlgo,
450    field: usize,
451    from: Option<usize>,
452    to: usize,
453    out: usize,
454}
455
456/// A record walk's per-record state: where each field started and the integers read.
457struct Frame {
458    starts: Vec<Option<usize>>,
459    ends: Vec<Option<usize>>,
460    ints: Vec<Option<i128>>,
461}
462
463impl Frame {
464    fn new(slots: usize) -> Self {
465        Self {
466            starts: vec![None; slots],
467            ends: vec![None; slots],
468            ints: vec![None; slots],
469        }
470    }
471
472    /// Forget slots `range`.
473    fn clear(&mut self, range: Range<usize>) {
474        for s in range {
475            self.starts[s] = None;
476            self.ends[s] = None;
477            self.ints[s] = None;
478        }
479    }
480}
481
482/// Why a record could not be read.
483enum Stop {
484    /// The bytes ran out before the record did.
485    Truncated,
486    /// Anything else: said as it is.
487    Said(String),
488}
489
490/// Records read through a spec, by walking them.
491pub struct FramedRecords {
492    bytes: Arc<Bytes>,
493    plan: Arc<Plan>,
494    chunks: Vec<Chunk>,
495    index: Arc<Index>,
496    /// Built by the walk at open, when the records are one run of the map.
497    table: Option<Arc<RowTable>>,
498    rows: usize,
499    schema: SchemaRef,
500    cache: Mutex<VecDeque<(usize, Arc<Vec<u8>>)>>,
501}
502
503impl std::fmt::Debug for FramedRecords {
504    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
505        f.debug_struct("FramedRecords")
506            .field("rows", &self.rows)
507            .field("chunks", &self.chunks.len())
508            .finish()
509    }
510}
511
512/// Whether `spec` needs records walked rather than the fixed reader.
513pub fn needed(spec: &Spec) -> bool {
514    fn walked(field: &Field) -> bool {
515        !field.is_fixed_width()
516            || field.size.as_ref().is_some_and(|s| !s.is_fixed())
517            || field.count.as_ref().is_some_and(|s| !s.is_fixed())
518            || field.delta != Delta::None
519            || !field.bits.is_empty()
520            || field.string_at.is_some()
521    }
522    let r = &spec.records;
523    r.framing != Framing::Fixed
524        || !r.variants.is_empty()
525        || r.checksum.is_some()
526        || r.ring.is_some()
527        || spec.blocks.is_some()
528        || spec.capture.is_some()
529        || formats::all_fields(r).any(walked)
530}
531
532/// What compiling the spec builds: columns, slots and running sums.
533struct Compiler<'a> {
534    spec: &'a Spec,
535    header: &'a HeaderValues,
536    data: &'a [u8],
537    slots: usize,
538    deltas: Vec<Delta>,
539}
540
541/// Names of a scope's fields and their slots.
542#[derive(Default, Clone)]
543struct Scope(Vec<(String, usize)>);
544
545impl Scope {
546    fn slot(&self, name: &str) -> Option<usize> {
547        self.0
548            .iter()
549            .rev()
550            .find(|(n, _)| n == name)
551            .map(|(_, s)| *s)
552    }
553}
554
555impl Compiler<'_> {
556    fn size(&self, amount: &Amount, scope: &Scope, what: &str) -> Result<SizeRef, String> {
557        Ok(match amount {
558            Amount::Record { field, adjust } => SizeRef::Slot {
559                slot: scope
560                    .slot(field)
561                    .ok_or_else(|| format!("{what}: `{field}` is not an earlier field"))?,
562                adjust: *adjust,
563            },
564            Amount::Rest => SizeRef::Rest,
565            other => SizeRef::Given(
566                usize::try_from(self.header.resolve(other, what)?)
567                    .map_err(|_| format!("{what}: too large"))?,
568            ),
569        })
570    }
571
572    /// The fields' plans; their columns are added to `columns`.
573    fn fields(
574        &mut self,
575        fields: &[Field],
576        scope: &mut Scope,
577        columns: &mut Vec<OutColumn>,
578        shared: bool,
579    ) -> Result<Vec<FieldPlan>, String> {
580        let mut out = Vec::new();
581        for field in fields {
582            out.push(self.field(field, scope, columns, shared)?);
583        }
584        Ok(out)
585    }
586
587    /// The column named `name`, new or (for fields variants share) the one there is.
588    fn column(
589        columns: &mut Vec<OutColumn>,
590        name: &str,
591        dtype: DataType,
592        proto: Sink,
593        shared: bool,
594    ) -> usize {
595        if shared && let Some(i) = columns.iter().position(|c| c.name == name) {
596            return i;
597        }
598        columns.push(OutColumn {
599            name: name.into(),
600            dtype,
601            proto,
602        });
603        columns.len() - 1
604    }
605
606    fn field(
607        &mut self,
608        field: &Field,
609        scope: &mut Scope,
610        columns: &mut Vec<OutColumn>,
611        shared: bool,
612    ) -> Result<FieldPlan, String> {
613        let name = field.name.clone().unwrap_or_default();
614        let slot = self.slots;
615        self.slots += 1;
616        let big = field.endian.unwrap_or(self.spec.endian) == formats::Endian::Big;
617        let count = field
618            .count
619            .as_ref()
620            .map(|c| self.size(c, scope, "count"))
621            .transpose()?;
622        let given_count = match count {
623            Some(SizeRef::Given(n)) => {
624                if n as u64 > MAX_ITEMS {
625                    return Err(format!(
626                        "field `{name}`: {n} values is more than {MAX_ITEMS}"
627                    ));
628                }
629                Some(n)
630            }
631            _ => None,
632        };
633        let mut plan = FieldPlan {
634            name: name.clone(),
635            slot,
636            kind: Kind::Pad {
637                size: SizeRef::Given(0),
638            },
639            count,
640            outs: Vec::new(),
641            bits: Vec::new(),
642            delta: field.delta,
643            delta_index: 0,
644            sentinel: None,
645        };
646        let layout = |width: usize, cells: usize| -> Result<ColumnLayout, String> {
647            let mut layout = formats::layout_of(
648                self.spec,
649                field,
650                &name,
651                formats::Place {
652                    start: 0,
653                    stride: None,
654                    width,
655                    count: cells,
656                },
657                self.header,
658            )?;
659            layout.big_endian = big;
660            Ok(layout)
661        };
662        let packed = |layout: ColumnLayout| {
663            let cell = layout.width * layout.count;
664            Sink::Packed {
665                layout,
666                cell,
667                buf: Vec::new(),
668                valid: Vec::new(),
669            }
670        };
671        // A fixed value, possibly counted: an Array, flattened columns, or a List.
672        let mut fixed_out = |this: &mut Self,
673                             plan: &mut FieldPlan,
674                             mut layout: ColumnLayout|
675         -> Result<(), String> {
676            let _ = this;
677            match (count, given_count) {
678                (None, _) => {
679                    let dtype = layout.dtype();
680                    plan.outs = vec![Self::column(columns, &name, dtype, packed(layout), shared)];
681                }
682                (Some(_), Some(n)) if field.flatten => {
683                    layout.count = 1;
684                    for i in 0..n {
685                        let n_name = format!("{name}_{i}");
686                        let mut l = layout.clone();
687                        l.name = n_name.as_str().into();
688                        let dtype = l.dtype();
689                        plan.outs
690                            .push(Self::column(columns, &n_name, dtype, packed(l), shared));
691                    }
692                }
693                (Some(_), Some(n)) => {
694                    layout.count = n.max(1);
695                    let dtype = layout.dtype();
696                    plan.outs = vec![Self::column(columns, &name, dtype, packed(layout), shared)];
697                }
698                (Some(_), None) => {
699                    layout.count = 1;
700                    let item = layout.dtype();
701                    let sink = Sink::List {
702                        inner: Box::new(packed(layout)),
703                        offsets: Vec::new(),
704                        valid: Vec::new(),
705                    };
706                    plan.outs = vec![Self::column(
707                        columns,
708                        &name,
709                        DataType::List(Box::new(item)),
710                        sink,
711                        shared,
712                    )];
713                }
714            }
715            Ok(())
716        };
717        match field.ty {
718            Type::Pad => {
719                let size = field
720                    .size
721                    .as_ref()
722                    .map_or(Ok(SizeRef::Given(0)), |s| self.size(s, scope, "size"))?;
723                plan.kind = Kind::Pad { size };
724            }
725            Type::Group => {
726                let mut inner_scope = Scope::default();
727                let mut inner_columns = Vec::new();
728                let first = self.slots;
729                let items =
730                    self.fields(&field.group, &mut inner_scope, &mut inner_columns, false)?;
731                let slots = self.slots - first;
732                let names: Vec<PlSmallStr> = inner_columns.iter().map(|c| c.name.clone()).collect();
733                let struct_dtype = DataType::Struct(
734                    inner_columns
735                        .iter()
736                        .map(|c| polars::prelude::Field::new(c.name.clone(), c.dtype.clone()))
737                        .collect(),
738                );
739                let sink = Sink::List {
740                    inner: Box::new(Sink::Struct {
741                        names,
742                        fields: inner_columns.into_iter().map(|c| c.proto).collect(),
743                        len: 0,
744                    }),
745                    offsets: Vec::new(),
746                    valid: Vec::new(),
747                };
748                plan.outs = vec![Self::column(
749                    columns,
750                    &name,
751                    DataType::List(Box::new(struct_dtype)),
752                    sink,
753                    shared,
754                )];
755                plan.kind = Kind::Group { items, slots };
756            }
757            Type::VarU | Type::VarS => {
758                let signed = field.ty == Type::VarS;
759                let mut l = layout(8, 1)?;
760                if field.delta != Delta::None {
761                    l.physical = Physical::Signed(8);
762                    plan.sentinel = sentinel(field, 64, signed);
763                    l.null = None;
764                }
765                plan.kind = Kind::Var { signed };
766                fixed_out(self, &mut plan, l)?;
767            }
768            Type::Strz => {
769                let max = field
770                    .size
771                    .as_ref()
772                    .map(|s| self.size(s, scope, "size"))
773                    .transpose()?;
774                plan.kind = Kind::Strz {
775                    max,
776                    encoding: field.encoding,
777                };
778                plan.outs = vec![Self::column(
779                    columns,
780                    &name,
781                    DataType::String,
782                    Sink::Text(Vec::new()),
783                    shared,
784                )];
785            }
786            Type::Str | Type::Bytes if !field.size.as_ref().is_some_and(Amount::is_fixed) => {
787                let size = self.size(
788                    field.size.as_ref().expect("checked at parse"),
789                    scope,
790                    "size",
791                )?;
792                let text = field.ty == Type::Str;
793                plan.kind = Kind::Sized {
794                    size,
795                    encoding: text.then_some(field.encoding),
796                };
797                let (dtype, sink) = if text {
798                    (DataType::String, Sink::Text(Vec::new()))
799                } else {
800                    (DataType::Binary, Sink::Binary(Vec::new()))
801                };
802                plan.outs = vec![Self::column(columns, &name, dtype, sink, shared)];
803            }
804            _ => {
805                let width = match (field.ty.width(), &field.size) {
806                    (Some(w), _) => w as usize,
807                    (None, Some(amount)) => usize::try_from(self.header.resolve(amount, "size")?)
808                        .map_err(|_| "size: too large".to_string())?,
809                    (None, None) => 0,
810                };
811                if width == 0 {
812                    return Err(format!("field `{name}` takes no bytes"));
813                }
814                let int = match field.ty {
815                    Type::Unsigned(_) => Some(IntRead { signed: false, big }),
816                    Type::Signed(_) => Some(IntRead { signed: true, big }),
817                    _ => None,
818                };
819                if let Some(section) = &field.string_at {
820                    let section = self.section(section)?;
821                    plan.kind = Kind::StringAt {
822                        width,
823                        big,
824                        section,
825                    };
826                    plan.outs = vec![Self::column(
827                        columns,
828                        &name,
829                        DataType::String,
830                        Sink::Text(Vec::new()),
831                        shared,
832                    )];
833                } else {
834                    let mut l = layout(width, 1)?;
835                    if field.delta != Delta::None {
836                        l.physical = Physical::Signed(8);
837                        l.width = 8;
838                        plan.sentinel =
839                            sentinel(field, width as u32 * 8, int.is_some_and(|i| i.signed));
840                        l.null = None;
841                    }
842                    plan.kind = Kind::Fixed { width, int };
843                    fixed_out(self, &mut plan, l)?;
844                }
845            }
846        }
847        if field.delta != Delta::None {
848            plan.delta_index = self.deltas.len();
849            self.deltas.push(field.delta);
850        }
851        for bit in &field.bits {
852            let sink = Sink::Bits {
853                width: bit.width,
854                labels: bit.labels.clone(),
855                values: Vec::new(),
856            };
857            let dtype = match (&bit.labels, bit.width) {
858                (Some(_), _) => DataType::String,
859                (None, 1) => DataType::Boolean,
860                (None, 2..=8) => DataType::UInt8,
861                (None, 9..=16) => DataType::UInt16,
862                (None, 17..=32) => DataType::UInt32,
863                _ => DataType::UInt64,
864            };
865            let col = Self::column(columns, &bit.name, dtype, sink, shared);
866            plan.bits.push((col, bit.bit, bit.width));
867        }
868        if field.name.is_some() {
869            scope.0.push((name, slot));
870        }
871        Ok(plan)
872    }
873
874    /// The bytes of section `name`, bounded by the file.
875    fn section(&self, name: &str) -> Result<Range<usize>, String> {
876        let section = self
877            .spec
878            .sections
879            .iter()
880            .find(|s| s.name == name)
881            .ok_or_else(|| format!("no section named {name}"))?;
882        let offset = self.header.resolve_any(&section.offset, "section offset")?;
883        let size = self.header.resolve_any(&section.size, "section size")?;
884        let end = offset.saturating_add(size);
885        if end > self.data.len() as u64 {
886            return Err(format!(
887                "section {name} runs from byte {offset} to {end}, past the file's {} bytes",
888                self.data.len()
889            ));
890        }
891        Ok(offset as usize..end as usize)
892    }
893}
894
895/// The raw value a delta field's `null` stands for.
896fn sentinel(field: &Field, bits: u32, signed: bool) -> Option<i128> {
897    use crate::formats::fixed_records::Null;
898    let bits = bits.clamp(1, 64);
899    match field.null? {
900        Null::Min if signed => Some(-(1i128 << (bits - 1))),
901        Null::Min => Some(0),
902        Null::Max if signed => Some((1i128 << (bits - 1)) - 1),
903        Null::Max => Some((1i128 << bits) - 1),
904        Null::Value(v) => Some(v),
905        Null::NaN => None,
906    }
907}
908
909/// An integer of `bytes`.
910fn int_of(bytes: &[u8], read: IntRead) -> i128 {
911    if read.signed {
912        i128::from(crate::formats::fixed_records::read_signed(bytes, read.big))
913    } else {
914        i128::from(crate::formats::fixed_records::read_unsigned(
915            bytes, read.big,
916        ))
917    }
918}
919
920/// A LEB128 value at the front of `bytes`, and the bytes it took.
921fn leb128(bytes: &[u8]) -> Option<(u64, usize)> {
922    let mut value = 0u64;
923    for (i, b) in bytes.iter().take(10).enumerate() {
924        let part = u64::from(b & 0x7f);
925        if i == 9 && part > 1 {
926            return None;
927        }
928        value |= part << (7 * i);
929        if b & 0x80 == 0 {
930            return Some((value, i + 1));
931        }
932    }
933    None
934}
935
936/// Text of `bytes` in `encoding`, trimmed as a fixed field's is.
937fn decode_text(bytes: &[u8], encoding: Encoding) -> String {
938    match encoding {
939        Encoding::Utf8 => crate::formats::fixed_records::text(bytes),
940        Encoding::Latin1 => crate::formats::fixed_records::latin1(bytes),
941        Encoding::Utf16Le => crate::formats::fixed_records::utf16(bytes, false),
942        Encoding::Utf16Be => crate::formats::fixed_records::utf16(bytes, true),
943    }
944}
945
946/// Where the NUL unit ends text at the front of `bytes`, if there is one.
947fn nul_at(bytes: &[u8], unit: usize) -> Option<usize> {
948    if unit == 1 {
949        memchr::memchr(0, bytes)
950    } else {
951        bytes
952            .chunks_exact(unit)
953            .position(|c| c.iter().all(|b| *b == 0))
954            .map(|i| i * unit)
955    }
956}
957
958/// What reading at a place found.
959#[derive(Debug, Clone, Copy, PartialEq, Eq)]
960enum Got {
961    /// A row.
962    Row,
963    /// A record of a variant not read.
964    Skipped,
965    /// No record left.
966    None,
967}
968
969/// What a walk pushes to: the sinks of one level's columns, and which it filled.
970struct Out<'s> {
971    sinks: &'s mut [Sink],
972    filled: &'s mut [bool],
973}
974
975impl Out<'_> {
976    /// Nulls in every column not filled, and a fresh start for the next row.
977    fn finish_row(&mut self) {
978        for (sink, filled) in self.sinks.iter_mut().zip(self.filled.iter_mut()) {
979            if !*filled {
980                sink.push_null();
981            }
982            *filled = false;
983        }
984    }
985}
986
987/// Reads records out of one run of bytes.
988struct Walker<'a> {
989    plan: &'a Plan,
990    /// The whole file, for sections.
991    file: &'a [u8],
992    frame: Frame,
993    acc: Vec<i128>,
994    /// The end of what the current record (or group item) may take.
995    end: usize,
996    /// Whether `end` is the record's own end, so a field past it is null.
997    bounded: bool,
998    /// A field did not fit, so the rest of the record is null.
999    short: bool,
1000    record_start: usize,
1001    /// The field that holds the record's size, and what to add to it.
1002    size_slot: Option<(usize, i64)>,
1003    /// Bytes skipped looking for sync markers.
1004    skipped: u64,
1005    /// The variant of the record read last, if it had one.
1006    variant: Option<usize>,
1007}
1008
1009impl<'a> Walker<'a> {
1010    fn new(plan: &'a Plan, file: &'a [u8]) -> Self {
1011        Self {
1012            plan,
1013            file,
1014            frame: Frame::new(plan.slots),
1015            acc: vec![0; plan.deltas.len()],
1016            end: 0,
1017            bounded: false,
1018            short: false,
1019            record_start: 0,
1020            size_slot: None,
1021            skipped: 0,
1022            variant: None,
1023        }
1024    }
1025
1026    fn size_of(&self, s: SizeRef, pos: usize, what: &str) -> Result<usize, Stop> {
1027        match s {
1028            SizeRef::Given(n) => Ok(n),
1029            SizeRef::Rest => Ok(self.end.saturating_sub(pos)),
1030            SizeRef::Slot { slot, adjust } => {
1031                let v = self.frame.ints[slot].ok_or(Stop::Truncated)? + i128::from(adjust);
1032                usize::try_from(v)
1033                    .ok()
1034                    .filter(|n| *n as u64 <= formats::MAX_SIZE)
1035                    .ok_or_else(|| {
1036                        Stop::Said(format!(
1037                            "`{what}` at byte {pos} is {v}, outside 0 to {}",
1038                            formats::MAX_SIZE
1039                        ))
1040                    })
1041            }
1042        }
1043    }
1044
1045    fn take(&self, pos: &mut usize, n: usize) -> Result<Range<usize>, Stop> {
1046        let start = *pos;
1047        let stop = start.checked_add(n).ok_or(Stop::Truncated)?;
1048        if stop > self.end {
1049            return Err(Stop::Truncated);
1050        }
1051        *pos = stop;
1052        Ok(start..stop)
1053    }
1054
1055    /// Read `fields` from `data` at `pos`, pushing to `out`. A field that does not fit
1056    /// is null once the record's end is known, and truncates the record before then.
1057    fn walk(
1058        &mut self,
1059        fields: &[FieldPlan],
1060        data: &[u8],
1061        pos: &mut usize,
1062        mut out: Option<&mut Out<'_>>,
1063    ) -> Result<(), Stop> {
1064        for f in fields {
1065            let res = if self.short {
1066                Err(Stop::Truncated)
1067            } else {
1068                self.field(f, data, pos, out.as_deref_mut())
1069            };
1070            match res {
1071                Ok(()) => {}
1072                Err(Stop::Truncated) if self.bounded => {
1073                    self.short = true;
1074                    if let Some(out) = out.as_deref_mut() {
1075                        null_outs(f, out);
1076                    }
1077                }
1078                Err(e) => return Err(e),
1079            }
1080            if let Some((slot, adjust)) = self.size_slot
1081                && slot == f.slot
1082                && !self.bounded
1083            {
1084                let start = self.record_start;
1085                let len = self.frame.ints[slot].ok_or(Stop::Truncated)? + i128::from(adjust);
1086                let total = usize::try_from(len)
1087                    .ok()
1088                    .filter(|t| *t as u64 <= formats::MAX_SIZE && *t >= *pos - start)
1089                    .ok_or_else(|| {
1090                        Stop::Said(format!(
1091                            "the record at byte {start} gives its size as {len}, outside {} to {}",
1092                            *pos - start,
1093                            formats::MAX_SIZE
1094                        ))
1095                    })?;
1096                let rec_end = start + total;
1097                if rec_end > self.end {
1098                    return Err(Stop::Truncated);
1099                }
1100                self.end = rec_end;
1101                self.bounded = true;
1102            }
1103        }
1104        Ok(())
1105    }
1106
1107    fn field(
1108        &mut self,
1109        f: &FieldPlan,
1110        data: &[u8],
1111        pos: &mut usize,
1112        out: Option<&mut Out<'_>>,
1113    ) -> Result<(), Stop> {
1114        self.frame.starts[f.slot] = Some(*pos);
1115        let count = match f.count {
1116            None => None,
1117            Some(c) => {
1118                let n = self.size_of(c, *pos, &f.name)?;
1119                if n as u64 > MAX_ITEMS {
1120                    return Err(Stop::Said(format!(
1121                        "`{}` at byte {} counts {n} values, more than {MAX_ITEMS}",
1122                        f.name, *pos
1123                    )));
1124                }
1125                Some(n)
1126            }
1127        };
1128        match &f.kind {
1129            Kind::Pad { size } => {
1130                let n = self.size_of(*size, *pos, "pad")?;
1131                self.take(pos, n)?;
1132            }
1133            Kind::Fixed { width, int } => {
1134                let cells = count.unwrap_or(1);
1135                let range = self.take(pos, width.checked_mul(cells).ok_or(Stop::Truncated)?)?;
1136                let bytes = &data[range];
1137                let raw = int.and_then(|r| (cells > 0).then(|| int_of(&bytes[..*width], r)));
1138                if count.is_none() {
1139                    self.frame.ints[f.slot] = raw;
1140                }
1141                self.value(f, raw, Some((bytes, *width)), count, out);
1142            }
1143            Kind::Var { signed } => {
1144                let cells = count.unwrap_or(1);
1145                let mut bytes = Vec::with_capacity(cells.min(64) * 8);
1146                let mut first = None;
1147                for _ in 0..cells {
1148                    let (v, n) =
1149                        leb128(&data[(*pos).min(self.end)..self.end]).ok_or(Stop::Truncated)?;
1150                    *pos += n;
1151                    let v = if *signed {
1152                        i128::from(((v >> 1) as i64) ^ -((v & 1) as i64))
1153                    } else {
1154                        i128::from(v)
1155                    };
1156                    first.get_or_insert(v);
1157                    if *signed {
1158                        bytes.extend((v as i64).to_le_bytes());
1159                    } else {
1160                        bytes.extend((v as u64).to_le_bytes());
1161                    }
1162                }
1163                if count.is_none() {
1164                    self.frame.ints[f.slot] = first;
1165                }
1166                self.value(f, first, Some((&bytes, 8)), count, out);
1167            }
1168            Kind::Sized { size, encoding } => {
1169                let n = self.size_of(*size, *pos, &f.name)?;
1170                let range = self.take(pos, n)?;
1171                if let Some(out) = out {
1172                    let col = f.outs[0];
1173                    out.filled[col] = true;
1174                    match (&mut out.sinks[col], encoding) {
1175                        (Sink::Text(v), Some(enc)) => {
1176                            v.push(Some(decode_text(&data[range], *enc)));
1177                        }
1178                        (Sink::Binary(v), None) => v.push(Some(data[range].to_vec())),
1179                        _ => {}
1180                    }
1181                }
1182            }
1183            Kind::Strz { max, encoding } => {
1184                let unit = encoding.unit();
1185                let text = match max {
1186                    Some(m) => {
1187                        let n = self.size_of(*m, *pos, &f.name)?;
1188                        let range = self.take(pos, n)?;
1189                        let bytes = &data[range];
1190                        let stop = nul_at(bytes, unit).unwrap_or(bytes.len());
1191                        decode_text(&bytes[..stop], *encoding)
1192                    }
1193                    None => {
1194                        let rest = &data[(*pos).min(self.end)..self.end];
1195                        let stop = nul_at(rest, unit).ok_or(Stop::Truncated)?;
1196                        let text = decode_text(&rest[..stop], *encoding);
1197                        *pos += stop + unit;
1198                        text
1199                    }
1200                };
1201                if let Some(out) = out {
1202                    let col = f.outs[0];
1203                    out.filled[col] = true;
1204                    if let Sink::Text(v) = &mut out.sinks[col] {
1205                        v.push(Some(text));
1206                    }
1207                }
1208            }
1209            Kind::StringAt {
1210                width,
1211                big,
1212                section,
1213            } => {
1214                let range = self.take(pos, *width)?;
1215                let offset = crate::formats::fixed_records::read_unsigned(&data[range], *big);
1216                self.frame.ints[f.slot] = Some(i128::from(offset));
1217                if let Some(out) = out {
1218                    let col = f.outs[0];
1219                    out.filled[col] = true;
1220                    let heap = &self.file[section.clone()];
1221                    let text = usize::try_from(offset)
1222                        .ok()
1223                        .filter(|o| *o < heap.len())
1224                        .map(|o| {
1225                            let rest = &heap[o..];
1226                            let stop = memchr::memchr(0, rest).unwrap_or(rest.len());
1227                            crate::formats::fixed_records::text(&rest[..stop])
1228                        });
1229                    if let Sink::Text(v) = &mut out.sinks[col] {
1230                        v.push(text);
1231                    }
1232                }
1233            }
1234            Kind::Group { items, slots } => {
1235                let n = count.unwrap_or(0);
1236                let first = items.first().map_or(0, |i| i.slot);
1237                let saved = (
1238                    self.end,
1239                    self.bounded,
1240                    self.short,
1241                    self.record_start,
1242                    self.size_slot,
1243                );
1244                self.size_slot = None;
1245                // Walked once to see the whole group fits, then again into the sinks, so a
1246                // group cut short leaves no half-pushed items.
1247                let start = *pos;
1248                let mut probe = start;
1249                let mut fits = Ok(());
1250                for _ in 0..n {
1251                    self.frame.clear(first..first + slots);
1252                    self.bounded = true;
1253                    self.short = false;
1254                    self.record_start = probe;
1255                    if let Err(e) = self.walk(items, data, &mut probe, None) {
1256                        fits = Err(e);
1257                        break;
1258                    }
1259                    if self.short {
1260                        fits = Err(Stop::Truncated);
1261                        break;
1262                    }
1263                }
1264                let (end, bounded, short, record_start, size_slot) = saved;
1265                self.end = end;
1266                if let Err(e) = fits {
1267                    self.bounded = bounded;
1268                    self.short = short;
1269                    self.record_start = record_start;
1270                    self.size_slot = size_slot;
1271                    return Err(e);
1272                }
1273                *pos = probe;
1274                if let Some(out) = out {
1275                    let col = f.outs[0];
1276                    out.filled[col] = true;
1277                    if let Sink::List {
1278                        inner,
1279                        offsets,
1280                        valid,
1281                    } = &mut out.sinks[col]
1282                        && let Sink::Struct { fields, len, .. } = inner.as_mut()
1283                    {
1284                        let mut at = start;
1285                        let mut filled = vec![false; fields.len()];
1286                        for _ in 0..n {
1287                            self.frame.clear(first..first + slots);
1288                            self.bounded = true;
1289                            self.short = false;
1290                            self.record_start = at;
1291                            let mut item_out = Out {
1292                                sinks: fields,
1293                                filled: &mut filled,
1294                            };
1295                            // It fitted above, and reads the same again.
1296                            let _ = self.walk(items, data, &mut at, Some(&mut item_out));
1297                            item_out.finish_row();
1298                            *len += 1;
1299                        }
1300                        offsets.push(offsets.last().copied().unwrap_or(0) + n as i64);
1301                        valid.push(true);
1302                    }
1303                }
1304                self.end = end;
1305                self.bounded = bounded;
1306                self.short = short;
1307                self.record_start = record_start;
1308                self.size_slot = size_slot;
1309            }
1310        }
1311        self.frame.ends[f.slot] = Some(*pos);
1312        Ok(())
1313    }
1314
1315    /// Push a number (and its bits, and its running sum) to `out`, or only keep the sum.
1316    fn value(
1317        &mut self,
1318        f: &FieldPlan,
1319        raw: Option<i128>,
1320        bytes: Option<(&[u8], usize)>,
1321        count: Option<usize>,
1322        out: Option<&mut Out<'_>>,
1323    ) {
1324        let summed = if f.delta == Delta::None {
1325            None
1326        } else {
1327            let raw = raw.unwrap_or(0);
1328            if Some(raw) == f.sentinel {
1329                Some(None)
1330            } else {
1331                let sum = self.acc[f.delta_index].wrapping_add(raw);
1332                self.acc[f.delta_index] = sum;
1333                Some(i64::try_from(sum).ok())
1334            }
1335        };
1336        let Some(out) = out else { return };
1337        match (summed, bytes) {
1338            (Some(sum), _) => {
1339                let col = f.outs[0];
1340                out.filled[col] = true;
1341                match sum {
1342                    Some(v) => out.sinks[col].push_bytes(&v.to_le_bytes()),
1343                    None => out.sinks[col].push_null(),
1344                }
1345            }
1346            (None, Some((bytes, width))) => push_fixed(f, bytes, width, count, out),
1347            (None, None) => {}
1348        }
1349        push_bits(f, raw, out);
1350    }
1351
1352    /// Read one record at `pos`, before `chunk_end`; `pos` ends after it.
1353    fn record(
1354        &mut self,
1355        data: &[u8],
1356        pos: &mut usize,
1357        chunk_end: usize,
1358        out: Option<&mut Out<'_>>,
1359        time: Option<i64>,
1360    ) -> Result<Got, Stop> {
1361        let Some(only) = self.plan.only else {
1362            return Ok(match self.record_inner(data, pos, chunk_end, out, time)? {
1363                None => Got::None,
1364                Some(_) => Got::Row,
1365            });
1366        };
1367        // One variant read alone: the record is walked to learn its variant, and walked
1368        // again into the columns when it is the one.
1369        let (start, acc, skipped) = (*pos, self.acc.clone(), self.skipped);
1370        match self.record_inner(data, pos, chunk_end, None, time)? {
1371            None => Ok(Got::None),
1372            Some(Some(v)) if v == only => {
1373                if out.is_some() {
1374                    *pos = start;
1375                    self.acc = acc;
1376                    self.skipped = skipped;
1377                    self.record_inner(data, pos, chunk_end, out, time)?;
1378                }
1379                Ok(Got::Row)
1380            }
1381            Some(_) => Ok(Got::Skipped),
1382        }
1383    }
1384
1385    /// Read one record: `None` when there is none left (for sync framing, no marker),
1386    /// else the index of its variant, if it has one.
1387    fn record_inner(
1388        &mut self,
1389        data: &[u8],
1390        pos: &mut usize,
1391        chunk_end: usize,
1392        mut out: Option<&mut Out<'_>>,
1393        time: Option<i64>,
1394    ) -> Result<Option<Option<usize>>, Stop> {
1395        let plan = self.plan;
1396        if *pos >= chunk_end {
1397            return Ok(None);
1398        }
1399        if !plan.sync.is_empty() {
1400            let rest = &data[*pos..chunk_end];
1401            if !rest.starts_with(&plan.sync) {
1402                match memchr::memmem::find(rest, &plan.sync) {
1403                    Some(at) => {
1404                        self.skipped += at as u64;
1405                        *pos += at;
1406                    }
1407                    None => {
1408                        self.skipped += rest.len() as u64;
1409                        *pos = chunk_end;
1410                        return Ok(None);
1411                    }
1412                }
1413            }
1414            *pos += plan.sync.len();
1415        }
1416        let first = plan.chunk_slots;
1417        self.frame.clear(first..plan.slots);
1418        self.record_start = *pos;
1419        self.end = chunk_end;
1420        self.bounded = false;
1421        self.short = false;
1422        self.size_slot = None;
1423        match plan.size {
1424            Some(SizeRef::Given(n)) => {
1425                let end = pos.checked_add(n).ok_or(Stop::Truncated)?;
1426                if end > chunk_end {
1427                    return Err(Stop::Truncated);
1428                }
1429                self.end = end;
1430                self.bounded = true;
1431            }
1432            Some(SizeRef::Slot { slot, adjust }) => self.size_slot = Some((slot, adjust)),
1433            _ => {}
1434        }
1435        let start = *pos;
1436        let mut p = start;
1437        self.walk(&plan.common, data, &mut p, out.as_deref_mut())?;
1438        let mut label = None;
1439        let mut chosen = None;
1440        if let Some((slot, text)) = plan.type_slot {
1441            let int = self.frame.ints[slot];
1442            let shown = if text {
1443                match (self.frame.starts[slot], self.frame.ends[slot]) {
1444                    (Some(a), Some(b)) => Some(crate::formats::fixed_records::text(&data[a..b])),
1445                    _ => None,
1446                }
1447            } else {
1448                None
1449            };
1450            let variant = plan
1451                .variants
1452                .iter()
1453                .enumerate()
1454                .find(|(_, v)| match (&shown, int) {
1455                    (Some(t), _) => v.texts.iter().any(|w| w == t),
1456                    (None, Some(i)) => v.ints.contains(&i),
1457                    _ => false,
1458                });
1459            match variant {
1460                Some((index, variant)) => {
1461                    chosen = Some(index);
1462                    if let Some(size) = variant.size
1463                        && !self.bounded
1464                    {
1465                        let n = self.size_of(size, p, "size")?;
1466                        let end = start.checked_add(n).ok_or(Stop::Truncated)?;
1467                        if end > chunk_end || end < p {
1468                            return Err(if end > chunk_end {
1469                                Stop::Truncated
1470                            } else {
1471                                Stop::Said(format!(
1472                                    "variant {} at byte {start} is {n} bytes, less than its common fields",
1473                                    variant.name
1474                                ))
1475                            });
1476                        }
1477                        self.end = end;
1478                        self.bounded = true;
1479                    }
1480                    self.walk(&variant.fields, data, &mut p, out.as_deref_mut())?;
1481                    label = Some(variant.name.clone());
1482                }
1483                None => {
1484                    let shown = shown.or_else(|| int.map(|i| i.to_string()));
1485                    if !self.bounded {
1486                        return Err(Stop::Said(format!(
1487                            "the record at byte {start} has type {}, which no variant names",
1488                            shown.as_deref().unwrap_or("(none)")
1489                        )));
1490                    }
1491                    label = shown.map(|s| Arc::from(format!("?{s}")));
1492                }
1493            }
1494        }
1495        self.variant = chosen;
1496        *pos = if self.bounded { self.end } else { p };
1497        if let Some((width, big)) = plan.suffix {
1498            let range = self.take_at(*pos, width, chunk_end)?;
1499            let again = crate::formats::fixed_records::read_unsigned(&data[range], big);
1500            let (slot, _) = self.size_slot.unwrap_or((usize::MAX, 0));
1501            let first = self.frame.ints.get(slot).copied().flatten();
1502            if first != Some(i128::from(again)) {
1503                return Err(Stop::Said(format!(
1504                    "the record at byte {start} ends with length {again}, not the {} it starts with",
1505                    first.map_or_else(|| "?".to_string(), |v| v.to_string())
1506                )));
1507            }
1508            *pos += width;
1509        }
1510        if let Some(out) = out {
1511            if let (Some(col), Some(label)) = (plan.type_out, label)
1512                && let Sink::Label(v) = &mut out.sinks[col]
1513            {
1514                v.push(Some(label));
1515                out.filled[col] = true;
1516            }
1517            if let Some(check) = &plan.checksum {
1518                let from = check.from.map_or(Some(start), |s| self.frame.starts[s]);
1519                let to = self.frame.starts[check.to];
1520                let stored = self.frame.ints[check.field];
1521                let ok = match (from, to, stored) {
1522                    (Some(a), Some(b), Some(v)) if a <= b => {
1523                        Some(i128::from(check.algo.compute(&data[a..b])) == v)
1524                    }
1525                    _ => None,
1526                };
1527                if let Sink::Flag(v) = &mut out.sinks[check.out] {
1528                    v.push(ok);
1529                    out.filled[check.out] = true;
1530                }
1531            }
1532            if let Some(col) = plan.time_out
1533                && let Sink::Time(v) = &mut out.sinks[col]
1534            {
1535                v.push(time);
1536                out.filled[col] = true;
1537            }
1538            out.finish_row();
1539        }
1540        Ok(Some(chosen))
1541    }
1542
1543    fn take_at(&self, pos: usize, n: usize, end: usize) -> Result<Range<usize>, Stop> {
1544        let stop = pos.checked_add(n).ok_or(Stop::Truncated)?;
1545        if stop > end {
1546            return Err(Stop::Truncated);
1547        }
1548        Ok(pos..stop)
1549    }
1550}
1551
1552/// The bit fields of `f`'s value `raw` to their columns.
1553fn push_bits(f: &FieldPlan, raw: Option<i128>, out: &mut Out<'_>) {
1554    for (col, bit, width) in &f.bits {
1555        out.filled[*col] = true;
1556        match (raw, &mut out.sinks[*col]) {
1557            (Some(raw), Sink::Bits { values, .. }) => {
1558                let mask = if *width >= 64 {
1559                    u64::MAX
1560                } else {
1561                    (1u64 << width) - 1
1562                };
1563                values.push(Some(((raw as u64) >> bit) & mask));
1564            }
1565            (_, sink) => sink.push_null(),
1566        }
1567    }
1568}
1569
1570/// Nulls in every column of `f`.
1571fn null_outs(f: &FieldPlan, out: &mut Out<'_>) {
1572    for col in f.outs.iter().chain(f.bits.iter().map(|(c, _, _)| c)) {
1573        if !out.filled[*col] {
1574            out.sinks[*col].push_null();
1575            out.filled[*col] = true;
1576        }
1577    }
1578}
1579
1580fn push_fixed(f: &FieldPlan, bytes: &[u8], width: usize, count: Option<usize>, out: &mut Out<'_>) {
1581    if f.outs.len() == 1 {
1582        let col = f.outs[0];
1583        out.filled[col] = true;
1584        match (&mut out.sinks[col], count) {
1585            (
1586                Sink::List {
1587                    inner,
1588                    offsets,
1589                    valid,
1590                },
1591                Some(n),
1592            ) => {
1593                for i in 0..n {
1594                    inner.push_bytes(&bytes[i * width..(i + 1) * width]);
1595                }
1596                offsets.push(offsets.last().copied().unwrap_or(0) + n as i64);
1597                valid.push(true);
1598            }
1599            (sink, _) => sink.push_bytes(bytes),
1600        }
1601    } else {
1602        for (i, col) in f.outs.iter().enumerate() {
1603            out.filled[*col] = true;
1604            out.sinks[*col].push_bytes(&bytes[i * width..(i + 1) * width]);
1605        }
1606    }
1607}
1608
1609/// Bytes of one chunk: a range of the map, or a decompressed block.
1610enum ChunkData {
1611    Map(Range<usize>),
1612    Owned(Arc<Vec<u8>>),
1613}
1614
1615/// Where a read is: in which chunk, where, and how many of its records it has taken.
1616struct Cursor {
1617    pos: usize,
1618    data_start: usize,
1619    end: usize,
1620    taken: u64,
1621    limit: Option<u64>,
1622    time: Option<i64>,
1623}
1624
1625/// The rows a read wants, and how many it has.
1626struct Want {
1627    start: u64,
1628    len: usize,
1629    produced: usize,
1630}
1631
1632impl Want {
1633    /// A null row at `row`, if it is one wanted.
1634    fn null_row(&mut self, row: u64, out: &mut Out<'_>) {
1635        if row >= self.start && self.produced < self.len {
1636            out.finish_row();
1637            self.produced += 1;
1638        }
1639    }
1640}
1641
1642/// What opening found to say: trailing bytes, skipped bytes, records cut short.
1643#[derive(Default)]
1644struct Found {
1645    notes: Vec<String>,
1646    skipped: u64,
1647}
1648
1649impl FramedRecords {
1650    /// Read `spec`'s records from `bytes[data]`, with `header` read already.
1651    pub fn open(
1652        spec: &Spec,
1653        bytes: Arc<Bytes>,
1654        header: &HeaderValues,
1655        data: Range<usize>,
1656        named: &str,
1657        path: Option<&std::path::Path>,
1658    ) -> Result<(Self, Vec<String>), String> {
1659        let file = bytes.as_slice();
1660        let mut compiler = Compiler {
1661            spec,
1662            header,
1663            data: file,
1664            slots: 0,
1665            deltas: Vec::new(),
1666        };
1667        let mut columns = Vec::new();
1668        let mut chunk_scope = Scope::default();
1669        let mut discard = Vec::new();
1670        let chunk_header = match &spec.capture {
1671            Some(capture) => {
1672                compiler.fields(&capture.header, &mut chunk_scope, &mut discard, false)?
1673            }
1674            None => Vec::new(),
1675        };
1676        let chunk_count = spec
1677            .capture
1678            .as_ref()
1679            .and_then(|c| c.count.as_ref())
1680            .and_then(|n| chunk_scope.slot(n));
1681        let chunk_slots = compiler.slots;
1682        let time_out = spec
1683            .capture
1684            .as_ref()
1685            .and_then(|c| c.time.as_ref())
1686            .map(|name| {
1687                Compiler::column(
1688                    &mut columns,
1689                    name,
1690                    DataType::Datetime(TimeUnit::Nanoseconds, None),
1691                    Sink::Time(Vec::new()),
1692                    false,
1693                )
1694            });
1695        let records = &spec.records;
1696        let mut scope = Scope::default();
1697        let common = compiler.fields(&records.fields, &mut scope, &mut columns, false)?;
1698        let type_slot = records.type_field.as_ref().map(|name| {
1699            let field = records
1700                .fields
1701                .iter()
1702                .find(|f| f.name.as_deref() == Some(name.as_str()))
1703                .expect("checked at parse");
1704            (
1705                scope.slot(name).expect("a common field"),
1706                field.ty.is_text(),
1707            )
1708        });
1709        let only = match &spec.variant {
1710            Some(name) => Some(
1711                records
1712                    .variants
1713                    .iter()
1714                    .position(|v| &v.name == name)
1715                    .ok_or_else(|| format!("no variant named {name}"))?,
1716            ),
1717            None => None,
1718        };
1719        let type_out = type_slot.filter(|_| only.is_none()).map(|_| {
1720            Compiler::column(
1721                &mut columns,
1722                "type",
1723                DataType::from_categories(Categories::global()),
1724                Sink::Label(Vec::new()),
1725                false,
1726            )
1727        });
1728        let size = records
1729            .size
1730            .as_ref()
1731            .map(|s| compiler.size(s, &scope, "record size"))
1732            .transpose()?;
1733        let mut variants = Vec::new();
1734        let mut unread = Vec::new();
1735        for (i, variant) in records.variants.iter().enumerate() {
1736            let mut vscope = scope.clone();
1737            let into = if only.is_some_and(|o| o != i) {
1738                &mut unread
1739            } else {
1740                &mut columns
1741            };
1742            let fields = compiler.fields(&variant.fields, &mut vscope, into, true)?;
1743            let size = variant
1744                .size
1745                .as_ref()
1746                .map(|s| compiler.size(s, &vscope, "variant size"))
1747                .transpose()?;
1748            if matches!(size, Some(SizeRef::Slot { .. })) {
1749                return Err(format!(
1750                    "variant {}: its size comes from the header or is written in the spec",
1751                    variant.name
1752                ));
1753            }
1754            let mut ints = Vec::new();
1755            let mut texts = Vec::new();
1756            for when in &variant.when {
1757                match when {
1758                    formats::Expected::Int(v) => ints.push(*v),
1759                    // Text fields read with their padding trimmed: `"fmt "` is `fmt`.
1760                    formats::Expected::Text(t) => {
1761                        texts.push(t.trim_end_matches(['\0', ' ']).to_string())
1762                    }
1763                }
1764            }
1765            variants.push(VariantPlan {
1766                name: Arc::from(variant.name.as_str()),
1767                ints,
1768                texts,
1769                fields,
1770                size,
1771            });
1772        }
1773        let suffix = if records.length_suffix {
1774            let Some(Amount::Record { field, .. }) = &records.size else {
1775                return Err("length_suffix: needs size from a field".into());
1776            };
1777            let target = records
1778                .fields
1779                .iter()
1780                .find(|f| f.name.as_deref() == Some(field.as_str()))
1781                .ok_or_else(|| format!("length_suffix: `{field}` is not a common field"))?;
1782            let width = target
1783                .ty
1784                .width()
1785                .ok_or("length_suffix: the length field is fixed-width")?;
1786            let big = target.endian.unwrap_or(spec.endian) == formats::Endian::Big;
1787            Some((width as usize, big))
1788        } else {
1789            None
1790        };
1791        let checksum = records
1792            .checksum
1793            .as_ref()
1794            .map(|c| -> Result<ChecksumPlan, String> {
1795                // A field of a variant has a slot of its own in each scope; the checksum
1796                // names one by name, so look in every variant's.
1797                let slot_of = |name: &str| -> Option<usize> {
1798                    scope.slot(name).or_else(|| {
1799                        variants
1800                            .iter()
1801                            .find_map(|v| v.fields.iter().find(|f| f.name == name).map(|f| f.slot))
1802                    })
1803                };
1804                let field = slot_of(&c.field).ok_or("checksum: no such field")?;
1805                let to =
1806                    c.to.as_deref()
1807                        .map_or(Some(field), slot_of)
1808                        .ok_or("checksum: no such field")?;
1809                let from = c
1810                    .from
1811                    .as_deref()
1812                    .map(slot_of)
1813                    .map(|s| s.ok_or("checksum: no such field"))
1814                    .transpose()?;
1815                let out = Compiler::column(
1816                    &mut columns,
1817                    "checksum_ok",
1818                    DataType::Boolean,
1819                    Sink::Flag(Vec::new()),
1820                    false,
1821                );
1822                Ok(ChecksumPlan {
1823                    algo: c.algo,
1824                    field,
1825                    from,
1826                    to,
1827                    out,
1828                })
1829            })
1830            .transpose()?;
1831        if columns.is_empty() {
1832            return Err("the records have no named field".into());
1833        }
1834        // Record size known up front: written, or what fixed fields take.
1835        let fixed_size = match size {
1836            Some(SizeRef::Given(n)) => Some(n),
1837            None if variants.is_empty() => common.iter().try_fold(0usize, |sum, f| {
1838                let width = match (&f.kind, f.count) {
1839                    (Kind::Fixed { width, .. }, None) => *width,
1840                    (Kind::Fixed { width, .. }, Some(SizeRef::Given(n))) => width.checked_mul(n)?,
1841                    (
1842                        Kind::Pad {
1843                            size: SizeRef::Given(n),
1844                        },
1845                        None,
1846                    ) => *n,
1847                    _ => return None,
1848                };
1849                sum.checked_add(width)
1850            }),
1851            _ => None,
1852        };
1853        let size = match (size, fixed_size, records.framing) {
1854            (None, Some(n), Framing::Fixed)
1855                if !common.iter().any(|f| matches!(f.kind, Kind::Group { .. })) =>
1856            {
1857                Some(SizeRef::Given(n))
1858            }
1859            (size, _, _) => size,
1860        };
1861        if matches!(size, Some(SizeRef::Given(0))) && records.sync.is_empty() {
1862            return Err("a record takes no bytes".into());
1863        }
1864        let plan = Plan {
1865            framing: records.framing,
1866            common,
1867            type_slot,
1868            variants,
1869            type_out,
1870            size,
1871            suffix,
1872            align: records.align as usize,
1873            sync: records.sync.clone(),
1874            checksum,
1875            chunk_header,
1876            chunk_count,
1877            chunk_slots,
1878            time_out,
1879            slots: compiler.slots,
1880            deltas: compiler.deltas,
1881            columns,
1882            only,
1883            sources: Vec::new(),
1884        };
1885        let plan = Plan {
1886            sources: sources(&plan),
1887            ..plan
1888        };
1889        let mut found = Found::default();
1890        let chunks = if let Some(capture) = &spec.capture {
1891            let _ = capture;
1892            crate::formats::framed_records::capture::packets(file, &mut found.notes)?
1893        } else if let Some(blocks) = &spec.blocks {
1894            list_blocks(spec, blocks, header, file, data.clone(), &mut found.notes)?
1895        } else {
1896            vec![Chunk {
1897                source: ChunkSource::Map(data.clone()),
1898                records: None,
1899                time_ns: None,
1900            }]
1901        };
1902        let schema: Schema = plan
1903            .columns
1904            .iter()
1905            .map(|c| polars::prelude::Field::new(c.name.clone(), c.dtype.clone()))
1906            .collect();
1907        let mut records_read = Self {
1908            bytes,
1909            plan: Arc::new(plan),
1910            chunks,
1911            index: Arc::new(Index::Walk(Vec::new())),
1912            table: None,
1913            rows: 0,
1914            schema: Arc::new(schema),
1915            cache: Mutex::new(VecDeque::new()),
1916        };
1917        let ring = records
1918            .ring
1919            .as_ref()
1920            .map(|r| header.resolve_any(r, "ring"))
1921            .transpose()?;
1922        let count = records
1923            .count
1924            .as_ref()
1925            .map(|c| header.resolve_any(c, "count"))
1926            .transpose()?;
1927        // The walk of this file with this spec, kept from an earlier open of it.
1928        let kept = path
1929            .and_then(crate::formats::indexed::peek::<KeptWalk>)
1930            .filter(|k| k.spec == *spec && k.spec.variant == spec.variant && k.data == data);
1931        if let Some(kept) = kept {
1932            records_read.index = kept.index.clone();
1933            records_read.table = kept.table.clone();
1934            records_read.rows = kept.rows;
1935            found.notes.extend(kept.notes.iter().cloned());
1936            return Ok((records_read, found.notes));
1937        }
1938        let before = found.notes.len();
1939        records_read.build_index(named, ring, count, &mut found)?;
1940        if found.skipped > 0 {
1941            found.notes.push(format!(
1942                "{} {} skipped between records, to the next sync marker",
1943                found.skipped,
1944                if found.skipped == 1 { "byte" } else { "bytes" }
1945            ));
1946        }
1947        if let Some(path) = path
1948            && matches!(*records_read.index, Index::Walk(_))
1949        {
1950            crate::formats::indexed::keep(
1951                path,
1952                Arc::new(KeptWalk {
1953                    spec: spec.clone(),
1954                    data,
1955                    index: records_read.index.clone(),
1956                    table: records_read.table.clone(),
1957                    rows: records_read.rows,
1958                    notes: found.notes[before..].to_vec(),
1959                }),
1960            );
1961        }
1962        Ok((records_read, found.notes))
1963    }
1964
1965    /// Whether a stride of `size` bytes finds every record: one run of fixed records.
1966    fn stride(&self) -> Option<usize> {
1967        let plan = &self.plan;
1968        match (plan.size, &self.chunks[..]) {
1969            (Some(SizeRef::Given(n)), [chunk])
1970                if plan.framing == Framing::Fixed
1971                    && plan.only.is_none()
1972                    && plan.sync.is_empty()
1973                    && plan.suffix.is_none()
1974                    // A running sum needs every record before the one read.
1975                    && plan.deltas.is_empty()
1976                    && matches!(chunk.source, ChunkSource::Map(_))
1977                    && plan.chunk_header.is_empty() =>
1978            {
1979                let align = plan.align.max(1);
1980                Some(n.div_ceil(align) * align)
1981            }
1982            _ => None,
1983        }
1984    }
1985
1986    fn build_index(
1987        &mut self,
1988        named: &str,
1989        ring: Option<u64>,
1990        limit: Option<u64>,
1991        found: &mut Found,
1992    ) -> Result<(), String> {
1993        let file = self.bytes.as_slice();
1994        if let Some(size) = self.stride() {
1995            let ChunkSource::Map(range) = self.chunks[0].source.clone() else {
1996                unreachable!("a stride is over the map")
1997            };
1998            let len = range.len();
1999            let mut rows = len / size;
2000            // The last record need not be padded out to the alignment.
2001            let unpadded = self.plan.size.map_or(size, |s| match s {
2002                SizeRef::Given(n) => n,
2003                _ => size,
2004            });
2005            if len % size >= unpadded {
2006                rows += 1;
2007            }
2008            if let Some(count) = limit {
2009                if count < rows as u64 {
2010                    rows = count as usize;
2011                } else if count > rows as u64 {
2012                    found.notes.push(format!(
2013                        "header says {count} records {} {rows} whole ones shown",
2014                        crate::glyphs::get().middot
2015                    ));
2016                }
2017            } else {
2018                let used = (rows * size).min(len);
2019                let trailing = len - used;
2020                if trailing > 0 && len % size < unpadded {
2021                    found
2022                        .notes
2023                        .push(trailing_note(named, &file[range.end - trailing..range.end]));
2024                }
2025            }
2026            let rows = rows.min(MAX_ROWS);
2027            let ring = match ring {
2028                Some(r) if rows > 0 => {
2029                    if r >= rows as u64 {
2030                        found.notes.push(format!(
2031                            "ring's oldest record {r} past the {rows} records {} read from the first",
2032                            crate::glyphs::get().middot
2033                        ));
2034                        0
2035                    } else {
2036                        r as usize
2037                    }
2038                }
2039                _ => 0,
2040            };
2041            self.index = Arc::new(Index::Stride {
2042                start: range.start,
2043                size,
2044                ring,
2045            });
2046            self.rows = rows;
2047            return Ok(());
2048        }
2049        let plan = self.plan.clone();
2050        let mut walker = Walker::new(&plan, file);
2051        let mut checkpoints = Vec::new();
2052        let mut rows: u64 = 0;
2053        let max_rows = MAX_ROWS as u64;
2054        let needs_walk_everything = plan.deltas.contains(&Delta::All);
2055        // Read once: the limits sit behind a lock, and this loop runs per record.
2056        let indexed_records = crate::limits::get().indexed_records;
2057        let mut table = (self.chunks.len() == 1
2058            && matches!(self.chunks[0].source, ChunkSource::Map(_))
2059            && plan.chunk_header.is_empty()
2060            && plan.variants.len() < usize::from(WALK))
2061        .then(|| RowTable {
2062            starts: crate::formats::indexed::Offsets::for_file(file.len()),
2063            tags: Vec::new(),
2064        });
2065        'chunks: for ci in 0..self.chunks.len() {
2066            if limit.is_some_and(|l| rows >= l) || rows >= max_rows {
2067                break;
2068            }
2069            reset_block_sums(&plan, &mut walker.acc);
2070            if let Some(n) = self.chunks[ci].records
2071                && !needs_walk_everything
2072            {
2073                checkpoints.push(Checkpoint {
2074                    row: rows,
2075                    chunk: ci as u32,
2076                    pos: u32::MAX,
2077                    taken: 0,
2078                    acc: walker.acc.clone().into_boxed_slice(),
2079                });
2080                rows = rows
2081                    .saturating_add(n)
2082                    .min(limit.unwrap_or(u64::MAX))
2083                    .min(max_rows);
2084                continue;
2085            }
2086            let (data, mut cursor) = match self.enter(ci, &mut walker) {
2087                Ok(c) => c,
2088                Err(e) => {
2089                    found.notes.push(format!("block {ci}: {e}; left out"));
2090                    continue;
2091                }
2092            };
2093            let chunk_bytes = self.slice(&data, file);
2094            loop {
2095                if cursor.limit.is_some_and(|l| cursor.taken >= l)
2096                    || limit.is_some_and(|l| rows >= l)
2097                    || rows >= max_rows
2098                {
2099                    break;
2100                }
2101                if cursor.taken.is_multiple_of(CHECKPOINT) {
2102                    checkpoints.push(Checkpoint {
2103                        row: rows,
2104                        chunk: ci as u32,
2105                        pos: cursor.pos as u32,
2106                        taken: cursor.taken as u32,
2107                        acc: walker.acc.clone().into_boxed_slice(),
2108                    });
2109                }
2110                let before = cursor.pos;
2111                match walker.record(chunk_bytes, &mut cursor.pos, cursor.end, None, None) {
2112                    Ok(got @ (Got::Row | Got::Skipped)) => {
2113                        if got == Got::Row {
2114                            rows += 1;
2115                            if let Some(t) = table.as_mut() {
2116                                if t.tags.len() >= indexed_records {
2117                                    table = None;
2118                                } else {
2119                                    let tag = match (walker.short, plan.type_slot) {
2120                                        (true, _) => WALK,
2121                                        (false, None) => 0,
2122                                        (false, Some(_)) => walker
2123                                            .variant
2124                                            .and_then(|v| u8::try_from(v).ok())
2125                                            .unwrap_or(WALK),
2126                                    };
2127                                    t.starts.push(walker.record_start - plan.sync.len());
2128                                    t.tags.push(tag);
2129                                }
2130                            }
2131                        }
2132                        cursor.taken += 1;
2133                        align(&mut cursor, plan.align);
2134                        if cursor.pos <= before {
2135                            found.notes.push(format!(
2136                                "zero-length record at byte {before} {} rest left out",
2137                                crate::glyphs::get().middot
2138                            ));
2139                            break 'chunks;
2140                        }
2141                    }
2142                    Ok(Got::None) => {
2143                        // A checkpoint pushed for a record that is not there.
2144                        if checkpoints.last().is_some_and(|c: &Checkpoint| {
2145                            c.row == rows && c.chunk == ci as u32 && c.pos == before as u32
2146                        }) {
2147                            checkpoints.pop();
2148                        }
2149                        break;
2150                    }
2151                    Err(stop) => {
2152                        if checkpoints.last().is_some_and(|c: &Checkpoint| {
2153                            c.row == rows && c.chunk == ci as u32 && c.pos == before as u32
2154                        }) {
2155                            checkpoints.pop();
2156                        }
2157                        let what = if self.chunks.len() > 1 {
2158                            format!("{named}, block {ci},")
2159                        } else {
2160                            named.to_string()
2161                        };
2162                        match stop {
2163                            Stop::Truncated => {
2164                                let rest = &chunk_bytes[before..cursor.end];
2165                                found.notes.push(trailing_note(&what, rest));
2166                            }
2167                            Stop::Said(said) => {
2168                                found.notes.push(format!(
2169                                    "{said} {} {} bytes from there left out",
2170                                    crate::glyphs::get().middot,
2171                                    cursor.end - before
2172                                ));
2173                            }
2174                        }
2175                        if self.chunks.len() == 1 {
2176                            break 'chunks;
2177                        }
2178                        break;
2179                    }
2180                }
2181                if cursor.pos >= cursor.end {
2182                    break;
2183                }
2184            }
2185            // A chunk that counts its records and holds fewer is filled with nulls.
2186            if let Some(l) = cursor.limit
2187                && cursor.taken < l
2188                && self.chunks[ci].records.is_some()
2189            {
2190                found.notes.push(format!(
2191                    "block {ci}: {l} records declared, {} found",
2192                    cursor.taken
2193                ));
2194            }
2195        }
2196        if let Some(l) = limit
2197            && rows < l
2198        {
2199            found.notes.push(format!(
2200                "header says {l} records {} {rows} shown",
2201                crate::glyphs::get().middot
2202            ));
2203        }
2204        found.skipped += walker.skipped;
2205        self.rows = rows as usize;
2206        self.index = Arc::new(Index::Walk(checkpoints));
2207        self.table = table.filter(|t| t.tags.len() == self.rows).map(|mut t| {
2208            t.starts.shrink();
2209            t.tags.shrink_to_fit();
2210            Arc::new(t)
2211        });
2212        Ok(())
2213    }
2214
2215    fn slice<'b>(&'b self, data: &'b ChunkData, file: &'b [u8]) -> &'b [u8] {
2216        match data {
2217            ChunkData::Map(_) => file,
2218            ChunkData::Owned(v) => v,
2219        }
2220    }
2221
2222    /// The bytes of chunk `i`: its range of the map, or its block decompressed.
2223    fn chunk_data(&self, i: usize) -> Result<ChunkData, String> {
2224        match &self.chunks[i].source {
2225            ChunkSource::Map(range) => Ok(ChunkData::Map(range.clone())),
2226            ChunkSource::Block {
2227                body,
2228                codec,
2229                uncompressed,
2230            } => {
2231                if let Ok(cache) = self.cache.lock()
2232                    && let Some((_, data)) = cache.iter().find(|(k, _)| *k == i)
2233                {
2234                    return Ok(ChunkData::Owned(data.clone()));
2235                }
2236                self.bytes.still_whole().map_err(|e| e.to_string())?;
2237                let raw = &self.bytes.as_slice()[body.clone()];
2238                let data = Arc::new(decompress(raw, *codec, *uncompressed)?);
2239                if let Ok(mut cache) = self.cache.lock() {
2240                    cache.push_front((i, data.clone()));
2241                    cache.truncate(CACHED_BLOCKS);
2242                }
2243                Ok(ChunkData::Owned(data))
2244            }
2245        }
2246    }
2247
2248    /// Chunk `i`'s bytes, and a cursor at its start, after its payload header.
2249    fn enter(&self, i: usize, walker: &mut Walker<'_>) -> Result<(ChunkData, Cursor), String> {
2250        let data = self.chunk_data(i)?;
2251        let (start, end) = match &data {
2252            ChunkData::Map(r) => (r.start, r.end),
2253            ChunkData::Owned(v) => (0, v.len()),
2254        };
2255        let mut cursor = Cursor {
2256            pos: start,
2257            data_start: start,
2258            end,
2259            taken: 0,
2260            limit: self.chunks[i].records,
2261            time: self.chunks[i].time_ns,
2262        };
2263        if !self.plan.chunk_header.is_empty() {
2264            let file = self.bytes.as_slice();
2265            let bytes = self.slice(&data, file);
2266            walker.frame.clear(0..self.plan.chunk_slots);
2267            walker.end = end;
2268            walker.bounded = false;
2269            walker.short = false;
2270            walker.size_slot = None;
2271            walker.record_start = start;
2272            let mut p = start;
2273            match walker.walk(&self.plan.chunk_header, bytes, &mut p, None) {
2274                Ok(()) => {
2275                    cursor.pos = p;
2276                    cursor.data_start = p;
2277                    if let Some(slot) = self.plan.chunk_count {
2278                        cursor.limit = walker.frame.ints[slot].map(|v| v.max(0) as u64);
2279                    }
2280                }
2281                Err(_) => {
2282                    // Too short for its header: nothing in it.
2283                    cursor.pos = end;
2284                    cursor.limit = Some(0);
2285                }
2286            }
2287        }
2288        Ok((data, cursor))
2289    }
2290
2291    /// Rows `[start, start + len)`, of the columns `wanted` names, decoded now.
2292    fn decode(&self, start: usize, len: usize, wanted: &[bool]) -> PolarsResult<DataFrame> {
2293        self.bytes.still_whole()?;
2294        let start = start.min(self.rows);
2295        let len = len.min(self.rows - start);
2296        let plan = &*self.plan;
2297        let mut sinks: Vec<Sink> = plan
2298            .columns
2299            .iter()
2300            .zip(wanted)
2301            .map(|(c, w)| if *w { c.proto.clone() } else { Sink::Skip })
2302            .collect();
2303        let mut filled = vec![false; sinks.len()];
2304        let file = self.bytes.as_slice();
2305        let mut walker = Walker::new(plan, file);
2306        match &*self.index {
2307            Index::Stride {
2308                start: base,
2309                size,
2310                ring,
2311            } => {
2312                let mut out = Out {
2313                    sinks: &mut sinks,
2314                    filled: &mut filled,
2315                };
2316                let ChunkSource::Map(range) = &self.chunks[0].source else {
2317                    unreachable!("a stride is over the map")
2318                };
2319                for row in start..start + len {
2320                    let k = (ring + row) % self.rows.max(1);
2321                    let mut pos = base + k * size;
2322                    let end = (pos + size).min(range.end);
2323                    if walker
2324                        .record(file, &mut pos, end, Some(&mut out), None)
2325                        .is_err()
2326                    {
2327                        out.finish_row();
2328                    }
2329                }
2330            }
2331            Index::Walk(checkpoints) => {
2332                let at = checkpoints.partition_point(|c| c.row <= start as u64);
2333                let mut want = Want {
2334                    start: start as u64,
2335                    len,
2336                    produced: 0,
2337                };
2338                let mut out = Out {
2339                    sinks: &mut sinks,
2340                    filled: &mut filled,
2341                };
2342                if let Some(cp) = at.checked_sub(1).map(|i| &checkpoints[i]) {
2343                    self.read_from(cp, &mut walker, &mut want, &mut out);
2344                }
2345                // Rows the index counted that a read now does not find: a file changed
2346                // under the map, or a block that no longer decompresses.
2347                for _ in want.produced..len {
2348                    out.finish_row();
2349                }
2350            }
2351        }
2352        let columns = sinks
2353            .into_iter()
2354            .zip(&plan.columns)
2355            .zip(wanted)
2356            .filter(|(_, w)| **w)
2357            .map(|((sink, c), _)| sink.finish(c.name.clone()).map(Column::from))
2358            .collect::<PolarsResult<Vec<_>>>()?;
2359        DataFrame::new(len, columns)
2360    }
2361
2362    /// From checkpoint `cp`, skip to the row `want` starts at and read its rows.
2363    fn read_from(
2364        &self,
2365        cp: &Checkpoint,
2366        walker: &mut Walker<'_>,
2367        want: &mut Want,
2368        out: &mut Out<'_>,
2369    ) {
2370        let file = self.bytes.as_slice();
2371        let plan = &*self.plan;
2372        walker.acc.copy_from_slice(&cp.acc);
2373        let mut row = cp.row;
2374        let mut ci = cp.chunk as usize;
2375        let mut first = true;
2376        while want.produced < want.len && ci < self.chunks.len() {
2377            if !first {
2378                reset_block_sums(plan, &mut walker.acc);
2379            }
2380            let (data, mut cursor) = match self.enter(ci, walker) {
2381                Ok(c) => c,
2382                Err(_) => {
2383                    // A block that will not decompress now: its promised rows are null.
2384                    let n = self.chunks[ci].records.unwrap_or(0);
2385                    for _ in 0..n {
2386                        want.null_row(row, out);
2387                        row += 1;
2388                    }
2389                    ci += 1;
2390                    first = false;
2391                    continue;
2392                }
2393            };
2394            if first && cp.pos != u32::MAX {
2395                cursor.pos = cp.pos as usize;
2396                cursor.taken = u64::from(cp.taken);
2397            }
2398            first = false;
2399            let bytes = self.slice(&data, file);
2400            while want.produced < want.len {
2401                if cursor.limit.is_some_and(|l| cursor.taken >= l) || cursor.pos >= cursor.end {
2402                    break;
2403                }
2404                let reading = row >= want.start;
2405                let got = walker.record(
2406                    bytes,
2407                    &mut cursor.pos,
2408                    cursor.end,
2409                    if reading { Some(&mut *out) } else { None },
2410                    cursor.time,
2411                );
2412                match got {
2413                    Ok(Got::Row) => {
2414                        if reading {
2415                            want.produced += 1;
2416                        }
2417                        row += 1;
2418                        cursor.taken += 1;
2419                        align(&mut cursor, plan.align);
2420                    }
2421                    Ok(Got::Skipped) => {
2422                        cursor.taken += 1;
2423                        align(&mut cursor, plan.align);
2424                    }
2425                    Ok(Got::None) | Err(_) => break,
2426                }
2427            }
2428            // A chunk that promised more records than it held: nulls for the rest.
2429            if let Some(l) = self.chunks[ci].records {
2430                while cursor.taken < l && want.produced < want.len {
2431                    want.null_row(row, out);
2432                    row += 1;
2433                    cursor.taken += 1;
2434                }
2435            }
2436            ci += 1;
2437        }
2438    }
2439
2440    /// Column `column` of `rows`, each read where the row table says its record starts.
2441    fn decode_from(
2442        &self,
2443        table: &RowTable,
2444        column: usize,
2445        rows: &[IdxSize],
2446    ) -> PolarsResult<Column> {
2447        self.bytes.still_whole()?;
2448        let plan = &*self.plan;
2449        let file = self.bytes.as_slice();
2450        let ChunkSource::Map(range) = &self.chunks[0].source else {
2451            unreachable!("a row table is over the map")
2452        };
2453        let mut sinks: Vec<Sink> = vec![Sink::Skip; plan.columns.len()];
2454        sinks[column] = plan.columns[column].proto.clone();
2455        let mut filled = vec![false; sinks.len()];
2456        let mut out = Out {
2457            sinks: &mut sinks,
2458            filled: &mut filled,
2459        };
2460        let mut walker = Walker::new(plan, file);
2461        if plan.type_out == Some(column) {
2462            // Each row's variant is its tag: the names are cast once and the rows
2463            // gathered by tag. A row read field by field says its own label.
2464            let mut labels: Vec<Option<Arc<str>>> =
2465                plan.variants.iter().map(|v| Some(v.name.clone())).collect();
2466            let codes: IdxCa = rows
2467                .iter()
2468                .map(|&row| {
2469                    let row = row as usize;
2470                    let tag = table.tags[row];
2471                    if tag != WALK {
2472                        return Some(IdxSize::from(tag));
2473                    }
2474                    let mut pos = table.starts.get(row);
2475                    let _ = walker.record(file, &mut pos, range.end, Some(&mut out), None);
2476                    let Sink::Label(said) = &mut out.sinks[column] else {
2477                        return None;
2478                    };
2479                    let label = said.pop().flatten()?;
2480                    labels.push(Some(label));
2481                    Some((labels.len() - 1) as IdxSize)
2482                })
2483                .collect();
2484            let name = plan.columns[column].name.clone();
2485            let labels = Sink::Label(labels).finish(name)?;
2486            return Ok(labels.take(&codes)?.into_column());
2487        }
2488        let sources = &plan.sources[column];
2489        for &row in rows {
2490            let row = row as usize;
2491            let start = table.starts.get(row);
2492            let tag = table.tags[row];
2493            let source = match tag {
2494                WALK => &Source::Walk,
2495                v => &sources[usize::from(v)],
2496            };
2497            match source {
2498                Source::Null => out.finish_row(),
2499                Source::Label => {
2500                    if let Sink::Label(v) = &mut out.sinks[column] {
2501                        v.push(plan.variants.get(usize::from(tag)).map(|v| v.name.clone()));
2502                    }
2503                    out.filled[column] = true;
2504                    out.finish_row();
2505                }
2506                Source::At { offset, field } => {
2507                    let Kind::Fixed { width, int } = field.kind else {
2508                        unreachable!("a value at a place is fixed")
2509                    };
2510                    let cells = match field.count {
2511                        Some(SizeRef::Given(n)) => Some(n),
2512                        _ => None,
2513                    };
2514                    let at = start + offset;
2515                    let bytes = at
2516                        .checked_add(width * cells.unwrap_or(1))
2517                        .filter(|end| *end <= range.end)
2518                        .map(|end| &file[at..end]);
2519                    if let Some(bytes) = bytes {
2520                        let raw = int
2521                            .and_then(|r| (cells != Some(0)).then(|| int_of(&bytes[..width], r)));
2522                        push_fixed(field, bytes, width, cells, &mut out);
2523                        push_bits(field, raw, &mut out);
2524                    }
2525                    out.finish_row();
2526                }
2527                Source::Walk | Source::Summed => {
2528                    let mut pos = start;
2529                    // A row the walk does not finish (the file changed under the map)
2530                    // is null rather than missing.
2531                    if !matches!(
2532                        walker.record(file, &mut pos, range.end, Some(&mut out), None),
2533                        Ok(Got::Row)
2534                    ) {
2535                        out.finish_row();
2536                    }
2537                }
2538            }
2539        }
2540        let name = plan.columns[column].name.clone();
2541        Ok(sinks.swap_remove(column).finish(name)?.into_column())
2542    }
2543
2544    pub fn rows(&self) -> usize {
2545        self.rows
2546    }
2547
2548    pub fn schema(&self) -> SchemaRef {
2549        self.schema.clone()
2550    }
2551
2552    pub fn sources(&self) -> &[Arc<Bytes>] {
2553        std::slice::from_ref(&self.bytes)
2554    }
2555}
2556
2557impl FramedRecords {
2558    /// The frame: decoded over a row index, only what a query reaches.
2559    pub fn lazy(self: &Arc<Self>) -> LazyFrame {
2560        crate::formats::row_index::lazy(self)
2561    }
2562
2563    /// Rows `[start, start + len)` as a frame of their own, read from the nearest
2564    /// checkpoint before them.
2565    pub fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
2566        let all = vec![true; self.plan.columns.len()];
2567        Ok(self.decode(start, len, &all)?.lazy())
2568    }
2569
2570    /// The first `rows` rows of every column, decoded now.
2571    pub fn collect(&self, rows: usize) -> PolarsResult<DataFrame> {
2572        let all = vec![true; self.plan.columns.len()];
2573        self.decode(0, rows, &all)
2574    }
2575}
2576
2577impl crate::formats::row_index::RowSource for FramedRecords {
2578    fn height(&self) -> usize {
2579        self.rows
2580    }
2581
2582    fn schema(&self) -> SchemaRef {
2583        self.schema.clone()
2584    }
2585
2586    // Each column decodes alone, so a query decodes only what it names, a morsel at a time;
2587    // with a row table each row reads from its record start, without one each column walks
2588    // the span its rows cover.
2589    fn decode(&self, column: usize, index: &IdxCa) -> PolarsResult<Column> {
2590        let rows = crate::formats::row_index::checked(index, self.rows)?;
2591        if let Some(table) = &self.table
2592            && let Some(sources) = self.plan.sources.get(column)
2593            && !sources.iter().any(|s| matches!(s, Source::Summed))
2594        {
2595            return self.decode_from(table, column, &rows);
2596        }
2597        let mut wanted = vec![false; self.plan.columns.len()];
2598        *wanted
2599            .get_mut(column)
2600            .ok_or_else(|| polars_err!(OutOfBounds: "no column {column}"))? = true;
2601        let (Some(&lo), Some(&hi)) = (rows.iter().min(), rows.iter().max()) else {
2602            return Ok(self.decode(0, 0, &wanted)?.columns()[0].clone());
2603        };
2604        let (lo, span) = (lo as usize, (hi - lo) as usize + 1);
2605        let values = self.decode(lo, span, &wanted)?.columns()[0].clone();
2606        let contiguous = rows.len() == span && rows.windows(2).all(|w| w[1] == w[0] + 1);
2607        if contiguous {
2608            return Ok(values);
2609        }
2610        let at = IdxCa::from_vec(
2611            PlSmallStr::EMPTY,
2612            rows.iter().map(|&r| r - lo as IdxSize).collect(),
2613        );
2614        values.take(&at)
2615    }
2616}
2617
2618impl crate::formats::pushdown::Windowed for FramedRecords {
2619    fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
2620        FramedRecords::window(self, start, len)
2621    }
2622}
2623
2624/// Put `cursor` at the next multiple of `align` from the chunk's start.
2625fn align(cursor: &mut Cursor, align: usize) {
2626    if align > 1 {
2627        let offset = cursor.pos - cursor.data_start;
2628        let rounded = offset.div_ceil(align).saturating_mul(align);
2629        cursor.pos = cursor.data_start.saturating_add(rounded).min(cursor.end);
2630    }
2631}
2632
2633/// Running sums that start again with each block.
2634fn reset_block_sums(plan: &Plan, acc: &mut [i128]) {
2635    for (sum, delta) in acc.iter_mut().zip(&plan.deltas) {
2636        if *delta == Delta::Block {
2637            *sum = 0;
2638        }
2639    }
2640}
2641
2642/// A warning naming the bytes at the end that are not a whole record.
2643fn trailing_note(what: &str, bytes: &[u8]) -> String {
2644    formats::trailing_note(what, bytes)
2645}
2646
2647impl Plan {
2648    /// A plan of `fields` alone, to read a block header or an index entry.
2649    fn bare(fields: Vec<FieldPlan>, slots: usize) -> Self {
2650        Self {
2651            framing: Framing::Fixed,
2652            common: fields,
2653            type_slot: None,
2654            variants: Vec::new(),
2655            type_out: None,
2656            size: None,
2657            suffix: None,
2658            align: 1,
2659            sync: Vec::new(),
2660            checksum: None,
2661            chunk_header: Vec::new(),
2662            chunk_count: None,
2663            chunk_slots: 0,
2664            time_out: None,
2665            slots,
2666            deltas: Vec::new(),
2667            columns: Vec::new(),
2668            only: None,
2669            sources: Vec::new(),
2670        }
2671    }
2672}
2673
2674/// Where each column's value is in a record of each variant: at a fixed place while
2675/// every field before it in the record is fixed in size, else found by a walk.
2676fn sources(plan: &Plan) -> Vec<Vec<Source>> {
2677    let variants = plan.variants.len().max(1);
2678    let mut out = vec![vec![Source::Null; variants]; plan.columns.len()];
2679    // Read alone, the other variants' fields fill no column of this table.
2680    for v in (0..variants).filter(|v| plan.only.is_none_or(|only| only == *v)) {
2681        let fields = plan
2682            .common
2683            .iter()
2684            .chain(plan.variants.get(v).into_iter().flat_map(|p| &p.fields));
2685        // A row's start is its sync marker's.
2686        let mut offset = Some(plan.sync.len());
2687        for f in fields {
2688            let width = match (&f.kind, f.count) {
2689                (Kind::Fixed { width, .. }, None) => Some(*width),
2690                (Kind::Fixed { width, .. }, Some(SizeRef::Given(n))) => width.checked_mul(n),
2691                (
2692                    Kind::Pad {
2693                        size: SizeRef::Given(n),
2694                    },
2695                    None,
2696                ) => Some(*n),
2697                _ => None,
2698            };
2699            let source = match (offset, &f.kind, width) {
2700                (Some(offset), Kind::Fixed { .. }, Some(_)) if f.delta == Delta::None => {
2701                    Source::At {
2702                        offset,
2703                        field: f.clone(),
2704                    }
2705                }
2706                _ => Source::Walk,
2707            };
2708            for &c in &f.outs {
2709                out[c][v] = if f.delta == Delta::None {
2710                    source.clone()
2711                } else {
2712                    Source::Summed
2713                };
2714            }
2715            for (c, _, _) in &f.bits {
2716                out[*c][v] = source.clone();
2717            }
2718            offset = offset.zip(width).and_then(|(o, w)| o.checked_add(w));
2719        }
2720    }
2721    if let Some(c) = plan.type_out {
2722        out[c].fill(Source::Label);
2723    }
2724    for c in plan.checksum.iter().map(|c| c.out).chain(plan.time_out) {
2725        out[c].fill(Source::Walk);
2726    }
2727    out
2728}
2729
2730/// Fields read once, for their values: a block header, an index entry.
2731struct Struct {
2732    plan: Plan,
2733    scope: Scope,
2734}
2735
2736impl Struct {
2737    fn compile(
2738        spec: &Spec,
2739        header: &HeaderValues,
2740        file: &[u8],
2741        fields: &[Field],
2742    ) -> Result<Self, String> {
2743        let mut compiler = Compiler {
2744            spec,
2745            header,
2746            data: file,
2747            slots: 0,
2748            deltas: Vec::new(),
2749        };
2750        let mut scope = Scope::default();
2751        let mut columns = Vec::new();
2752        let plans = compiler.fields(fields, &mut scope, &mut columns, false)?;
2753        Ok(Self {
2754            plan: Plan::bare(plans, compiler.slots),
2755            scope,
2756        })
2757    }
2758
2759    /// Read the fields at `pos`, before `end`: the frame, and where they end.
2760    fn read(&self, file: &[u8], pos: usize, end: usize) -> Result<(Frame, usize), Stop> {
2761        let mut walker = Walker::new(&self.plan, file);
2762        walker.end = end;
2763        walker.record_start = pos;
2764        let mut p = pos;
2765        walker.walk(&self.plan.common, file, &mut p, None)?;
2766        Ok((walker.frame, p))
2767    }
2768
2769    fn value(&self, frame: &Frame, name: &str) -> Option<i128> {
2770        frame.ints[self.scope.slot(name)?]
2771    }
2772}
2773
2774/// The blocks of `data`: from the index the file keeps, or by walking their headers.
2775fn list_blocks(
2776    spec: &Spec,
2777    blocks: &formats::Blocks,
2778    header: &HeaderValues,
2779    file: &[u8],
2780    data: Range<usize>,
2781    notes: &mut Vec<String>,
2782) -> Result<Vec<Chunk>, String> {
2783    let head = Struct::compile(spec, header, file, &blocks.header)?;
2784    let size = {
2785        let compiler = Compiler {
2786            spec,
2787            header,
2788            data: file,
2789            slots: 0,
2790            deltas: Vec::new(),
2791        };
2792        compiler.size(&blocks.size, &head.scope, "block size")?
2793    };
2794    let codec_of = |frame: &Frame| -> Result<Compression, String> {
2795        match &blocks.codec {
2796            Codec::Fixed(c) => Ok(*c),
2797            Codec::ByField { field, values } => {
2798                let code = head
2799                    .value(frame, field)
2800                    .ok_or("the block header has no codec")?;
2801                i64::try_from(code)
2802                    .ok()
2803                    .and_then(|c| values.get(&c).copied())
2804                    .ok_or_else(|| format!("codec {code} is not one compression names"))
2805            }
2806        }
2807    };
2808    let mut chunks = Vec::new();
2809    let read_block = |at: usize,
2810                      end: usize,
2811                      rows: Option<u64>,
2812                      chunks: &mut Vec<Chunk>,
2813                      notes: &mut Vec<String>|
2814     -> Result<Option<usize>, String> {
2815        let (frame, body_start) = match head.read(file, at, end) {
2816            Ok(x) => x,
2817            Err(_) => {
2818                notes.push(format!(
2819                    "the block header at byte {at} runs past the data; the rest is left out"
2820                ));
2821                return Ok(None);
2822            }
2823        };
2824        let body_len = match size {
2825            SizeRef::Given(n) => n,
2826            SizeRef::Slot { slot, adjust } => {
2827                let v = frame.ints[slot].unwrap_or(0) + i128::from(adjust);
2828                usize::try_from(v)
2829                    .ok()
2830                    .filter(|n| *n <= MAX_BLOCK)
2831                    .ok_or_else(|| {
2832                        format!(
2833                            "the block at byte {at} gives its size as {v}, outside 0 to {MAX_BLOCK}"
2834                        )
2835                    })?
2836            }
2837            SizeRef::Rest => end - body_start,
2838        };
2839        let body_end = body_start.checked_add(body_len).filter(|e| *e <= end);
2840        let Some(body_end) = body_end else {
2841            notes.push(format!(
2842                "the block at byte {at} is {body_len} bytes and runs past the data; the rest is left out"
2843            ));
2844            return Ok(None);
2845        };
2846        let codec = match codec_of(&frame) {
2847            Ok(c) => c,
2848            Err(e) => {
2849                notes.push(format!("the block at byte {at}: {e}; left out"));
2850                return Ok(Some(body_end));
2851            }
2852        };
2853        let records = rows.or_else(|| {
2854            blocks
2855                .records
2856                .as_ref()
2857                .and_then(|r| head.value(&frame, r))
2858                .map(|v| v.max(0) as u64)
2859        });
2860        let uncompressed = blocks
2861            .uncompressed
2862            .as_ref()
2863            .and_then(|u| head.value(&frame, u))
2864            .and_then(|v| usize::try_from(v).ok());
2865        chunks.push(Chunk {
2866            source: if codec == Compression::None {
2867                ChunkSource::Map(body_start..body_end)
2868            } else {
2869                ChunkSource::Block {
2870                    body: body_start..body_end,
2871                    codec,
2872                    uncompressed,
2873                }
2874            },
2875            records,
2876            time_ns: None,
2877        });
2878        Ok(Some(body_end))
2879    };
2880    match &blocks.index {
2881        Some(index) => {
2882            let entry = Struct::compile(spec, header, file, &index.fields)?;
2883            let at = usize::try_from(header.resolve_any(&index.at, "index at")?)
2884                .map_err(|_| "index: too far")?;
2885            let count = header.resolve_any(&index.count, "index count")?;
2886            let width = formats::fields_width(&index.fields).unwrap_or(1).max(1) as usize;
2887            let fits = file.len().saturating_sub(at) / width;
2888            if count > fits as u64 {
2889                return Err(format!(
2890                    "the block index at byte {at} says {count} entries; the file has room for {fits}"
2891                ));
2892            }
2893            let mut pos = at;
2894            for _ in 0..count {
2895                let (frame, next) = head_or(entry.read(file, pos, file.len()))?;
2896                pos = next;
2897                let offset = entry.value(&frame, "offset").unwrap_or(-1);
2898                let rows = entry.value(&frame, "rows").map(|v| v.max(0) as u64);
2899                let Some(offset) = usize::try_from(offset).ok().filter(|o| *o < file.len()) else {
2900                    notes.push(format!(
2901                        "an index entry points at byte {offset}, outside the file; left out"
2902                    ));
2903                    continue;
2904                };
2905                read_block(offset, file.len(), rows, &mut chunks, notes)?;
2906            }
2907        }
2908        None => {
2909            let mut pos = data.start;
2910            while pos < data.end {
2911                match read_block(pos, data.end, None, &mut chunks, notes)? {
2912                    Some(next) if next > pos => pos = next,
2913                    Some(_) => {
2914                        notes.push(format!(
2915                            "zero-length block at byte {pos} {} rest left out",
2916                            crate::glyphs::get().middot
2917                        ));
2918                        break;
2919                    }
2920                    None => break,
2921                }
2922            }
2923        }
2924    }
2925    Ok(chunks)
2926}
2927
2928fn head_or(read: Result<(Frame, usize), Stop>) -> Result<(Frame, usize), String> {
2929    read.map_err(|_| "an index entry runs past the end of the file".to_string())
2930}
2931
2932/// `raw` decompressed by `codec`, at most [`MAX_BLOCK`] bytes.
2933pub fn decompress(
2934    raw: &[u8],
2935    codec: Compression,
2936    uncompressed: Option<usize>,
2937) -> Result<Vec<u8>, String> {
2938    use std::io::Read;
2939    let read_all = |mut reader: Box<dyn Read + '_>| -> Result<Vec<u8>, String> {
2940        let mut out = Vec::new();
2941        reader
2942            .by_ref()
2943            .take(MAX_BLOCK as u64 + 1)
2944            .read_to_end(&mut out)
2945            .map_err(|e| format!("{} block: {e}", codec.name()))?;
2946        if out.len() > MAX_BLOCK {
2947            return Err(format!(
2948                "a block decompresses to more than {MAX_BLOCK} bytes"
2949            ));
2950        }
2951        Ok(out)
2952    };
2953    match codec {
2954        Compression::None => Ok(raw.to_vec()),
2955        Compression::Gzip => read_all(Box::new(flate2::read::MultiGzDecoder::new(raw))),
2956        Compression::Deflate => read_all(Box::new(flate2::read::DeflateDecoder::new(raw))),
2957        Compression::Zlib => read_all(Box::new(flate2::read::ZlibDecoder::new(raw))),
2958        Compression::Zstd => read_all(Box::new(
2959            zstd::Decoder::new(raw).map_err(|e| format!("zstd block: {e}"))?,
2960        )),
2961        Compression::Lz4 => read_all(Box::new(
2962            lz4::Decoder::new(raw).map_err(|e| format!("lz4 block: {e}"))?,
2963        )),
2964        Compression::Lz4Block => {
2965            let size = uncompressed.filter(|n| *n <= MAX_BLOCK).ok_or_else(|| {
2966                format!("an lz4 block needs its decompressed size, at most {MAX_BLOCK}")
2967            })?;
2968            lz4::block::decompress(raw, Some(size as i32)).map_err(|e| format!("lz4 block: {e}"))
2969        }
2970        Compression::Snappy => {
2971            let size = snap::raw::decompress_len(raw).map_err(|e| format!("snappy block: {e}"))?;
2972            if size > MAX_BLOCK {
2973                return Err(format!(
2974                    "a block decompresses to more than {MAX_BLOCK} bytes"
2975                ));
2976            }
2977            snap::raw::Decoder::new()
2978                .decompress_vec(raw)
2979                .map_err(|e| format!("snappy block: {e}"))
2980        }
2981        Compression::SnappyFramed => read_all(Box::new(snap::read::FrameDecoder::new(raw))),
2982        Compression::Brotli => read_all(Box::new(brotli::Decompressor::new(raw, 4096))),
2983        Compression::Bzip2 => read_all(Box::new(bzip2::read::BzDecoder::new(raw))),
2984        Compression::Xz => read_all(Box::new(xz2::read::XzDecoder::new(raw))),
2985    }
2986}
2987
2988/// Packet captures (pcap and pcapng): the UDP payload of each packet, with its time.
2989pub mod capture {
2990    use super::{Chunk, ChunkSource};
2991
2992    /// The most packets one capture is read for.
2993    const MAX_PACKETS: usize = 64 << 20;
2994
2995    fn u16_at(b: &[u8], at: usize, big: bool) -> Option<u16> {
2996        let raw: [u8; 2] = b.get(at..at + 2)?.try_into().ok()?;
2997        Some(if big {
2998            u16::from_be_bytes(raw)
2999        } else {
3000            u16::from_le_bytes(raw)
3001        })
3002    }
3003
3004    fn u32_at(b: &[u8], at: usize, big: bool) -> Option<u32> {
3005        let raw: [u8; 4] = b.get(at..at + 4)?.try_into().ok()?;
3006        Some(if big {
3007            u32::from_be_bytes(raw)
3008        } else {
3009            u32::from_le_bytes(raw)
3010        })
3011    }
3012
3013    /// Whether `head` starts a pcap or pcapng file.
3014    pub fn is_capture(head: &[u8]) -> bool {
3015        matches!(
3016            head.get(..4),
3017            Some(
3018                [0xd4, 0xc3, 0xb2, 0xa1]
3019                    | [0xa1, 0xb2, 0xc3, 0xd4]
3020                    | [0x4d, 0x3c, 0xb2, 0xa1]
3021                    | [0xa1, 0xb2, 0x3c, 0x4d]
3022                    | [0x0a, 0x0d, 0x0d, 0x0a]
3023            )
3024        )
3025    }
3026
3027    /// Where the UDP payload of a frame of `link` type is in `frame`, relative to it.
3028    pub fn udp_payload(frame: &[u8], link: u32) -> Option<std::ops::Range<usize>> {
3029        // The link layer: where the IP packet starts and which protocol it is.
3030        let (mut at, mut ethertype) = match link {
3031            // Ethernet.
3032            1 => (14, u16_at(frame, 12, true)?),
3033            // Raw IP.
3034            101 | 12 | 14 => (0, 0),
3035            228 => (0, 0x0800),
3036            229 => (0, 0x86dd),
3037            // Linux cooked captures.
3038            113 => (16, u16_at(frame, 14, true)?),
3039            276 => (20, u16_at(frame, 0, true)?),
3040            // BSD loopback: a four-byte address family in host order.
3041            0 | 108 => {
3042                let family = u32_at(frame, 0, false)?;
3043                let family = if family > 0xffff {
3044                    family.swap_bytes()
3045                } else {
3046                    family
3047                };
3048                (4, if family == 2 { 0x0800 } else { 0x86dd })
3049            }
3050            _ => return None,
3051        };
3052        // VLAN tags.
3053        while ethertype == 0x8100 || ethertype == 0x88a8 {
3054            ethertype = u16_at(frame, at + 2, true)?;
3055            at += 4;
3056        }
3057        if ethertype == 0 {
3058            ethertype = match frame.get(at)? >> 4 {
3059                4 => 0x0800,
3060                6 => 0x86dd,
3061                _ => return None,
3062            };
3063        }
3064        let (udp, ip_end) = match ethertype {
3065            0x0800 => {
3066                let ihl = usize::from(frame.get(at)? & 0x0f) * 4;
3067                let total = usize::from(u16_at(frame, at + 2, true)?);
3068                let flags = u16_at(frame, at + 6, true)?;
3069                // Fragments other than a whole datagram are not read.
3070                if ihl < 20 || *frame.get(at + 9)? != 17 || flags & 0x3fff != 0 {
3071                    return None;
3072                }
3073                (at + ihl, (at + total).min(frame.len()))
3074            }
3075            0x86dd => {
3076                if *frame.get(at + 6)? != 17 {
3077                    return None;
3078                }
3079                let payload = usize::from(u16_at(frame, at + 4, true)?);
3080                (at + 40, (at + 40 + payload).min(frame.len()))
3081            }
3082            _ => return None,
3083        };
3084        let len = usize::from(u16_at(frame, udp + 4, true)?);
3085        let start = udp + 8;
3086        let end = (udp + len).min(ip_end);
3087        (len >= 8 && start <= end).then_some(start..end)
3088    }
3089
3090    /// The UDP payloads of the capture in `file`, each with its time.
3091    pub(super) fn packets(file: &[u8], notes: &mut Vec<String>) -> Result<Vec<Chunk>, String> {
3092        let mut out = Vec::new();
3093        let mut other = 0u64;
3094        let magic = file
3095            .get(..4)
3096            .ok_or("the capture is shorter than its header")?;
3097        if magic == [0x0a, 0x0d, 0x0d, 0x0a] {
3098            pcapng(file, &mut out, &mut other)?;
3099        } else {
3100            pcap(file, &mut out, &mut other)?;
3101        }
3102        if other > 0 {
3103            notes.push(format!(
3104                "{other} {} in the capture {} not UDP, left out",
3105                if other == 1 { "packet" } else { "packets" },
3106                if other == 1 { "is" } else { "are" }
3107            ));
3108        }
3109        Ok(out)
3110    }
3111
3112    fn push(
3113        out: &mut Vec<Chunk>,
3114        base: usize,
3115        payload: std::ops::Range<usize>,
3116        time_ns: Option<i64>,
3117    ) {
3118        out.push(Chunk {
3119            source: ChunkSource::Map(base + payload.start..base + payload.end),
3120            records: None,
3121            time_ns,
3122        });
3123    }
3124
3125    fn pcap(file: &[u8], out: &mut Vec<Chunk>, other: &mut u64) -> Result<(), String> {
3126        let (big, nanos) = match file.get(..4) {
3127            Some([0xd4, 0xc3, 0xb2, 0xa1]) => (false, false),
3128            Some([0xa1, 0xb2, 0xc3, 0xd4]) => (true, false),
3129            Some([0x4d, 0x3c, 0xb2, 0xa1]) => (false, true),
3130            Some([0xa1, 0xb2, 0x3c, 0x4d]) => (true, true),
3131            _ => return Err("not a pcap or pcapng capture: its magic is not one".into()),
3132        };
3133        let link =
3134            u32_at(file, 20, big).ok_or("the capture is shorter than its header")? & 0x0fff_ffff;
3135        let mut at = 24usize;
3136        while at + 16 <= file.len() && out.len() < MAX_PACKETS {
3137            let secs = i64::from(u32_at(file, at, big).unwrap_or(0));
3138            let frac = i64::from(u32_at(file, at + 4, big).unwrap_or(0));
3139            let caplen = u32_at(file, at + 8, big).unwrap_or(0) as usize;
3140            let start = at + 16;
3141            let Some(end) = start.checked_add(caplen).filter(|e| *e <= file.len()) else {
3142                break;
3143            };
3144            let time = secs
3145                .checked_mul(1_000_000_000)
3146                .and_then(|s| s.checked_add(if nanos { frac } else { frac * 1000 }));
3147            match udp_payload(&file[start..end], link) {
3148                Some(payload) => push(out, start, payload, time),
3149                None => *other += 1,
3150            }
3151            at = end;
3152        }
3153        Ok(())
3154    }
3155
3156    fn pcapng(file: &[u8], out: &mut Vec<Chunk>, other: &mut u64) -> Result<(), String> {
3157        let mut at = 0usize;
3158        let mut big = false;
3159        // Each interface's link type and ticks per second.
3160        let mut interfaces: Vec<(u32, u64)> = Vec::new();
3161        while at + 12 <= file.len() && out.len() < MAX_PACKETS {
3162            let kind = u32_at(file, at, big).unwrap_or(0);
3163            if kind == 0x0a0d_0d0a {
3164                big = match file.get(at + 8..at + 12) {
3165                    Some([0x1a, 0x2b, 0x3c, 0x4d]) => true,
3166                    Some([0x4d, 0x3c, 0x2b, 0x1a]) => false,
3167                    _ => return Err("a pcapng section header has no byte-order magic".into()),
3168                };
3169                interfaces.clear();
3170            }
3171            let len = u32_at(file, at + 4, big).unwrap_or(0) as usize;
3172            if len < 12 || !len.is_multiple_of(4) || at + len > file.len() {
3173                break;
3174            }
3175            let body = &file[at + 8..at + len - 4];
3176            match kind {
3177                1 => {
3178                    let link = u32::from(u16_at(body, 0, big).unwrap_or(0));
3179                    let mut ticks = 1_000_000u64;
3180                    // Options: if_tsresol (9) sets the resolution.
3181                    let mut o = 8;
3182                    while o + 4 <= body.len() {
3183                        let code = u16_at(body, o, big).unwrap_or(0);
3184                        let olen = usize::from(u16_at(body, o + 2, big).unwrap_or(0));
3185                        if code == 0 {
3186                            break;
3187                        }
3188                        if code == 9 && olen >= 1 {
3189                            let r = body[o + 4];
3190                            ticks = if r & 0x80 != 0 {
3191                                1u64.checked_shl(u32::from(r & 0x7f)).unwrap_or(1_000_000)
3192                            } else {
3193                                10u64.checked_pow(u32::from(r)).unwrap_or(1_000_000)
3194                            };
3195                        }
3196                        o += 4 + olen.div_ceil(4) * 4;
3197                    }
3198                    interfaces.push((link, ticks.max(1)));
3199                }
3200                6 => {
3201                    let iface = u32_at(body, 0, big).unwrap_or(0) as usize;
3202                    let high = u64::from(u32_at(body, 4, big).unwrap_or(0));
3203                    let low = u64::from(u32_at(body, 8, big).unwrap_or(0));
3204                    let caplen = u32_at(body, 12, big).unwrap_or(0) as usize;
3205                    let Some(frame) = body.get(20..20 + caplen) else {
3206                        at += len;
3207                        continue;
3208                    };
3209                    let (link, ticks) = interfaces.get(iface).copied().unwrap_or((1, 1_000_000));
3210                    let stamp = (high << 32) | low;
3211                    let time =
3212                        i64::try_from(u128::from(stamp) * 1_000_000_000 / u128::from(ticks)).ok();
3213                    match udp_payload(frame, link) {
3214                        Some(payload) => push(out, at + 8 + 20, payload, time),
3215                        None => *other += 1,
3216                    }
3217                }
3218                3 => {
3219                    let (link, _) = interfaces.first().copied().unwrap_or((1, 1_000_000));
3220                    let frame = &body[4.min(body.len())..];
3221                    match udp_payload(frame, link) {
3222                        Some(payload) => push(out, at + 8 + 4, payload, None),
3223                        None => *other += 1,
3224                    }
3225                }
3226                _ => {}
3227            }
3228            at += len;
3229        }
3230        Ok(())
3231    }
3232}
3233
3234#[cfg(test)]
3235mod tests;