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