1use 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
21const CHECKPOINT: u64 = 1024;
23
24pub const MAX_BLOCK: usize = 256 << 20;
26
27const CACHED_BLOCKS: usize = 8;
29
30const MAX_ITEMS: u64 = 1 << 20;
32
33const MAX_ROWS: usize = IdxSize::MAX as usize;
35
36#[derive(Debug, Clone, Copy)]
38enum SizeRef {
39 Given(usize),
40 Slot { slot: usize, adjust: i64 },
41 Rest,
42}
43
44#[derive(Debug, Clone, Copy)]
46struct IntRead {
47 signed: bool,
48 big: bool,
49}
50
51#[derive(Debug, Clone)]
53struct FieldPlan {
54 name: String,
55 slot: usize,
56 kind: Kind,
57 count: Option<SizeRef>,
59 outs: Vec<usize>,
61 bits: Vec<(usize, u32, u32)>,
63 delta: Delta,
64 delta_index: usize,
66 sentinel: Option<i128>,
68}
69
70#[derive(Debug, Clone)]
71enum Kind {
72 Fixed {
75 width: usize,
76 int: Option<IntRead>,
77 },
78 Sized {
80 size: SizeRef,
81 encoding: Option<Encoding>,
82 },
83 Strz {
85 max: Option<SizeRef>,
86 encoding: Encoding,
87 },
88 Var {
89 signed: bool,
90 },
91 Pad {
92 size: SizeRef,
93 },
94 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#[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#[derive(Debug, Clone)]
118enum Sink {
119 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 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 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 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 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
302fn 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#[derive(Debug, Clone)]
326struct OutColumn {
327 name: PlSmallStr,
328 dtype: DataType,
329 proto: Sink,
330}
331
332#[derive(Debug, Clone)]
334enum ChunkSource {
335 Map(Range<usize>),
337 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: Option<u64>,
350 time_ns: Option<i64>,
352}
353
354#[derive(Debug, Clone)]
356struct Checkpoint {
357 row: u64,
358 chunk: u32,
359 pos: u32,
360 taken: u32,
362 acc: Box<[i128]>,
363}
364
365#[derive(Debug)]
367enum Index {
368 Stride {
370 start: usize,
371 size: usize,
372 ring: usize,
373 },
374 Walk(Vec<Checkpoint>),
375}
376
377const WALK: u8 = u8::MAX;
379
380#[derive(Debug)]
384struct RowTable {
385 starts: crate::formats::indexed::Offsets,
387 tags: Vec<u8>,
389}
390
391struct 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#[derive(Debug, Clone)]
404enum Source {
405 Null,
407 At { offset: usize, field: FieldPlan },
409 Label,
411 Walk,
413 Summed,
415}
416
417#[derive(Debug)]
419struct Plan {
420 framing: Framing,
421 common: Vec<FieldPlan>,
422 type_slot: Option<(usize, bool)>,
424 variants: Vec<VariantPlan>,
425 type_out: Option<usize>,
426 size: Option<SizeRef>,
428 suffix: Option<(usize, bool)>,
430 align: usize,
431 sync: Vec<u8>,
432 checksum: Option<ChecksumPlan>,
433 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 only: Option<usize>,
443 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
456struct 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 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
482enum Stop {
484 Truncated,
486 Said(String),
488}
489
490pub struct FramedRecords {
492 bytes: Arc<Bytes>,
493 plan: Arc<Plan>,
494 chunks: Vec<Chunk>,
495 index: Arc<Index>,
496 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
512pub 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
532struct Compiler<'a> {
534 spec: &'a Spec,
535 header: &'a HeaderValues,
536 data: &'a [u8],
537 slots: usize,
538 deltas: Vec<Delta>,
539}
540
541#[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 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 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 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 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(§ion.offset, "section offset")?;
883 let size = self.header.resolve_any(§ion.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
895fn 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
909fn 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
920fn 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
936fn 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
946fn 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
960enum Got {
961 Row,
963 Skipped,
965 None,
967}
968
969struct Out<'s> {
971 sinks: &'s mut [Sink],
972 filled: &'s mut [bool],
973}
974
975impl Out<'_> {
976 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
987struct Walker<'a> {
989 plan: &'a Plan,
990 file: &'a [u8],
992 frame: Frame,
993 acc: Vec<i128>,
994 end: usize,
996 bounded: bool,
998 short: bool,
1000 record_start: usize,
1001 size_slot: Option<(usize, i64)>,
1003 skipped: u64,
1005 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 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 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 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 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 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 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 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
1552fn 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
1570fn 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
1609enum ChunkData {
1611 Map(Range<usize>),
1612 Owned(Arc<Vec<u8>>),
1613}
1614
1615struct Cursor {
1617 pos: usize,
1618 data_start: usize,
1619 end: usize,
1620 taken: u64,
1621 limit: Option<u64>,
1622 time: Option<i64>,
1623}
1624
1625struct Want {
1627 start: u64,
1628 len: usize,
1629 produced: usize,
1630}
1631
1632impl Want {
1633 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#[derive(Default)]
1644struct Found {
1645 notes: Vec<String>,
1646 skipped: u64,
1647}
1648
1649impl FramedRecords {
1650 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 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 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 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 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 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 && 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 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 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 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 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 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 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 cursor.pos = end;
2284 cursor.limit = Some(0);
2285 }
2286 }
2287 }
2288 Ok((data, cursor))
2289 }
2290
2291 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 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 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 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 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 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 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 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 pub fn lazy(self: &Arc<Self>) -> LazyFrame {
2560 crate::formats::row_index::lazy(self)
2561 }
2562
2563 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 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 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
2624fn 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
2633fn 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
2642fn trailing_note(what: &str, bytes: &[u8]) -> String {
2644 formats::trailing_note(what, bytes)
2645}
2646
2647impl Plan {
2648 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
2674fn 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 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 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
2730struct 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 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
2774fn 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
2932pub 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
2988pub mod capture {
2990 use super::{Chunk, ChunkSource};
2991
2992 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 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 pub fn udp_payload(frame: &[u8], link: u32) -> Option<std::ops::Range<usize>> {
3029 let (mut at, mut ethertype) = match link {
3031 1 => (14, u16_at(frame, 12, true)?),
3033 101 | 12 | 14 => (0, 0),
3035 228 => (0, 0x0800),
3036 229 => (0, 0x86dd),
3037 113 => (16, u16_at(frame, 14, true)?),
3039 276 => (20, u16_at(frame, 0, true)?),
3040 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 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 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 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 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 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;