use crate::fixed_records::{Bytes, ColumnLayout, Physical};
use crate::formats::{
self, Amount, Codec, Compression, Delta, Encoding, Field, Framing, HeaderValues, Spec, Type,
};
use polars::prelude::*;
use std::collections::VecDeque;
use std::ops::Range;
use std::sync::{Arc, Mutex};
const CHECKPOINT: u64 = 1024;
pub const MAX_BLOCK: usize = 256 << 20;
const CACHED_BLOCKS: usize = 8;
const MAX_ITEMS: u64 = 1 << 20;
const MAX_ROWS: usize = IdxSize::MAX as usize;
#[derive(Debug, Clone, Copy)]
enum SizeRef {
Given(usize),
Slot { slot: usize, adjust: i64 },
Rest,
}
#[derive(Debug, Clone, Copy)]
struct IntRead {
signed: bool,
big: bool,
}
#[derive(Debug, Clone)]
struct FieldPlan {
name: String,
slot: usize,
kind: Kind,
count: Option<SizeRef>,
outs: Vec<usize>,
bits: Vec<(usize, u32, u32)>,
delta: Delta,
delta_index: usize,
sentinel: Option<i128>,
}
#[derive(Debug, Clone)]
enum Kind {
Fixed {
width: usize,
int: Option<IntRead>,
},
Sized {
size: SizeRef,
encoding: Option<Encoding>,
},
Strz {
max: Option<SizeRef>,
encoding: Encoding,
},
Var {
signed: bool,
},
Pad {
size: SizeRef,
},
StringAt {
width: usize,
big: bool,
section: Range<usize>,
},
Group {
items: Vec<FieldPlan>,
slots: usize,
},
}
#[derive(Debug, Clone)]
struct VariantPlan {
name: Arc<str>,
ints: Vec<i128>,
texts: Vec<String>,
fields: Vec<FieldPlan>,
size: Option<SizeRef>,
}
#[derive(Debug, Clone)]
enum Sink {
Packed {
layout: ColumnLayout,
cell: usize,
buf: Vec<u8>,
valid: Vec<bool>,
},
Text(Vec<Option<String>>),
Binary(Vec<Option<Vec<u8>>>),
Bits {
width: u32,
labels: Option<Arc<std::collections::BTreeMap<i64, String>>>,
values: Vec<Option<u64>>,
},
Flag(Vec<Option<bool>>),
Label(Vec<Option<Arc<str>>>),
Time(Vec<Option<i64>>),
List {
inner: Box<Sink>,
offsets: Vec<i64>,
valid: Vec<bool>,
},
Struct {
names: Vec<PlSmallStr>,
fields: Vec<Sink>,
len: usize,
},
Skip,
}
impl Sink {
fn push_null(&mut self) {
match self {
Self::Packed {
cell, buf, valid, ..
} => {
buf.resize(buf.len() + *cell, 0);
valid.push(false);
}
Self::Text(v) => v.push(None),
Self::Label(v) => v.push(None),
Self::Binary(v) => v.push(None),
Self::Bits { values, .. } => values.push(None),
Self::Flag(v) => v.push(None),
Self::Time(v) => v.push(None),
Self::List { offsets, valid, .. } => {
offsets.push(*offsets.last().unwrap_or(&0));
valid.push(false);
}
Self::Struct { fields, len, .. } => {
for f in fields {
f.push_null();
}
*len += 1;
}
Self::Skip => {}
}
}
fn push_bytes(&mut self, bytes: &[u8]) {
if let Self::Packed { buf, valid, .. } = self {
buf.extend_from_slice(bytes);
valid.push(true);
}
}
fn finish(self, name: PlSmallStr) -> PolarsResult<Series> {
Ok(match self {
Self::Packed {
layout,
cell,
buf,
valid,
} => {
let rows = valid.len();
let layout = ColumnLayout {
start: 0,
stride: cell,
..layout
};
let column = crate::fixed_records::decode(&buf, &layout, rows)?;
let series = column
.as_materialized_series()
.clone()
.with_name(name.clone());
if valid.iter().all(|v| *v) {
series
} else {
let mask: BooleanChunked = valid.into_iter().collect();
let nulls = Series::full_null(PlSmallStr::EMPTY, rows, series.dtype());
series.zip_with(&mask, &nulls)?.with_name(name)
}
}
Self::Text(v) => StringChunked::from_iter_options(name, v.into_iter()).into_series(),
Self::Label(v) => {
let mut distinct: Vec<Arc<str>> = Vec::new();
let mut by_text = std::collections::HashMap::new();
let codes: IdxCa = v
.iter()
.map(|label| {
let label = label.as_ref()?;
if let Some(i) =
distinct.iter().take(16).position(|d| Arc::ptr_eq(d, label))
{
return Some(i as IdxSize);
}
Some(*by_text.entry(label.clone()).or_insert_with(|| {
distinct.push(label.clone());
(distinct.len() - 1) as IdxSize
}))
})
.collect();
StringChunked::from_iter_values(name, distinct.iter().map(|d| &**d))
.into_series()
.cast(&DataType::from_categories(Categories::global()))?
.take(&codes)?
}
Self::Binary(v) => BinaryChunked::from_iter_options(name, v.into_iter()).into_series(),
Self::Bits {
width,
labels,
values,
} => match (labels, width) {
(Some(labels), _) => StringChunked::from_iter_options(
name,
values.into_iter().map(|v| {
v.map(|v| {
i64::try_from(v)
.ok()
.and_then(|k| labels.get(&k).cloned())
.unwrap_or_else(|| v.to_string())
})
}),
)
.into_series(),
(None, 1) => BooleanChunked::from_iter_options(
name,
values.into_iter().map(|v| v.map(|v| v != 0)),
)
.into_series(),
(None, w) => {
let wide =
UInt64Chunked::from_iter_options(name, values.into_iter()).into_series();
let dtype = match w {
2..=8 => DataType::UInt8,
9..=16 => DataType::UInt16,
17..=32 => DataType::UInt32,
_ => DataType::UInt64,
};
wide.strict_cast(&dtype)?
}
},
Self::Flag(v) => BooleanChunked::from_iter_options(name, v.into_iter()).into_series(),
Self::Time(v) => Int64Chunked::from_iter_options(name, v.into_iter())
.into_datetime(TimeUnit::Nanoseconds, None)
.into_series(),
Self::List {
inner,
offsets,
valid,
} => {
let values = inner.finish(PlSmallStr::from_static("item"))?.rechunk();
list_series(name, values, offsets, valid)?
}
Self::Struct { names, fields, len } => {
let series = names
.into_iter()
.zip(fields)
.map(|(n, f)| f.finish(n))
.collect::<PolarsResult<Vec<_>>>()?;
StructChunked::from_series(name, len, series.iter())?.into_series()
}
Self::Skip => Series::new_empty(name, &DataType::Null),
})
}
}
fn list_series(
name: PlSmallStr,
values: Series,
offsets: Vec<i64>,
valid: Vec<bool>,
) -> PolarsResult<Series> {
use polars_arrow::array::ListArray;
use polars_arrow::bitmap::Bitmap;
use polars_arrow::offset::OffsetsBuffer;
let inner_dtype = values.dtype().clone();
let array = values.to_arrow(0, CompatLevel::newest());
let mut all = Vec::with_capacity(offsets.len() + 1);
all.push(0i64);
all.extend(offsets);
let offsets = OffsetsBuffer::<i64>::try_from(all)?;
let validity = (!valid.iter().all(|v| *v)).then(|| Bitmap::from_iter(valid));
let dtype = ListArray::<i64>::default_datatype(array.dtype().clone());
let list = ListArray::<i64>::try_new(dtype, offsets, array, validity)?;
Series::from_arrow(name, Box::new(list))?.cast(&DataType::List(Box::new(inner_dtype)))
}
#[derive(Debug, Clone)]
struct OutColumn {
name: PlSmallStr,
dtype: DataType,
proto: Sink,
}
#[derive(Debug, Clone)]
enum ChunkSource {
Map(Range<usize>),
Block {
body: Range<usize>,
codec: Compression,
uncompressed: Option<usize>,
},
}
#[derive(Debug, Clone)]
struct Chunk {
source: ChunkSource,
records: Option<u64>,
time_ns: Option<i64>,
}
#[derive(Debug, Clone)]
struct Checkpoint {
row: u64,
chunk: u32,
pos: u32,
taken: u32,
acc: Box<[i128]>,
}
#[derive(Debug)]
enum Index {
Stride {
start: usize,
size: usize,
ring: usize,
},
Walk(Vec<Checkpoint>),
}
const WALK: u8 = u8::MAX;
#[derive(Debug)]
struct RowTable {
starts: crate::indexed::Offsets,
tags: Vec<u8>,
}
struct KeptWalk {
spec: Spec,
data: Range<usize>,
index: Arc<Index>,
table: Option<Arc<RowTable>>,
rows: usize,
notes: Vec<String>,
}
#[derive(Debug, Clone)]
enum Source {
Null,
At { offset: usize, field: FieldPlan },
Label,
Walk,
Summed,
}
#[derive(Debug)]
struct Plan {
framing: Framing,
common: Vec<FieldPlan>,
type_slot: Option<(usize, bool)>,
variants: Vec<VariantPlan>,
type_out: Option<usize>,
size: Option<SizeRef>,
suffix: Option<(usize, bool)>,
align: usize,
sync: Vec<u8>,
checksum: Option<ChecksumPlan>,
chunk_header: Vec<FieldPlan>,
chunk_count: Option<usize>,
chunk_slots: usize,
time_out: Option<usize>,
slots: usize,
deltas: Vec<Delta>,
columns: Vec<OutColumn>,
only: Option<usize>,
sources: Vec<Vec<Source>>,
}
#[derive(Debug, Clone)]
struct ChecksumPlan {
algo: formats::ChecksumAlgo,
field: usize,
from: Option<usize>,
to: usize,
out: usize,
}
struct Frame {
starts: Vec<Option<usize>>,
ends: Vec<Option<usize>>,
ints: Vec<Option<i128>>,
}
impl Frame {
fn new(slots: usize) -> Self {
Self {
starts: vec![None; slots],
ends: vec![None; slots],
ints: vec![None; slots],
}
}
fn clear(&mut self, range: Range<usize>) {
for s in range {
self.starts[s] = None;
self.ends[s] = None;
self.ints[s] = None;
}
}
}
enum Stop {
Truncated,
Said(String),
}
pub struct FramedRecords {
bytes: Arc<Bytes>,
plan: Arc<Plan>,
chunks: Vec<Chunk>,
index: Arc<Index>,
table: Option<Arc<RowTable>>,
rows: usize,
schema: SchemaRef,
cache: Mutex<VecDeque<(usize, Arc<Vec<u8>>)>>,
}
impl std::fmt::Debug for FramedRecords {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("FramedRecords")
.field("rows", &self.rows)
.field("chunks", &self.chunks.len())
.finish()
}
}
pub fn needed(spec: &Spec) -> bool {
fn walked(field: &Field) -> bool {
!field.is_fixed_width()
|| field.size.as_ref().is_some_and(|s| !s.is_fixed())
|| field.count.as_ref().is_some_and(|s| !s.is_fixed())
|| field.delta != Delta::None
|| !field.bits.is_empty()
|| field.string_at.is_some()
}
let r = &spec.records;
r.framing != Framing::Fixed
|| !r.variants.is_empty()
|| r.checksum.is_some()
|| r.ring.is_some()
|| spec.blocks.is_some()
|| spec.capture.is_some()
|| formats::all_fields(r).any(walked)
}
struct Compiler<'a> {
spec: &'a Spec,
header: &'a HeaderValues,
data: &'a [u8],
slots: usize,
deltas: Vec<Delta>,
}
#[derive(Default, Clone)]
struct Scope(Vec<(String, usize)>);
impl Scope {
fn slot(&self, name: &str) -> Option<usize> {
self.0
.iter()
.rev()
.find(|(n, _)| n == name)
.map(|(_, s)| *s)
}
}
impl Compiler<'_> {
fn size(&self, amount: &Amount, scope: &Scope, what: &str) -> Result<SizeRef, String> {
Ok(match amount {
Amount::Record { field, adjust } => SizeRef::Slot {
slot: scope
.slot(field)
.ok_or_else(|| format!("{what}: `{field}` is not an earlier field"))?,
adjust: *adjust,
},
Amount::Rest => SizeRef::Rest,
other => SizeRef::Given(
usize::try_from(self.header.resolve(other, what)?)
.map_err(|_| format!("{what}: too large"))?,
),
})
}
fn fields(
&mut self,
fields: &[Field],
scope: &mut Scope,
columns: &mut Vec<OutColumn>,
shared: bool,
) -> Result<Vec<FieldPlan>, String> {
let mut out = Vec::new();
for field in fields {
out.push(self.field(field, scope, columns, shared)?);
}
Ok(out)
}
fn column(
columns: &mut Vec<OutColumn>,
name: &str,
dtype: DataType,
proto: Sink,
shared: bool,
) -> usize {
if shared && let Some(i) = columns.iter().position(|c| c.name == name) {
return i;
}
columns.push(OutColumn {
name: name.into(),
dtype,
proto,
});
columns.len() - 1
}
fn field(
&mut self,
field: &Field,
scope: &mut Scope,
columns: &mut Vec<OutColumn>,
shared: bool,
) -> Result<FieldPlan, String> {
let name = field.name.clone().unwrap_or_default();
let slot = self.slots;
self.slots += 1;
let big = field.endian.unwrap_or(self.spec.endian) == formats::Endian::Big;
let count = field
.count
.as_ref()
.map(|c| self.size(c, scope, "count"))
.transpose()?;
let given_count = match count {
Some(SizeRef::Given(n)) => {
if n as u64 > MAX_ITEMS {
return Err(format!(
"field `{name}`: {n} values is more than {MAX_ITEMS}"
));
}
Some(n)
}
_ => None,
};
let mut plan = FieldPlan {
name: name.clone(),
slot,
kind: Kind::Pad {
size: SizeRef::Given(0),
},
count,
outs: Vec::new(),
bits: Vec::new(),
delta: field.delta,
delta_index: 0,
sentinel: None,
};
let layout = |width: usize, cells: usize| -> Result<ColumnLayout, String> {
let mut layout = formats::layout_of(
self.spec,
field,
&name,
formats::Place {
start: 0,
stride: None,
width,
count: cells,
},
self.header,
)?;
layout.big_endian = big;
Ok(layout)
};
let packed = |layout: ColumnLayout| {
let cell = layout.width * layout.count;
Sink::Packed {
layout,
cell,
buf: Vec::new(),
valid: Vec::new(),
}
};
let mut fixed_out = |this: &mut Self,
plan: &mut FieldPlan,
mut layout: ColumnLayout|
-> Result<(), String> {
let _ = this;
match (count, given_count) {
(None, _) => {
let dtype = layout.dtype();
plan.outs = vec![Self::column(columns, &name, dtype, packed(layout), shared)];
}
(Some(_), Some(n)) if field.flatten => {
layout.count = 1;
for i in 0..n {
let n_name = format!("{name}_{i}");
let mut l = layout.clone();
l.name = n_name.as_str().into();
let dtype = l.dtype();
plan.outs
.push(Self::column(columns, &n_name, dtype, packed(l), shared));
}
}
(Some(_), Some(n)) => {
layout.count = n.max(1);
let dtype = layout.dtype();
plan.outs = vec![Self::column(columns, &name, dtype, packed(layout), shared)];
}
(Some(_), None) => {
layout.count = 1;
let item = layout.dtype();
let sink = Sink::List {
inner: Box::new(packed(layout)),
offsets: Vec::new(),
valid: Vec::new(),
};
plan.outs = vec![Self::column(
columns,
&name,
DataType::List(Box::new(item)),
sink,
shared,
)];
}
}
Ok(())
};
match field.ty {
Type::Pad => {
let size = field
.size
.as_ref()
.map_or(Ok(SizeRef::Given(0)), |s| self.size(s, scope, "size"))?;
plan.kind = Kind::Pad { size };
}
Type::Group => {
let mut inner_scope = Scope::default();
let mut inner_columns = Vec::new();
let first = self.slots;
let items =
self.fields(&field.group, &mut inner_scope, &mut inner_columns, false)?;
let slots = self.slots - first;
let names: Vec<PlSmallStr> = inner_columns.iter().map(|c| c.name.clone()).collect();
let struct_dtype = DataType::Struct(
inner_columns
.iter()
.map(|c| polars::prelude::Field::new(c.name.clone(), c.dtype.clone()))
.collect(),
);
let sink = Sink::List {
inner: Box::new(Sink::Struct {
names,
fields: inner_columns.into_iter().map(|c| c.proto).collect(),
len: 0,
}),
offsets: Vec::new(),
valid: Vec::new(),
};
plan.outs = vec![Self::column(
columns,
&name,
DataType::List(Box::new(struct_dtype)),
sink,
shared,
)];
plan.kind = Kind::Group { items, slots };
}
Type::VarU | Type::VarS => {
let signed = field.ty == Type::VarS;
let mut l = layout(8, 1)?;
if field.delta != Delta::None {
l.physical = Physical::Signed(8);
plan.sentinel = sentinel(field, 64, signed);
l.null = None;
}
plan.kind = Kind::Var { signed };
fixed_out(self, &mut plan, l)?;
}
Type::Strz => {
let max = field
.size
.as_ref()
.map(|s| self.size(s, scope, "size"))
.transpose()?;
plan.kind = Kind::Strz {
max,
encoding: field.encoding,
};
plan.outs = vec![Self::column(
columns,
&name,
DataType::String,
Sink::Text(Vec::new()),
shared,
)];
}
Type::Str | Type::Bytes if !field.size.as_ref().is_some_and(Amount::is_fixed) => {
let size = self.size(
field.size.as_ref().expect("checked at parse"),
scope,
"size",
)?;
let text = field.ty == Type::Str;
plan.kind = Kind::Sized {
size,
encoding: text.then_some(field.encoding),
};
let (dtype, sink) = if text {
(DataType::String, Sink::Text(Vec::new()))
} else {
(DataType::Binary, Sink::Binary(Vec::new()))
};
plan.outs = vec![Self::column(columns, &name, dtype, sink, shared)];
}
_ => {
let width = match (field.ty.width(), &field.size) {
(Some(w), _) => w as usize,
(None, Some(amount)) => usize::try_from(self.header.resolve(amount, "size")?)
.map_err(|_| "size: too large".to_string())?,
(None, None) => 0,
};
if width == 0 {
return Err(format!("field `{name}` takes no bytes"));
}
let int = match field.ty {
Type::Unsigned(_) => Some(IntRead { signed: false, big }),
Type::Signed(_) => Some(IntRead { signed: true, big }),
_ => None,
};
if let Some(section) = &field.string_at {
let section = self.section(section)?;
plan.kind = Kind::StringAt {
width,
big,
section,
};
plan.outs = vec![Self::column(
columns,
&name,
DataType::String,
Sink::Text(Vec::new()),
shared,
)];
} else {
let mut l = layout(width, 1)?;
if field.delta != Delta::None {
l.physical = Physical::Signed(8);
l.width = 8;
plan.sentinel =
sentinel(field, width as u32 * 8, int.is_some_and(|i| i.signed));
l.null = None;
}
plan.kind = Kind::Fixed { width, int };
fixed_out(self, &mut plan, l)?;
}
}
}
if field.delta != Delta::None {
plan.delta_index = self.deltas.len();
self.deltas.push(field.delta);
}
for bit in &field.bits {
let sink = Sink::Bits {
width: bit.width,
labels: bit.labels.clone(),
values: Vec::new(),
};
let dtype = match (&bit.labels, bit.width) {
(Some(_), _) => DataType::String,
(None, 1) => DataType::Boolean,
(None, 2..=8) => DataType::UInt8,
(None, 9..=16) => DataType::UInt16,
(None, 17..=32) => DataType::UInt32,
_ => DataType::UInt64,
};
let col = Self::column(columns, &bit.name, dtype, sink, shared);
plan.bits.push((col, bit.bit, bit.width));
}
if field.name.is_some() {
scope.0.push((name, slot));
}
Ok(plan)
}
fn section(&self, name: &str) -> Result<Range<usize>, String> {
let section = self
.spec
.sections
.iter()
.find(|s| s.name == name)
.ok_or_else(|| format!("no section named {name}"))?;
let offset = self.header.resolve_any(§ion.offset, "section offset")?;
let size = self.header.resolve_any(§ion.size, "section size")?;
let end = offset.saturating_add(size);
if end > self.data.len() as u64 {
return Err(format!(
"section {name} runs from byte {offset} to {end}, past the file's {} bytes",
self.data.len()
));
}
Ok(offset as usize..end as usize)
}
}
fn sentinel(field: &Field, bits: u32, signed: bool) -> Option<i128> {
use crate::fixed_records::Null;
let bits = bits.clamp(1, 64);
match field.null? {
Null::Min if signed => Some(-(1i128 << (bits - 1))),
Null::Min => Some(0),
Null::Max if signed => Some((1i128 << (bits - 1)) - 1),
Null::Max => Some((1i128 << bits) - 1),
Null::Value(v) => Some(v),
Null::NaN => None,
}
}
fn int_of(bytes: &[u8], read: IntRead) -> i128 {
if read.signed {
i128::from(crate::fixed_records::read_signed(bytes, read.big))
} else {
i128::from(crate::fixed_records::read_unsigned(bytes, read.big))
}
}
fn leb128(bytes: &[u8]) -> Option<(u64, usize)> {
let mut value = 0u64;
for (i, b) in bytes.iter().take(10).enumerate() {
let part = u64::from(b & 0x7f);
if i == 9 && part > 1 {
return None;
}
value |= part << (7 * i);
if b & 0x80 == 0 {
return Some((value, i + 1));
}
}
None
}
fn decode_text(bytes: &[u8], encoding: Encoding) -> String {
match encoding {
Encoding::Utf8 => crate::fixed_records::text(bytes),
Encoding::Latin1 => crate::fixed_records::latin1(bytes),
Encoding::Utf16Le => crate::fixed_records::utf16(bytes, false),
Encoding::Utf16Be => crate::fixed_records::utf16(bytes, true),
}
}
fn nul_at(bytes: &[u8], unit: usize) -> Option<usize> {
if unit == 1 {
memchr::memchr(0, bytes)
} else {
bytes
.chunks_exact(unit)
.position(|c| c.iter().all(|b| *b == 0))
.map(|i| i * unit)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Got {
Row,
Skipped,
None,
}
struct Out<'s> {
sinks: &'s mut [Sink],
filled: &'s mut [bool],
}
impl Out<'_> {
fn finish_row(&mut self) {
for (sink, filled) in self.sinks.iter_mut().zip(self.filled.iter_mut()) {
if !*filled {
sink.push_null();
}
*filled = false;
}
}
}
struct Walker<'a> {
plan: &'a Plan,
file: &'a [u8],
frame: Frame,
acc: Vec<i128>,
end: usize,
bounded: bool,
short: bool,
record_start: usize,
size_slot: Option<(usize, i64)>,
skipped: u64,
variant: Option<usize>,
}
impl<'a> Walker<'a> {
fn new(plan: &'a Plan, file: &'a [u8]) -> Self {
Self {
plan,
file,
frame: Frame::new(plan.slots),
acc: vec![0; plan.deltas.len()],
end: 0,
bounded: false,
short: false,
record_start: 0,
size_slot: None,
skipped: 0,
variant: None,
}
}
fn size_of(&self, s: SizeRef, pos: usize, what: &str) -> Result<usize, Stop> {
match s {
SizeRef::Given(n) => Ok(n),
SizeRef::Rest => Ok(self.end.saturating_sub(pos)),
SizeRef::Slot { slot, adjust } => {
let v = self.frame.ints[slot].ok_or(Stop::Truncated)? + i128::from(adjust);
usize::try_from(v)
.ok()
.filter(|n| *n as u64 <= formats::MAX_SIZE)
.ok_or_else(|| {
Stop::Said(format!(
"`{what}` at byte {pos} is {v}, outside 0 to {}",
formats::MAX_SIZE
))
})
}
}
}
fn take(&self, pos: &mut usize, n: usize) -> Result<Range<usize>, Stop> {
let start = *pos;
let stop = start.checked_add(n).ok_or(Stop::Truncated)?;
if stop > self.end {
return Err(Stop::Truncated);
}
*pos = stop;
Ok(start..stop)
}
fn walk(
&mut self,
fields: &[FieldPlan],
data: &[u8],
pos: &mut usize,
mut out: Option<&mut Out<'_>>,
) -> Result<(), Stop> {
for f in fields {
let res = if self.short {
Err(Stop::Truncated)
} else {
self.field(f, data, pos, out.as_deref_mut())
};
match res {
Ok(()) => {}
Err(Stop::Truncated) if self.bounded => {
self.short = true;
if let Some(out) = out.as_deref_mut() {
null_outs(f, out);
}
}
Err(e) => return Err(e),
}
if let Some((slot, adjust)) = self.size_slot
&& slot == f.slot
&& !self.bounded
{
let start = self.record_start;
let len = self.frame.ints[slot].ok_or(Stop::Truncated)? + i128::from(adjust);
let total = usize::try_from(len)
.ok()
.filter(|t| *t as u64 <= formats::MAX_SIZE && *t >= *pos - start)
.ok_or_else(|| {
Stop::Said(format!(
"the record at byte {start} gives its size as {len}, outside {} to {}",
*pos - start,
formats::MAX_SIZE
))
})?;
let rec_end = start + total;
if rec_end > self.end {
return Err(Stop::Truncated);
}
self.end = rec_end;
self.bounded = true;
}
}
Ok(())
}
fn field(
&mut self,
f: &FieldPlan,
data: &[u8],
pos: &mut usize,
out: Option<&mut Out<'_>>,
) -> Result<(), Stop> {
self.frame.starts[f.slot] = Some(*pos);
let count = match f.count {
None => None,
Some(c) => {
let n = self.size_of(c, *pos, &f.name)?;
if n as u64 > MAX_ITEMS {
return Err(Stop::Said(format!(
"`{}` at byte {} counts {n} values, more than {MAX_ITEMS}",
f.name, *pos
)));
}
Some(n)
}
};
match &f.kind {
Kind::Pad { size } => {
let n = self.size_of(*size, *pos, "pad")?;
self.take(pos, n)?;
}
Kind::Fixed { width, int } => {
let cells = count.unwrap_or(1);
let range = self.take(pos, width.checked_mul(cells).ok_or(Stop::Truncated)?)?;
let bytes = &data[range];
let raw = int.and_then(|r| (cells > 0).then(|| int_of(&bytes[..*width], r)));
if count.is_none() {
self.frame.ints[f.slot] = raw;
}
self.value(f, raw, Some((bytes, *width)), count, out);
}
Kind::Var { signed } => {
let cells = count.unwrap_or(1);
let mut bytes = Vec::with_capacity(cells.min(64) * 8);
let mut first = None;
for _ in 0..cells {
let (v, n) =
leb128(&data[(*pos).min(self.end)..self.end]).ok_or(Stop::Truncated)?;
*pos += n;
let v = if *signed {
i128::from(((v >> 1) as i64) ^ -((v & 1) as i64))
} else {
i128::from(v)
};
first.get_or_insert(v);
if *signed {
bytes.extend((v as i64).to_le_bytes());
} else {
bytes.extend((v as u64).to_le_bytes());
}
}
if count.is_none() {
self.frame.ints[f.slot] = first;
}
self.value(f, first, Some((&bytes, 8)), count, out);
}
Kind::Sized { size, encoding } => {
let n = self.size_of(*size, *pos, &f.name)?;
let range = self.take(pos, n)?;
if let Some(out) = out {
let col = f.outs[0];
out.filled[col] = true;
match (&mut out.sinks[col], encoding) {
(Sink::Text(v), Some(enc)) => {
v.push(Some(decode_text(&data[range], *enc)));
}
(Sink::Binary(v), None) => v.push(Some(data[range].to_vec())),
_ => {}
}
}
}
Kind::Strz { max, encoding } => {
let unit = encoding.unit();
let text = match max {
Some(m) => {
let n = self.size_of(*m, *pos, &f.name)?;
let range = self.take(pos, n)?;
let bytes = &data[range];
let stop = nul_at(bytes, unit).unwrap_or(bytes.len());
decode_text(&bytes[..stop], *encoding)
}
None => {
let rest = &data[(*pos).min(self.end)..self.end];
let stop = nul_at(rest, unit).ok_or(Stop::Truncated)?;
let text = decode_text(&rest[..stop], *encoding);
*pos += stop + unit;
text
}
};
if let Some(out) = out {
let col = f.outs[0];
out.filled[col] = true;
if let Sink::Text(v) = &mut out.sinks[col] {
v.push(Some(text));
}
}
}
Kind::StringAt {
width,
big,
section,
} => {
let range = self.take(pos, *width)?;
let offset = crate::fixed_records::read_unsigned(&data[range], *big);
self.frame.ints[f.slot] = Some(i128::from(offset));
if let Some(out) = out {
let col = f.outs[0];
out.filled[col] = true;
let heap = &self.file[section.clone()];
let text = usize::try_from(offset)
.ok()
.filter(|o| *o < heap.len())
.map(|o| {
let rest = &heap[o..];
let stop = memchr::memchr(0, rest).unwrap_or(rest.len());
crate::fixed_records::text(&rest[..stop])
});
if let Sink::Text(v) = &mut out.sinks[col] {
v.push(text);
}
}
}
Kind::Group { items, slots } => {
let n = count.unwrap_or(0);
let first = items.first().map_or(0, |i| i.slot);
let saved = (
self.end,
self.bounded,
self.short,
self.record_start,
self.size_slot,
);
self.size_slot = None;
let start = *pos;
let mut probe = start;
let mut fits = Ok(());
for _ in 0..n {
self.frame.clear(first..first + slots);
self.bounded = true;
self.short = false;
self.record_start = probe;
if let Err(e) = self.walk(items, data, &mut probe, None) {
fits = Err(e);
break;
}
if self.short {
fits = Err(Stop::Truncated);
break;
}
}
let (end, bounded, short, record_start, size_slot) = saved;
self.end = end;
if let Err(e) = fits {
self.bounded = bounded;
self.short = short;
self.record_start = record_start;
self.size_slot = size_slot;
return Err(e);
}
*pos = probe;
if let Some(out) = out {
let col = f.outs[0];
out.filled[col] = true;
if let Sink::List {
inner,
offsets,
valid,
} = &mut out.sinks[col]
&& let Sink::Struct { fields, len, .. } = inner.as_mut()
{
let mut at = start;
let mut filled = vec![false; fields.len()];
for _ in 0..n {
self.frame.clear(first..first + slots);
self.bounded = true;
self.short = false;
self.record_start = at;
let mut item_out = Out {
sinks: fields,
filled: &mut filled,
};
let _ = self.walk(items, data, &mut at, Some(&mut item_out));
item_out.finish_row();
*len += 1;
}
offsets.push(offsets.last().copied().unwrap_or(0) + n as i64);
valid.push(true);
}
}
self.end = end;
self.bounded = bounded;
self.short = short;
self.record_start = record_start;
self.size_slot = size_slot;
}
}
self.frame.ends[f.slot] = Some(*pos);
Ok(())
}
fn value(
&mut self,
f: &FieldPlan,
raw: Option<i128>,
bytes: Option<(&[u8], usize)>,
count: Option<usize>,
out: Option<&mut Out<'_>>,
) {
let summed = if f.delta == Delta::None {
None
} else {
let raw = raw.unwrap_or(0);
if Some(raw) == f.sentinel {
Some(None)
} else {
let sum = self.acc[f.delta_index].wrapping_add(raw);
self.acc[f.delta_index] = sum;
Some(i64::try_from(sum).ok())
}
};
let Some(out) = out else { return };
match (summed, bytes) {
(Some(sum), _) => {
let col = f.outs[0];
out.filled[col] = true;
match sum {
Some(v) => out.sinks[col].push_bytes(&v.to_le_bytes()),
None => out.sinks[col].push_null(),
}
}
(None, Some((bytes, width))) => push_fixed(f, bytes, width, count, out),
(None, None) => {}
}
push_bits(f, raw, out);
}
fn record(
&mut self,
data: &[u8],
pos: &mut usize,
chunk_end: usize,
out: Option<&mut Out<'_>>,
time: Option<i64>,
) -> Result<Got, Stop> {
let Some(only) = self.plan.only else {
return Ok(match self.record_inner(data, pos, chunk_end, out, time)? {
None => Got::None,
Some(_) => Got::Row,
});
};
let (start, acc, skipped) = (*pos, self.acc.clone(), self.skipped);
match self.record_inner(data, pos, chunk_end, None, time)? {
None => Ok(Got::None),
Some(Some(v)) if v == only => {
if out.is_some() {
*pos = start;
self.acc = acc;
self.skipped = skipped;
self.record_inner(data, pos, chunk_end, out, time)?;
}
Ok(Got::Row)
}
Some(_) => Ok(Got::Skipped),
}
}
fn record_inner(
&mut self,
data: &[u8],
pos: &mut usize,
chunk_end: usize,
mut out: Option<&mut Out<'_>>,
time: Option<i64>,
) -> Result<Option<Option<usize>>, Stop> {
let plan = self.plan;
if *pos >= chunk_end {
return Ok(None);
}
if !plan.sync.is_empty() {
let rest = &data[*pos..chunk_end];
if !rest.starts_with(&plan.sync) {
match memchr::memmem::find(rest, &plan.sync) {
Some(at) => {
self.skipped += at as u64;
*pos += at;
}
None => {
self.skipped += rest.len() as u64;
*pos = chunk_end;
return Ok(None);
}
}
}
*pos += plan.sync.len();
}
let first = plan.chunk_slots;
self.frame.clear(first..plan.slots);
self.record_start = *pos;
self.end = chunk_end;
self.bounded = false;
self.short = false;
self.size_slot = None;
match plan.size {
Some(SizeRef::Given(n)) => {
let end = pos.checked_add(n).ok_or(Stop::Truncated)?;
if end > chunk_end {
return Err(Stop::Truncated);
}
self.end = end;
self.bounded = true;
}
Some(SizeRef::Slot { slot, adjust }) => self.size_slot = Some((slot, adjust)),
_ => {}
}
let start = *pos;
let mut p = start;
self.walk(&plan.common, data, &mut p, out.as_deref_mut())?;
let mut label = None;
let mut chosen = None;
if let Some((slot, text)) = plan.type_slot {
let int = self.frame.ints[slot];
let shown = if text {
match (self.frame.starts[slot], self.frame.ends[slot]) {
(Some(a), Some(b)) => Some(crate::fixed_records::text(&data[a..b])),
_ => None,
}
} else {
None
};
let variant = plan
.variants
.iter()
.enumerate()
.find(|(_, v)| match (&shown, int) {
(Some(t), _) => v.texts.iter().any(|w| w == t),
(None, Some(i)) => v.ints.contains(&i),
_ => false,
});
match variant {
Some((index, variant)) => {
chosen = Some(index);
if let Some(size) = variant.size
&& !self.bounded
{
let n = self.size_of(size, p, "size")?;
let end = start.checked_add(n).ok_or(Stop::Truncated)?;
if end > chunk_end || end < p {
return Err(if end > chunk_end {
Stop::Truncated
} else {
Stop::Said(format!(
"variant {} at byte {start} is {n} bytes, less than its common fields",
variant.name
))
});
}
self.end = end;
self.bounded = true;
}
self.walk(&variant.fields, data, &mut p, out.as_deref_mut())?;
label = Some(variant.name.clone());
}
None => {
let shown = shown.or_else(|| int.map(|i| i.to_string()));
if !self.bounded {
return Err(Stop::Said(format!(
"the record at byte {start} has type {}, which no variant names",
shown.as_deref().unwrap_or("(none)")
)));
}
label = shown.map(|s| Arc::from(format!("?{s}")));
}
}
}
self.variant = chosen;
*pos = if self.bounded { self.end } else { p };
if let Some((width, big)) = plan.suffix {
let range = self.take_at(*pos, width, chunk_end)?;
let again = crate::fixed_records::read_unsigned(&data[range], big);
let (slot, _) = self.size_slot.unwrap_or((usize::MAX, 0));
let first = self.frame.ints.get(slot).copied().flatten();
if first != Some(i128::from(again)) {
return Err(Stop::Said(format!(
"the record at byte {start} ends with length {again}, not the {} it starts with",
first.map_or_else(|| "?".to_string(), |v| v.to_string())
)));
}
*pos += width;
}
if let Some(out) = out {
if let (Some(col), Some(label)) = (plan.type_out, label)
&& let Sink::Label(v) = &mut out.sinks[col]
{
v.push(Some(label));
out.filled[col] = true;
}
if let Some(check) = &plan.checksum {
let from = check.from.map_or(Some(start), |s| self.frame.starts[s]);
let to = self.frame.starts[check.to];
let stored = self.frame.ints[check.field];
let ok = match (from, to, stored) {
(Some(a), Some(b), Some(v)) if a <= b => {
Some(i128::from(check.algo.compute(&data[a..b])) == v)
}
_ => None,
};
if let Sink::Flag(v) = &mut out.sinks[check.out] {
v.push(ok);
out.filled[check.out] = true;
}
}
if let Some(col) = plan.time_out
&& let Sink::Time(v) = &mut out.sinks[col]
{
v.push(time);
out.filled[col] = true;
}
out.finish_row();
}
Ok(Some(chosen))
}
fn take_at(&self, pos: usize, n: usize, end: usize) -> Result<Range<usize>, Stop> {
let stop = pos.checked_add(n).ok_or(Stop::Truncated)?;
if stop > end {
return Err(Stop::Truncated);
}
Ok(pos..stop)
}
}
fn push_bits(f: &FieldPlan, raw: Option<i128>, out: &mut Out<'_>) {
for (col, bit, width) in &f.bits {
out.filled[*col] = true;
match (raw, &mut out.sinks[*col]) {
(Some(raw), Sink::Bits { values, .. }) => {
let mask = if *width >= 64 {
u64::MAX
} else {
(1u64 << width) - 1
};
values.push(Some(((raw as u64) >> bit) & mask));
}
(_, sink) => sink.push_null(),
}
}
}
fn null_outs(f: &FieldPlan, out: &mut Out<'_>) {
for col in f.outs.iter().chain(f.bits.iter().map(|(c, _, _)| c)) {
if !out.filled[*col] {
out.sinks[*col].push_null();
out.filled[*col] = true;
}
}
}
fn push_fixed(f: &FieldPlan, bytes: &[u8], width: usize, count: Option<usize>, out: &mut Out<'_>) {
if f.outs.len() == 1 {
let col = f.outs[0];
out.filled[col] = true;
match (&mut out.sinks[col], count) {
(
Sink::List {
inner,
offsets,
valid,
},
Some(n),
) => {
for i in 0..n {
inner.push_bytes(&bytes[i * width..(i + 1) * width]);
}
offsets.push(offsets.last().copied().unwrap_or(0) + n as i64);
valid.push(true);
}
(sink, _) => sink.push_bytes(bytes),
}
} else {
for (i, col) in f.outs.iter().enumerate() {
out.filled[*col] = true;
out.sinks[*col].push_bytes(&bytes[i * width..(i + 1) * width]);
}
}
}
enum ChunkData {
Map(Range<usize>),
Owned(Arc<Vec<u8>>),
}
struct Cursor {
pos: usize,
data_start: usize,
end: usize,
taken: u64,
limit: Option<u64>,
time: Option<i64>,
}
struct Want {
start: u64,
len: usize,
produced: usize,
}
impl Want {
fn null_row(&mut self, row: u64, out: &mut Out<'_>) {
if row >= self.start && self.produced < self.len {
out.finish_row();
self.produced += 1;
}
}
}
#[derive(Default)]
struct Found {
notes: Vec<String>,
skipped: u64,
}
impl FramedRecords {
pub fn open(
spec: &Spec,
bytes: Arc<Bytes>,
header: &HeaderValues,
data: Range<usize>,
named: &str,
path: Option<&std::path::Path>,
) -> Result<(Self, Vec<String>), String> {
let file = bytes.as_slice();
let mut compiler = Compiler {
spec,
header,
data: file,
slots: 0,
deltas: Vec::new(),
};
let mut columns = Vec::new();
let mut chunk_scope = Scope::default();
let mut discard = Vec::new();
let chunk_header = match &spec.capture {
Some(capture) => {
compiler.fields(&capture.header, &mut chunk_scope, &mut discard, false)?
}
None => Vec::new(),
};
let chunk_count = spec
.capture
.as_ref()
.and_then(|c| c.count.as_ref())
.and_then(|n| chunk_scope.slot(n));
let chunk_slots = compiler.slots;
let time_out = spec
.capture
.as_ref()
.and_then(|c| c.time.as_ref())
.map(|name| {
Compiler::column(
&mut columns,
name,
DataType::Datetime(TimeUnit::Nanoseconds, None),
Sink::Time(Vec::new()),
false,
)
});
let records = &spec.records;
let mut scope = Scope::default();
let common = compiler.fields(&records.fields, &mut scope, &mut columns, false)?;
let type_slot = records.type_field.as_ref().map(|name| {
let field = records
.fields
.iter()
.find(|f| f.name.as_deref() == Some(name.as_str()))
.expect("checked at parse");
(
scope.slot(name).expect("a common field"),
field.ty.is_text(),
)
});
let only = match &spec.variant {
Some(name) => Some(
records
.variants
.iter()
.position(|v| &v.name == name)
.ok_or_else(|| format!("no variant named {name}"))?,
),
None => None,
};
let type_out = type_slot.filter(|_| only.is_none()).map(|_| {
Compiler::column(
&mut columns,
"type",
DataType::from_categories(Categories::global()),
Sink::Label(Vec::new()),
false,
)
});
let size = records
.size
.as_ref()
.map(|s| compiler.size(s, &scope, "record size"))
.transpose()?;
let mut variants = Vec::new();
let mut unread = Vec::new();
for (i, variant) in records.variants.iter().enumerate() {
let mut vscope = scope.clone();
let into = if only.is_some_and(|o| o != i) {
&mut unread
} else {
&mut columns
};
let fields = compiler.fields(&variant.fields, &mut vscope, into, true)?;
let size = variant
.size
.as_ref()
.map(|s| compiler.size(s, &vscope, "variant size"))
.transpose()?;
if matches!(size, Some(SizeRef::Slot { .. })) {
return Err(format!(
"variant {}: its size comes from the header or is written in the spec",
variant.name
));
}
let mut ints = Vec::new();
let mut texts = Vec::new();
for when in &variant.when {
match when {
formats::Expected::Int(v) => ints.push(*v),
formats::Expected::Text(t) => {
texts.push(t.trim_end_matches(['\0', ' ']).to_string())
}
}
}
variants.push(VariantPlan {
name: Arc::from(variant.name.as_str()),
ints,
texts,
fields,
size,
});
}
let suffix = if records.length_suffix {
let Some(Amount::Record { field, .. }) = &records.size else {
return Err("length_suffix: needs size from a field".into());
};
let target = records
.fields
.iter()
.find(|f| f.name.as_deref() == Some(field.as_str()))
.ok_or_else(|| format!("length_suffix: `{field}` is not a common field"))?;
let width = target
.ty
.width()
.ok_or("length_suffix: the length field is fixed-width")?;
let big = target.endian.unwrap_or(spec.endian) == formats::Endian::Big;
Some((width as usize, big))
} else {
None
};
let checksum = records
.checksum
.as_ref()
.map(|c| -> Result<ChecksumPlan, String> {
let slot_of = |name: &str| -> Option<usize> {
scope.slot(name).or_else(|| {
variants
.iter()
.find_map(|v| v.fields.iter().find(|f| f.name == name).map(|f| f.slot))
})
};
let field = slot_of(&c.field).ok_or("checksum: no such field")?;
let to =
c.to.as_deref()
.map_or(Some(field), slot_of)
.ok_or("checksum: no such field")?;
let from = c
.from
.as_deref()
.map(slot_of)
.map(|s| s.ok_or("checksum: no such field"))
.transpose()?;
let out = Compiler::column(
&mut columns,
"checksum_ok",
DataType::Boolean,
Sink::Flag(Vec::new()),
false,
);
Ok(ChecksumPlan {
algo: c.algo,
field,
from,
to,
out,
})
})
.transpose()?;
if columns.is_empty() {
return Err("the records have no named field".into());
}
let fixed_size = match size {
Some(SizeRef::Given(n)) => Some(n),
None if variants.is_empty() => common.iter().try_fold(0usize, |sum, f| {
let width = match (&f.kind, f.count) {
(Kind::Fixed { width, .. }, None) => *width,
(Kind::Fixed { width, .. }, Some(SizeRef::Given(n))) => width.checked_mul(n)?,
(
Kind::Pad {
size: SizeRef::Given(n),
},
None,
) => *n,
_ => return None,
};
sum.checked_add(width)
}),
_ => None,
};
let size = match (size, fixed_size, records.framing) {
(None, Some(n), Framing::Fixed)
if !common.iter().any(|f| matches!(f.kind, Kind::Group { .. })) =>
{
Some(SizeRef::Given(n))
}
(size, _, _) => size,
};
if matches!(size, Some(SizeRef::Given(0))) && records.sync.is_empty() {
return Err("a record takes no bytes".into());
}
let plan = Plan {
framing: records.framing,
common,
type_slot,
variants,
type_out,
size,
suffix,
align: records.align as usize,
sync: records.sync.clone(),
checksum,
chunk_header,
chunk_count,
chunk_slots,
time_out,
slots: compiler.slots,
deltas: compiler.deltas,
columns,
only,
sources: Vec::new(),
};
let plan = Plan {
sources: sources(&plan),
..plan
};
let mut found = Found::default();
let chunks = if let Some(capture) = &spec.capture {
let _ = capture;
crate::framed_records::capture::packets(file, &mut found.notes)?
} else if let Some(blocks) = &spec.blocks {
list_blocks(spec, blocks, header, file, data.clone(), &mut found.notes)?
} else {
vec![Chunk {
source: ChunkSource::Map(data.clone()),
records: None,
time_ns: None,
}]
};
let schema: Schema = plan
.columns
.iter()
.map(|c| polars::prelude::Field::new(c.name.clone(), c.dtype.clone()))
.collect();
let mut records_read = Self {
bytes,
plan: Arc::new(plan),
chunks,
index: Arc::new(Index::Walk(Vec::new())),
table: None,
rows: 0,
schema: Arc::new(schema),
cache: Mutex::new(VecDeque::new()),
};
let ring = records
.ring
.as_ref()
.map(|r| header.resolve_any(r, "ring"))
.transpose()?;
let count = records
.count
.as_ref()
.map(|c| header.resolve_any(c, "count"))
.transpose()?;
let kept = path
.and_then(crate::indexed::peek::<KeptWalk>)
.filter(|k| k.spec == *spec && k.spec.variant == spec.variant && k.data == data);
if let Some(kept) = kept {
records_read.index = kept.index.clone();
records_read.table = kept.table.clone();
records_read.rows = kept.rows;
found.notes.extend(kept.notes.iter().cloned());
return Ok((records_read, found.notes));
}
let before = found.notes.len();
records_read.build_index(named, ring, count, &mut found)?;
if found.skipped > 0 {
found.notes.push(format!(
"{} {} skipped between records, to the next sync marker",
found.skipped,
if found.skipped == 1 { "byte" } else { "bytes" }
));
}
if let Some(path) = path
&& matches!(*records_read.index, Index::Walk(_))
{
crate::indexed::keep(
path,
Arc::new(KeptWalk {
spec: spec.clone(),
data,
index: records_read.index.clone(),
table: records_read.table.clone(),
rows: records_read.rows,
notes: found.notes[before..].to_vec(),
}),
);
}
Ok((records_read, found.notes))
}
fn stride(&self) -> Option<usize> {
let plan = &self.plan;
match (plan.size, &self.chunks[..]) {
(Some(SizeRef::Given(n)), [chunk])
if plan.framing == Framing::Fixed
&& plan.only.is_none()
&& plan.sync.is_empty()
&& plan.suffix.is_none()
&& plan.deltas.is_empty()
&& matches!(chunk.source, ChunkSource::Map(_))
&& plan.chunk_header.is_empty() =>
{
let align = plan.align.max(1);
Some(n.div_ceil(align) * align)
}
_ => None,
}
}
fn build_index(
&mut self,
named: &str,
ring: Option<u64>,
limit: Option<u64>,
found: &mut Found,
) -> Result<(), String> {
let file = self.bytes.as_slice();
if let Some(size) = self.stride() {
let ChunkSource::Map(range) = self.chunks[0].source.clone() else {
unreachable!("a stride is over the map")
};
let len = range.len();
let mut rows = len / size;
let unpadded = self.plan.size.map_or(size, |s| match s {
SizeRef::Given(n) => n,
_ => size,
});
if len % size >= unpadded {
rows += 1;
}
if let Some(count) = limit {
if count < rows as u64 {
rows = count as usize;
} else if count > rows as u64 {
found.notes.push(format!(
"header says {count} records {} {rows} whole ones shown",
crate::glyphs::get().middot
));
}
} else {
let used = (rows * size).min(len);
let trailing = len - used;
if trailing > 0 && len % size < unpadded {
found
.notes
.push(trailing_note(named, &file[range.end - trailing..range.end]));
}
}
let rows = rows.min(MAX_ROWS);
let ring = match ring {
Some(r) if rows > 0 => {
if r >= rows as u64 {
found.notes.push(format!(
"ring's oldest record {r} past the {rows} records {} read from the first",
crate::glyphs::get().middot
));
0
} else {
r as usize
}
}
_ => 0,
};
self.index = Arc::new(Index::Stride {
start: range.start,
size,
ring,
});
self.rows = rows;
return Ok(());
}
let plan = self.plan.clone();
let mut walker = Walker::new(&plan, file);
let mut checkpoints = Vec::new();
let mut rows: u64 = 0;
let max_rows = MAX_ROWS as u64;
let needs_walk_everything = plan.deltas.contains(&Delta::All);
let mut table = (self.chunks.len() == 1
&& matches!(self.chunks[0].source, ChunkSource::Map(_))
&& plan.chunk_header.is_empty()
&& plan.variants.len() < usize::from(WALK))
.then(|| RowTable {
starts: crate::indexed::Offsets::for_file(file.len()),
tags: Vec::new(),
});
'chunks: for ci in 0..self.chunks.len() {
if limit.is_some_and(|l| rows >= l) || rows >= max_rows {
break;
}
reset_block_sums(&plan, &mut walker.acc);
if let Some(n) = self.chunks[ci].records
&& !needs_walk_everything
{
checkpoints.push(Checkpoint {
row: rows,
chunk: ci as u32,
pos: u32::MAX,
taken: 0,
acc: walker.acc.clone().into_boxed_slice(),
});
rows = rows
.saturating_add(n)
.min(limit.unwrap_or(u64::MAX))
.min(max_rows);
continue;
}
let (data, mut cursor) = match self.enter(ci, &mut walker) {
Ok(c) => c,
Err(e) => {
found.notes.push(format!("block {ci}: {e}; left out"));
continue;
}
};
let chunk_bytes = self.slice(&data, file);
loop {
if cursor.limit.is_some_and(|l| cursor.taken >= l)
|| limit.is_some_and(|l| rows >= l)
|| rows >= max_rows
{
break;
}
if cursor.taken.is_multiple_of(CHECKPOINT) {
checkpoints.push(Checkpoint {
row: rows,
chunk: ci as u32,
pos: cursor.pos as u32,
taken: cursor.taken as u32,
acc: walker.acc.clone().into_boxed_slice(),
});
}
let before = cursor.pos;
match walker.record(chunk_bytes, &mut cursor.pos, cursor.end, None, None) {
Ok(got @ (Got::Row | Got::Skipped)) => {
if got == Got::Row {
rows += 1;
if let Some(t) = table.as_mut() {
if t.tags.len() >= crate::indexed::MAX_RECORDS {
table = None;
} else {
let tag = match (walker.short, plan.type_slot) {
(true, _) => WALK,
(false, None) => 0,
(false, Some(_)) => walker
.variant
.and_then(|v| u8::try_from(v).ok())
.unwrap_or(WALK),
};
t.starts.push(walker.record_start - plan.sync.len());
t.tags.push(tag);
}
}
}
cursor.taken += 1;
align(&mut cursor, plan.align);
if cursor.pos <= before {
found.notes.push(format!(
"zero-length record at byte {before} {} rest left out",
crate::glyphs::get().middot
));
break 'chunks;
}
}
Ok(Got::None) => {
if checkpoints.last().is_some_and(|c: &Checkpoint| {
c.row == rows && c.chunk == ci as u32 && c.pos == before as u32
}) {
checkpoints.pop();
}
break;
}
Err(stop) => {
if checkpoints.last().is_some_and(|c: &Checkpoint| {
c.row == rows && c.chunk == ci as u32 && c.pos == before as u32
}) {
checkpoints.pop();
}
let what = if self.chunks.len() > 1 {
format!("{named}, block {ci},")
} else {
named.to_string()
};
match stop {
Stop::Truncated => {
let rest = &chunk_bytes[before..cursor.end];
found.notes.push(trailing_note(&what, rest));
}
Stop::Said(said) => {
found.notes.push(format!(
"{said} {} {} bytes from there left out",
crate::glyphs::get().middot,
cursor.end - before
));
}
}
if self.chunks.len() == 1 {
break 'chunks;
}
break;
}
}
if cursor.pos >= cursor.end {
break;
}
}
if let Some(l) = cursor.limit
&& cursor.taken < l
&& self.chunks[ci].records.is_some()
{
found.notes.push(format!(
"block {ci}: {l} records declared, {} found",
cursor.taken
));
}
}
if let Some(l) = limit
&& rows < l
{
found.notes.push(format!(
"header says {l} records {} {rows} shown",
crate::glyphs::get().middot
));
}
found.skipped += walker.skipped;
self.rows = rows as usize;
self.index = Arc::new(Index::Walk(checkpoints));
self.table = table.filter(|t| t.tags.len() == self.rows).map(|mut t| {
t.starts.shrink();
t.tags.shrink_to_fit();
Arc::new(t)
});
Ok(())
}
fn slice<'b>(&'b self, data: &'b ChunkData, file: &'b [u8]) -> &'b [u8] {
match data {
ChunkData::Map(_) => file,
ChunkData::Owned(v) => v,
}
}
fn chunk_data(&self, i: usize) -> Result<ChunkData, String> {
match &self.chunks[i].source {
ChunkSource::Map(range) => Ok(ChunkData::Map(range.clone())),
ChunkSource::Block {
body,
codec,
uncompressed,
} => {
if let Ok(cache) = self.cache.lock()
&& let Some((_, data)) = cache.iter().find(|(k, _)| *k == i)
{
return Ok(ChunkData::Owned(data.clone()));
}
self.bytes.still_whole().map_err(|e| e.to_string())?;
let raw = &self.bytes.as_slice()[body.clone()];
let data = Arc::new(decompress(raw, *codec, *uncompressed)?);
if let Ok(mut cache) = self.cache.lock() {
cache.push_front((i, data.clone()));
cache.truncate(CACHED_BLOCKS);
}
Ok(ChunkData::Owned(data))
}
}
}
fn enter(&self, i: usize, walker: &mut Walker<'_>) -> Result<(ChunkData, Cursor), String> {
let data = self.chunk_data(i)?;
let (start, end) = match &data {
ChunkData::Map(r) => (r.start, r.end),
ChunkData::Owned(v) => (0, v.len()),
};
let mut cursor = Cursor {
pos: start,
data_start: start,
end,
taken: 0,
limit: self.chunks[i].records,
time: self.chunks[i].time_ns,
};
if !self.plan.chunk_header.is_empty() {
let file = self.bytes.as_slice();
let bytes = self.slice(&data, file);
walker.frame.clear(0..self.plan.chunk_slots);
walker.end = end;
walker.bounded = false;
walker.short = false;
walker.size_slot = None;
walker.record_start = start;
let mut p = start;
match walker.walk(&self.plan.chunk_header, bytes, &mut p, None) {
Ok(()) => {
cursor.pos = p;
cursor.data_start = p;
if let Some(slot) = self.plan.chunk_count {
cursor.limit = walker.frame.ints[slot].map(|v| v.max(0) as u64);
}
}
Err(_) => {
cursor.pos = end;
cursor.limit = Some(0);
}
}
}
Ok((data, cursor))
}
fn decode(&self, start: usize, len: usize, wanted: &[bool]) -> PolarsResult<DataFrame> {
self.bytes.still_whole()?;
let start = start.min(self.rows);
let len = len.min(self.rows - start);
let plan = &*self.plan;
let mut sinks: Vec<Sink> = plan
.columns
.iter()
.zip(wanted)
.map(|(c, w)| if *w { c.proto.clone() } else { Sink::Skip })
.collect();
let mut filled = vec![false; sinks.len()];
let file = self.bytes.as_slice();
let mut walker = Walker::new(plan, file);
match &*self.index {
Index::Stride {
start: base,
size,
ring,
} => {
let mut out = Out {
sinks: &mut sinks,
filled: &mut filled,
};
let ChunkSource::Map(range) = &self.chunks[0].source else {
unreachable!("a stride is over the map")
};
for row in start..start + len {
let k = (ring + row) % self.rows.max(1);
let mut pos = base + k * size;
let end = (pos + size).min(range.end);
if walker
.record(file, &mut pos, end, Some(&mut out), None)
.is_err()
{
out.finish_row();
}
}
}
Index::Walk(checkpoints) => {
let at = checkpoints.partition_point(|c| c.row <= start as u64);
let mut want = Want {
start: start as u64,
len,
produced: 0,
};
let mut out = Out {
sinks: &mut sinks,
filled: &mut filled,
};
if let Some(cp) = at.checked_sub(1).map(|i| &checkpoints[i]) {
self.read_from(cp, &mut walker, &mut want, &mut out);
}
for _ in want.produced..len {
out.finish_row();
}
}
}
let columns = sinks
.into_iter()
.zip(&plan.columns)
.zip(wanted)
.filter(|(_, w)| **w)
.map(|((sink, c), _)| sink.finish(c.name.clone()).map(Column::from))
.collect::<PolarsResult<Vec<_>>>()?;
DataFrame::new(len, columns)
}
fn read_from(
&self,
cp: &Checkpoint,
walker: &mut Walker<'_>,
want: &mut Want,
out: &mut Out<'_>,
) {
let file = self.bytes.as_slice();
let plan = &*self.plan;
walker.acc.copy_from_slice(&cp.acc);
let mut row = cp.row;
let mut ci = cp.chunk as usize;
let mut first = true;
while want.produced < want.len && ci < self.chunks.len() {
if !first {
reset_block_sums(plan, &mut walker.acc);
}
let (data, mut cursor) = match self.enter(ci, walker) {
Ok(c) => c,
Err(_) => {
let n = self.chunks[ci].records.unwrap_or(0);
for _ in 0..n {
want.null_row(row, out);
row += 1;
}
ci += 1;
first = false;
continue;
}
};
if first && cp.pos != u32::MAX {
cursor.pos = cp.pos as usize;
cursor.taken = u64::from(cp.taken);
}
first = false;
let bytes = self.slice(&data, file);
while want.produced < want.len {
if cursor.limit.is_some_and(|l| cursor.taken >= l) || cursor.pos >= cursor.end {
break;
}
let reading = row >= want.start;
let got = walker.record(
bytes,
&mut cursor.pos,
cursor.end,
if reading { Some(&mut *out) } else { None },
cursor.time,
);
match got {
Ok(Got::Row) => {
if reading {
want.produced += 1;
}
row += 1;
cursor.taken += 1;
align(&mut cursor, plan.align);
}
Ok(Got::Skipped) => {
cursor.taken += 1;
align(&mut cursor, plan.align);
}
Ok(Got::None) | Err(_) => break,
}
}
if let Some(l) = self.chunks[ci].records {
while cursor.taken < l && want.produced < want.len {
want.null_row(row, out);
row += 1;
cursor.taken += 1;
}
}
ci += 1;
}
}
fn decode_from(
&self,
table: &RowTable,
column: usize,
rows: &[IdxSize],
) -> PolarsResult<Column> {
self.bytes.still_whole()?;
let plan = &*self.plan;
let file = self.bytes.as_slice();
let ChunkSource::Map(range) = &self.chunks[0].source else {
unreachable!("a row table is over the map")
};
let mut sinks: Vec<Sink> = vec![Sink::Skip; plan.columns.len()];
sinks[column] = plan.columns[column].proto.clone();
let mut filled = vec![false; sinks.len()];
let mut out = Out {
sinks: &mut sinks,
filled: &mut filled,
};
let mut walker = Walker::new(plan, file);
if plan.type_out == Some(column) {
let mut labels: Vec<Option<Arc<str>>> =
plan.variants.iter().map(|v| Some(v.name.clone())).collect();
let codes: IdxCa = rows
.iter()
.map(|&row| {
let row = row as usize;
let tag = table.tags[row];
if tag != WALK {
return Some(IdxSize::from(tag));
}
let mut pos = table.starts.get(row);
let _ = walker.record(file, &mut pos, range.end, Some(&mut out), None);
let Sink::Label(said) = &mut out.sinks[column] else {
return None;
};
let label = said.pop().flatten()?;
labels.push(Some(label));
Some((labels.len() - 1) as IdxSize)
})
.collect();
let name = plan.columns[column].name.clone();
let labels = Sink::Label(labels).finish(name)?;
return Ok(labels.take(&codes)?.into_column());
}
let sources = &plan.sources[column];
for &row in rows {
let row = row as usize;
let start = table.starts.get(row);
let tag = table.tags[row];
let source = match tag {
WALK => &Source::Walk,
v => &sources[usize::from(v)],
};
match source {
Source::Null => out.finish_row(),
Source::Label => {
if let Sink::Label(v) = &mut out.sinks[column] {
v.push(plan.variants.get(usize::from(tag)).map(|v| v.name.clone()));
}
out.filled[column] = true;
out.finish_row();
}
Source::At { offset, field } => {
let Kind::Fixed { width, int } = field.kind else {
unreachable!("a value at a place is fixed")
};
let cells = match field.count {
Some(SizeRef::Given(n)) => Some(n),
_ => None,
};
let at = start + offset;
let bytes = at
.checked_add(width * cells.unwrap_or(1))
.filter(|end| *end <= range.end)
.map(|end| &file[at..end]);
if let Some(bytes) = bytes {
let raw = int
.and_then(|r| (cells != Some(0)).then(|| int_of(&bytes[..width], r)));
push_fixed(field, bytes, width, cells, &mut out);
push_bits(field, raw, &mut out);
}
out.finish_row();
}
Source::Walk | Source::Summed => {
let mut pos = start;
if !matches!(
walker.record(file, &mut pos, range.end, Some(&mut out), None),
Ok(Got::Row)
) {
out.finish_row();
}
}
}
}
let name = plan.columns[column].name.clone();
Ok(sinks.swap_remove(column).finish(name)?.into_column())
}
pub fn rows(&self) -> usize {
self.rows
}
pub fn schema(&self) -> SchemaRef {
self.schema.clone()
}
pub fn sources(&self) -> &[Arc<Bytes>] {
std::slice::from_ref(&self.bytes)
}
}
impl FramedRecords {
pub fn lazy(self: &Arc<Self>) -> LazyFrame {
crate::row_index::lazy(self)
}
pub fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
let all = vec![true; self.plan.columns.len()];
Ok(self.decode(start, len, &all)?.lazy())
}
pub fn collect(&self, rows: usize) -> PolarsResult<DataFrame> {
let all = vec![true; self.plan.columns.len()];
self.decode(0, rows, &all)
}
}
impl crate::row_index::RowSource for FramedRecords {
fn height(&self) -> usize {
self.rows
}
fn schema(&self) -> SchemaRef {
self.schema.clone()
}
fn decode(&self, column: usize, index: &IdxCa) -> PolarsResult<Column> {
let rows = crate::row_index::checked(index, self.rows)?;
if let Some(table) = &self.table
&& let Some(sources) = self.plan.sources.get(column)
&& !sources.iter().any(|s| matches!(s, Source::Summed))
{
return self.decode_from(table, column, &rows);
}
let mut wanted = vec![false; self.plan.columns.len()];
*wanted
.get_mut(column)
.ok_or_else(|| polars_err!(OutOfBounds: "no column {column}"))? = true;
let (Some(&lo), Some(&hi)) = (rows.iter().min(), rows.iter().max()) else {
return Ok(self.decode(0, 0, &wanted)?.columns()[0].clone());
};
let (lo, span) = (lo as usize, (hi - lo) as usize + 1);
let values = self.decode(lo, span, &wanted)?.columns()[0].clone();
let contiguous = rows.len() == span && rows.windows(2).all(|w| w[1] == w[0] + 1);
if contiguous {
return Ok(values);
}
let at = IdxCa::from_vec(
PlSmallStr::EMPTY,
rows.iter().map(|&r| r - lo as IdxSize).collect(),
);
values.take(&at)
}
}
impl crate::pushdown::Windowed for FramedRecords {
fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
FramedRecords::window(self, start, len)
}
}
fn align(cursor: &mut Cursor, align: usize) {
if align > 1 {
let offset = cursor.pos - cursor.data_start;
let rounded = offset.div_ceil(align).saturating_mul(align);
cursor.pos = cursor.data_start.saturating_add(rounded).min(cursor.end);
}
}
fn reset_block_sums(plan: &Plan, acc: &mut [i128]) {
for (sum, delta) in acc.iter_mut().zip(&plan.deltas) {
if *delta == Delta::Block {
*sum = 0;
}
}
}
fn trailing_note(what: &str, bytes: &[u8]) -> String {
formats::trailing_note(what, bytes)
}
impl Plan {
fn bare(fields: Vec<FieldPlan>, slots: usize) -> Self {
Self {
framing: Framing::Fixed,
common: fields,
type_slot: None,
variants: Vec::new(),
type_out: None,
size: None,
suffix: None,
align: 1,
sync: Vec::new(),
checksum: None,
chunk_header: Vec::new(),
chunk_count: None,
chunk_slots: 0,
time_out: None,
slots,
deltas: Vec::new(),
columns: Vec::new(),
only: None,
sources: Vec::new(),
}
}
}
fn sources(plan: &Plan) -> Vec<Vec<Source>> {
let variants = plan.variants.len().max(1);
let mut out = vec![vec![Source::Null; variants]; plan.columns.len()];
for v in (0..variants).filter(|v| plan.only.is_none_or(|only| only == *v)) {
let fields = plan
.common
.iter()
.chain(plan.variants.get(v).into_iter().flat_map(|p| &p.fields));
let mut offset = Some(plan.sync.len());
for f in fields {
let width = match (&f.kind, f.count) {
(Kind::Fixed { width, .. }, None) => Some(*width),
(Kind::Fixed { width, .. }, Some(SizeRef::Given(n))) => width.checked_mul(n),
(
Kind::Pad {
size: SizeRef::Given(n),
},
None,
) => Some(*n),
_ => None,
};
let source = match (offset, &f.kind, width) {
(Some(offset), Kind::Fixed { .. }, Some(_)) if f.delta == Delta::None => {
Source::At {
offset,
field: f.clone(),
}
}
_ => Source::Walk,
};
for &c in &f.outs {
out[c][v] = if f.delta == Delta::None {
source.clone()
} else {
Source::Summed
};
}
for (c, _, _) in &f.bits {
out[*c][v] = source.clone();
}
offset = offset.zip(width).and_then(|(o, w)| o.checked_add(w));
}
}
if let Some(c) = plan.type_out {
out[c].fill(Source::Label);
}
for c in plan.checksum.iter().map(|c| c.out).chain(plan.time_out) {
out[c].fill(Source::Walk);
}
out
}
struct Struct {
plan: Plan,
scope: Scope,
}
impl Struct {
fn compile(
spec: &Spec,
header: &HeaderValues,
file: &[u8],
fields: &[Field],
) -> Result<Self, String> {
let mut compiler = Compiler {
spec,
header,
data: file,
slots: 0,
deltas: Vec::new(),
};
let mut scope = Scope::default();
let mut columns = Vec::new();
let plans = compiler.fields(fields, &mut scope, &mut columns, false)?;
Ok(Self {
plan: Plan::bare(plans, compiler.slots),
scope,
})
}
fn read(&self, file: &[u8], pos: usize, end: usize) -> Result<(Frame, usize), Stop> {
let mut walker = Walker::new(&self.plan, file);
walker.end = end;
walker.record_start = pos;
let mut p = pos;
walker.walk(&self.plan.common, file, &mut p, None)?;
Ok((walker.frame, p))
}
fn value(&self, frame: &Frame, name: &str) -> Option<i128> {
frame.ints[self.scope.slot(name)?]
}
}
fn list_blocks(
spec: &Spec,
blocks: &formats::Blocks,
header: &HeaderValues,
file: &[u8],
data: Range<usize>,
notes: &mut Vec<String>,
) -> Result<Vec<Chunk>, String> {
let head = Struct::compile(spec, header, file, &blocks.header)?;
let size = {
let compiler = Compiler {
spec,
header,
data: file,
slots: 0,
deltas: Vec::new(),
};
compiler.size(&blocks.size, &head.scope, "block size")?
};
let codec_of = |frame: &Frame| -> Result<Compression, String> {
match &blocks.codec {
Codec::Fixed(c) => Ok(*c),
Codec::ByField { field, values } => {
let code = head
.value(frame, field)
.ok_or("the block header has no codec")?;
i64::try_from(code)
.ok()
.and_then(|c| values.get(&c).copied())
.ok_or_else(|| format!("codec {code} is not one compression names"))
}
}
};
let mut chunks = Vec::new();
let read_block = |at: usize,
end: usize,
rows: Option<u64>,
chunks: &mut Vec<Chunk>,
notes: &mut Vec<String>|
-> Result<Option<usize>, String> {
let (frame, body_start) = match head.read(file, at, end) {
Ok(x) => x,
Err(_) => {
notes.push(format!(
"the block header at byte {at} runs past the data; the rest is left out"
));
return Ok(None);
}
};
let body_len = match size {
SizeRef::Given(n) => n,
SizeRef::Slot { slot, adjust } => {
let v = frame.ints[slot].unwrap_or(0) + i128::from(adjust);
usize::try_from(v)
.ok()
.filter(|n| *n <= MAX_BLOCK)
.ok_or_else(|| {
format!(
"the block at byte {at} gives its size as {v}, outside 0 to {MAX_BLOCK}"
)
})?
}
SizeRef::Rest => end - body_start,
};
let body_end = body_start.checked_add(body_len).filter(|e| *e <= end);
let Some(body_end) = body_end else {
notes.push(format!(
"the block at byte {at} is {body_len} bytes and runs past the data; the rest is left out"
));
return Ok(None);
};
let codec = match codec_of(&frame) {
Ok(c) => c,
Err(e) => {
notes.push(format!("the block at byte {at}: {e}; left out"));
return Ok(Some(body_end));
}
};
let records = rows.or_else(|| {
blocks
.records
.as_ref()
.and_then(|r| head.value(&frame, r))
.map(|v| v.max(0) as u64)
});
let uncompressed = blocks
.uncompressed
.as_ref()
.and_then(|u| head.value(&frame, u))
.and_then(|v| usize::try_from(v).ok());
chunks.push(Chunk {
source: if codec == Compression::None {
ChunkSource::Map(body_start..body_end)
} else {
ChunkSource::Block {
body: body_start..body_end,
codec,
uncompressed,
}
},
records,
time_ns: None,
});
Ok(Some(body_end))
};
match &blocks.index {
Some(index) => {
let entry = Struct::compile(spec, header, file, &index.fields)?;
let at = usize::try_from(header.resolve_any(&index.at, "index at")?)
.map_err(|_| "index: too far")?;
let count = header.resolve_any(&index.count, "index count")?;
let width = formats::fields_width(&index.fields).unwrap_or(1).max(1) as usize;
let fits = file.len().saturating_sub(at) / width;
if count > fits as u64 {
return Err(format!(
"the block index at byte {at} says {count} entries; the file has room for {fits}"
));
}
let mut pos = at;
for _ in 0..count {
let (frame, next) = head_or(entry.read(file, pos, file.len()))?;
pos = next;
let offset = entry.value(&frame, "offset").unwrap_or(-1);
let rows = entry.value(&frame, "rows").map(|v| v.max(0) as u64);
let Some(offset) = usize::try_from(offset).ok().filter(|o| *o < file.len()) else {
notes.push(format!(
"an index entry points at byte {offset}, outside the file; left out"
));
continue;
};
read_block(offset, file.len(), rows, &mut chunks, notes)?;
}
}
None => {
let mut pos = data.start;
while pos < data.end {
match read_block(pos, data.end, None, &mut chunks, notes)? {
Some(next) if next > pos => pos = next,
Some(_) => {
notes.push(format!(
"zero-length block at byte {pos} {} rest left out",
crate::glyphs::get().middot
));
break;
}
None => break,
}
}
}
}
Ok(chunks)
}
fn head_or(read: Result<(Frame, usize), Stop>) -> Result<(Frame, usize), String> {
read.map_err(|_| "an index entry runs past the end of the file".to_string())
}
pub fn decompress(
raw: &[u8],
codec: Compression,
uncompressed: Option<usize>,
) -> Result<Vec<u8>, String> {
use std::io::Read;
let read_all = |mut reader: Box<dyn Read + '_>| -> Result<Vec<u8>, String> {
let mut out = Vec::new();
reader
.by_ref()
.take(MAX_BLOCK as u64 + 1)
.read_to_end(&mut out)
.map_err(|e| format!("{} block: {e}", codec.name()))?;
if out.len() > MAX_BLOCK {
return Err(format!(
"a block decompresses to more than {MAX_BLOCK} bytes"
));
}
Ok(out)
};
match codec {
Compression::None => Ok(raw.to_vec()),
Compression::Gzip => read_all(Box::new(flate2::read::MultiGzDecoder::new(raw))),
Compression::Deflate => read_all(Box::new(flate2::read::DeflateDecoder::new(raw))),
Compression::Zlib => read_all(Box::new(flate2::read::ZlibDecoder::new(raw))),
Compression::Zstd => read_all(Box::new(
zstd::Decoder::new(raw).map_err(|e| format!("zstd block: {e}"))?,
)),
Compression::Lz4 => read_all(Box::new(
lz4::Decoder::new(raw).map_err(|e| format!("lz4 block: {e}"))?,
)),
Compression::Lz4Block => {
let size = uncompressed.filter(|n| *n <= MAX_BLOCK).ok_or_else(|| {
format!("an lz4 block needs its decompressed size, at most {MAX_BLOCK}")
})?;
lz4::block::decompress(raw, Some(size as i32)).map_err(|e| format!("lz4 block: {e}"))
}
Compression::Snappy => {
let size = snap::raw::decompress_len(raw).map_err(|e| format!("snappy block: {e}"))?;
if size > MAX_BLOCK {
return Err(format!(
"a block decompresses to more than {MAX_BLOCK} bytes"
));
}
snap::raw::Decoder::new()
.decompress_vec(raw)
.map_err(|e| format!("snappy block: {e}"))
}
Compression::SnappyFramed => read_all(Box::new(snap::read::FrameDecoder::new(raw))),
Compression::Brotli => read_all(Box::new(brotli::Decompressor::new(raw, 4096))),
Compression::Bzip2 => read_all(Box::new(bzip2::read::BzDecoder::new(raw))),
Compression::Xz => read_all(Box::new(xz2::read::XzDecoder::new(raw))),
}
}
pub mod capture {
use super::{Chunk, ChunkSource};
const MAX_PACKETS: usize = 64 << 20;
fn u16_at(b: &[u8], at: usize, big: bool) -> Option<u16> {
let raw: [u8; 2] = b.get(at..at + 2)?.try_into().ok()?;
Some(if big {
u16::from_be_bytes(raw)
} else {
u16::from_le_bytes(raw)
})
}
fn u32_at(b: &[u8], at: usize, big: bool) -> Option<u32> {
let raw: [u8; 4] = b.get(at..at + 4)?.try_into().ok()?;
Some(if big {
u32::from_be_bytes(raw)
} else {
u32::from_le_bytes(raw)
})
}
pub fn is_capture(head: &[u8]) -> bool {
matches!(
head.get(..4),
Some(
[0xd4, 0xc3, 0xb2, 0xa1]
| [0xa1, 0xb2, 0xc3, 0xd4]
| [0x4d, 0x3c, 0xb2, 0xa1]
| [0xa1, 0xb2, 0x3c, 0x4d]
| [0x0a, 0x0d, 0x0d, 0x0a]
)
)
}
pub fn udp_payload(frame: &[u8], link: u32) -> Option<std::ops::Range<usize>> {
let (mut at, mut ethertype) = match link {
1 => (14, u16_at(frame, 12, true)?),
101 | 12 | 14 => (0, 0),
228 => (0, 0x0800),
229 => (0, 0x86dd),
113 => (16, u16_at(frame, 14, true)?),
276 => (20, u16_at(frame, 0, true)?),
0 | 108 => {
let family = u32_at(frame, 0, false)?;
let family = if family > 0xffff {
family.swap_bytes()
} else {
family
};
(4, if family == 2 { 0x0800 } else { 0x86dd })
}
_ => return None,
};
while ethertype == 0x8100 || ethertype == 0x88a8 {
ethertype = u16_at(frame, at + 2, true)?;
at += 4;
}
if ethertype == 0 {
ethertype = match frame.get(at)? >> 4 {
4 => 0x0800,
6 => 0x86dd,
_ => return None,
};
}
let (udp, ip_end) = match ethertype {
0x0800 => {
let ihl = usize::from(frame.get(at)? & 0x0f) * 4;
let total = usize::from(u16_at(frame, at + 2, true)?);
let flags = u16_at(frame, at + 6, true)?;
if ihl < 20 || *frame.get(at + 9)? != 17 || flags & 0x3fff != 0 {
return None;
}
(at + ihl, (at + total).min(frame.len()))
}
0x86dd => {
if *frame.get(at + 6)? != 17 {
return None;
}
let payload = usize::from(u16_at(frame, at + 4, true)?);
(at + 40, (at + 40 + payload).min(frame.len()))
}
_ => return None,
};
let len = usize::from(u16_at(frame, udp + 4, true)?);
let start = udp + 8;
let end = (udp + len).min(ip_end);
(len >= 8 && start <= end).then_some(start..end)
}
pub(super) fn packets(file: &[u8], notes: &mut Vec<String>) -> Result<Vec<Chunk>, String> {
let mut out = Vec::new();
let mut other = 0u64;
let magic = file
.get(..4)
.ok_or("the capture is shorter than its header")?;
if magic == [0x0a, 0x0d, 0x0d, 0x0a] {
pcapng(file, &mut out, &mut other)?;
} else {
pcap(file, &mut out, &mut other)?;
}
if other > 0 {
notes.push(format!(
"{other} {} in the capture {} not UDP, left out",
if other == 1 { "packet" } else { "packets" },
if other == 1 { "is" } else { "are" }
));
}
Ok(out)
}
fn push(
out: &mut Vec<Chunk>,
base: usize,
payload: std::ops::Range<usize>,
time_ns: Option<i64>,
) {
out.push(Chunk {
source: ChunkSource::Map(base + payload.start..base + payload.end),
records: None,
time_ns,
});
}
fn pcap(file: &[u8], out: &mut Vec<Chunk>, other: &mut u64) -> Result<(), String> {
let (big, nanos) = match file.get(..4) {
Some([0xd4, 0xc3, 0xb2, 0xa1]) => (false, false),
Some([0xa1, 0xb2, 0xc3, 0xd4]) => (true, false),
Some([0x4d, 0x3c, 0xb2, 0xa1]) => (false, true),
Some([0xa1, 0xb2, 0x3c, 0x4d]) => (true, true),
_ => return Err("not a pcap or pcapng capture: its magic is not one".into()),
};
let link =
u32_at(file, 20, big).ok_or("the capture is shorter than its header")? & 0x0fff_ffff;
let mut at = 24usize;
while at + 16 <= file.len() && out.len() < MAX_PACKETS {
let secs = i64::from(u32_at(file, at, big).unwrap_or(0));
let frac = i64::from(u32_at(file, at + 4, big).unwrap_or(0));
let caplen = u32_at(file, at + 8, big).unwrap_or(0) as usize;
let start = at + 16;
let Some(end) = start.checked_add(caplen).filter(|e| *e <= file.len()) else {
break;
};
let time = secs
.checked_mul(1_000_000_000)
.and_then(|s| s.checked_add(if nanos { frac } else { frac * 1000 }));
match udp_payload(&file[start..end], link) {
Some(payload) => push(out, start, payload, time),
None => *other += 1,
}
at = end;
}
Ok(())
}
fn pcapng(file: &[u8], out: &mut Vec<Chunk>, other: &mut u64) -> Result<(), String> {
let mut at = 0usize;
let mut big = false;
let mut interfaces: Vec<(u32, u64)> = Vec::new();
while at + 12 <= file.len() && out.len() < MAX_PACKETS {
let kind = u32_at(file, at, big).unwrap_or(0);
if kind == 0x0a0d_0d0a {
big = match file.get(at + 8..at + 12) {
Some([0x1a, 0x2b, 0x3c, 0x4d]) => true,
Some([0x4d, 0x3c, 0x2b, 0x1a]) => false,
_ => return Err("a pcapng section header has no byte-order magic".into()),
};
interfaces.clear();
}
let len = u32_at(file, at + 4, big).unwrap_or(0) as usize;
if len < 12 || !len.is_multiple_of(4) || at + len > file.len() {
break;
}
let body = &file[at + 8..at + len - 4];
match kind {
1 => {
let link = u32::from(u16_at(body, 0, big).unwrap_or(0));
let mut ticks = 1_000_000u64;
let mut o = 8;
while o + 4 <= body.len() {
let code = u16_at(body, o, big).unwrap_or(0);
let olen = usize::from(u16_at(body, o + 2, big).unwrap_or(0));
if code == 0 {
break;
}
if code == 9 && olen >= 1 {
let r = body[o + 4];
ticks = if r & 0x80 != 0 {
1u64.checked_shl(u32::from(r & 0x7f)).unwrap_or(1_000_000)
} else {
10u64.checked_pow(u32::from(r)).unwrap_or(1_000_000)
};
}
o += 4 + olen.div_ceil(4) * 4;
}
interfaces.push((link, ticks.max(1)));
}
6 => {
let iface = u32_at(body, 0, big).unwrap_or(0) as usize;
let high = u64::from(u32_at(body, 4, big).unwrap_or(0));
let low = u64::from(u32_at(body, 8, big).unwrap_or(0));
let caplen = u32_at(body, 12, big).unwrap_or(0) as usize;
let Some(frame) = body.get(20..20 + caplen) else {
at += len;
continue;
};
let (link, ticks) = interfaces.get(iface).copied().unwrap_or((1, 1_000_000));
let stamp = (high << 32) | low;
let time =
i64::try_from(u128::from(stamp) * 1_000_000_000 / u128::from(ticks)).ok();
match udp_payload(frame, link) {
Some(payload) => push(out, at + 8 + 20, payload, time),
None => *other += 1,
}
}
3 => {
let (link, _) = interfaces.first().copied().unwrap_or((1, 1_000_000));
let frame = &body[4.min(body.len())..];
match udp_payload(frame, link) {
Some(payload) => push(out, at + 8 + 4, payload, None),
None => *other += 1,
}
}
_ => {}
}
at += len;
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use crate::fixed_records::Bytes;
use crate::formats::{Opened, Spec};
use polars::prelude::*;
use std::sync::Arc;
fn open(spec: &str, bytes: Vec<u8>) -> Opened {
let spec = Spec::parse(spec, None).unwrap_or_else(|e| panic!("{e}"));
spec.open_rows(Arc::new(Bytes::Owned(bytes)), "f")
.unwrap_or_else(|e| panic!("{e}"))
}
fn all(opened: &Opened) -> DataFrame {
opened
.records
.clone()
.into_lazy()
.unwrap()
.collect()
.unwrap()
}
fn cell(df: &DataFrame, column: &str, row: usize) -> String {
match df.column(column).unwrap().get(row).unwrap() {
AnyValue::String(s) => s.to_string(),
AnyValue::StringOwned(s) => s.to_string(),
AnyValue::Categorical(..) | AnyValue::CategoricalOwned(..) => df
.column(column)
.unwrap()
.cast(&DataType::String)
.unwrap()
.get(row)
.unwrap()
.to_string()
.trim_matches('"')
.to_string(),
v => v.to_string(),
}
}
fn windows_agree(opened: &Opened) {
let whole = all(opened);
let rows = opened.records.rows();
assert_eq!(whole.height(), rows);
let walked = opened.records.window(0, rows).unwrap().collect().unwrap();
assert!(walked.equals_missing(&whole), "{walked}\n{whole}");
for start in [
0,
1,
rows / 2,
rows.saturating_sub(3),
1023,
1024,
1025,
2049,
] {
if start >= rows {
continue;
}
let window = opened.records.window(start, 3).unwrap().collect().unwrap();
let expected = whole.slice(start as i64, 3);
assert!(
window.equals_missing(&expected),
"window at {start}:\n{window}\n{expected}"
);
}
}
const ITCH: &str = r#"
name = "t.itch"
endian = "be"
[records]
framing = "length_prefixed"
size = "len"
size_adjust = 2
fields = [{ name = "len", type = "u2" }, { name = "kind", type = "str", size = 1 }]
type = "kind"
[[variants]]
name = "add"
when = "A"
fields = [{ name = "ref", type = "u8" }, { name = "shares", type = "u4" }, { name = "stock", type = "str", size = 8 }, { name = "price", type = "u4", scale = 4 }]
[[variants]]
name = "exec"
when = ["E", "C"]
fields = [{ name = "ref", type = "u8" }, { name = "shares", type = "u4" }]
"#;
fn itch_message(kind: u8, body: &[u8]) -> Vec<u8> {
let mut out = ((body.len() + 1) as u16).to_be_bytes().to_vec();
out.push(kind);
out.extend(body);
out
}
fn add(r: u64, shares: u32, stock: &str, price: u32) -> Vec<u8> {
let mut body = r.to_be_bytes().to_vec();
body.extend(shares.to_be_bytes());
let mut s = stock.as_bytes().to_vec();
s.resize(8, b' ');
body.extend(s);
body.extend(price.to_be_bytes());
itch_message(b'A', &body)
}
fn exec(r: u64, shares: u32) -> Vec<u8> {
let mut body = r.to_be_bytes().to_vec();
body.extend(shares.to_be_bytes());
itch_message(b'E', &body)
}
#[test]
fn length_prefixed_variants_make_one_table_with_a_type_column() {
let mut bytes = Vec::new();
for i in 0..3000u64 {
if i % 3 == 0 {
bytes.extend(add(i, 100, "AAPL", 1_234_500));
} else {
bytes.extend(exec(i, 7));
}
}
bytes.extend(itch_message(b'Z', &[1, 2, 3]));
let opened = open(ITCH, bytes);
assert!(opened.notes.is_empty(), "{:?}", opened.notes);
let df = all(&opened);
assert_eq!(df.height(), 3001);
assert_eq!(
df.get_column_names(),
["len", "kind", "type", "ref", "shares", "stock", "price"]
);
assert_eq!(cell(&df, "type", 0), "add");
assert_eq!(cell(&df, "type", 1), "exec");
assert_eq!(cell(&df, "stock", 0), "AAPL");
assert_eq!(cell(&df, "price", 0), "123.4500");
assert_eq!(cell(&df, "stock", 1), "null");
assert_eq!(cell(&df, "shares", 1), "7");
assert_eq!(cell(&df, "type", 3000), "?Z");
assert_eq!(cell(&df, "ref", 3000), "null");
windows_agree(&opened);
let one = opened
.records
.clone()
.into_lazy()
.unwrap()
.select([col("ref")])
.collect()
.unwrap();
assert_eq!(one.width(), 1);
assert_eq!(cell(&one, "ref", 2999), "2999");
}
#[test]
fn one_variant_reads_alone() {
let mut bytes = Vec::new();
for i in 0..2500u64 {
if i % 5 == 0 {
bytes.extend(add(i, 100, "MSFT", 1));
} else {
bytes.extend(exec(i, 7));
}
}
let spec = Spec::parse(ITCH, None)
.unwrap()
.with_variant("add")
.unwrap();
let opened = spec.open_rows(Arc::new(Bytes::Owned(bytes)), "f").unwrap();
let df = all(&opened);
assert_eq!(df.height(), 500);
assert_eq!(
df.get_column_names(),
["len", "kind", "ref", "shares", "stock", "price"]
);
assert_eq!(cell(&df, "ref", 499), "2495");
windows_agree(&opened);
assert!(
Spec::parse(ITCH, None)
.unwrap()
.with_variant("nope")
.is_err()
);
}
#[test]
fn a_variant_sizes_its_record_and_an_unknown_type_stops_the_read() {
let spec = r#"
name = "t.var"
[records]
framing = "variant"
type = { field = "kind", type = "u1" }
[[variants]]
name = "a"
when = 1
fields = [{ name = "x", type = "u2" }]
[[variants]]
name = "b"
when = 2
fields = [{ name = "y", type = "u4" }, { name = "z", type = "u1" }]
"#;
let bytes = vec![1, 5, 0, 2, 9, 0, 0, 0, 3, 1, 6, 0, 7, 0xaa];
let opened = open(spec, bytes);
let df = all(&opened);
assert_eq!(df.height(), 3);
assert_eq!(cell(&df, "x", 2), "6");
assert_eq!(cell(&df, "y", 1), "9");
assert_eq!(cell(&df, "z", 1), "3");
assert!(opened.notes[0].contains("type 7"), "{:?}", opened.notes);
windows_agree(&opened);
}
#[test]
fn sync_markers_find_frames_through_garbage() {
let spec = r#"
name = "t.sync"
[records]
framing = "sync"
sync = "1ACFFC1D"
fields = [{ name = "seq", type = "u2" }, { name = "v", type = "u1" }]
"#;
let mut bytes = vec![0xff, 0xee];
for i in 0..3u16 {
bytes.extend([0x1a, 0xcf, 0xfc, 0x1d]);
bytes.extend(i.to_le_bytes());
bytes.push(i as u8 * 10);
if i == 1 {
bytes.extend([1, 2, 3]);
}
}
let opened = open(spec, bytes);
let df = all(&opened);
assert_eq!(df.height(), 3);
assert_eq!(cell(&df, "v", 2), "20");
windows_agree(&opened);
assert!(
opened.notes.iter().any(|n| n.starts_with("5 bytes")),
"{:?}",
opened.notes
);
}
#[test]
fn checksums_bits_and_groups() {
let spec = r#"
name = "t.ck"
[records]
framing = "length_prefixed"
size = "len"
fields = [
{ name = "len", type = "u1" },
{ name = "status", type = "u2", bits = [{ name = "valid", bit = 0 }, { name = "mode", bit = 4, width = 3, enum = { 0 = "IDLE", 1 = "RUN" } }] },
{ name = "n", type = "u1" },
{ name = "levels", group = { count = "n", fields = [{ name = "px", type = "s2" }, { name = "qty", type = "u1" }] } },
{ name = "crc", type = "u2" },
]
checksum = { algo = "crc16-ccitt", field = "crc", from = "status" }
"#;
let record = |status: u16, levels: &[(i16, u8)], corrupt: bool| {
let mut body = status.to_le_bytes().to_vec();
body.push(levels.len() as u8);
for (px, q) in levels {
body.extend(px.to_le_bytes());
body.push(*q);
}
let mut crc = crate::formats::ChecksumAlgo::Crc16Ccitt.compute(&body) as u16;
if corrupt {
crc ^= 1;
}
let mut out = vec![(1 + body.len() + 2) as u8];
out.extend(body);
out.extend(crc.to_le_bytes());
out
};
let mut bytes = record(0x0011, &[(-5, 1), (7, 2)], false);
bytes.extend(record(0x0000, &[], true));
let opened = open(spec, bytes);
let df = all(&opened);
assert_eq!(df.height(), 2);
assert_eq!(cell(&df, "valid", 0), "true");
assert_eq!(cell(&df, "mode", 0), "RUN");
assert_eq!(cell(&df, "mode", 1), "IDLE");
assert_eq!(cell(&df, "checksum_ok", 0), "true");
assert_eq!(cell(&df, "checksum_ok", 1), "false");
let levels = df.column("levels").unwrap();
assert!(
matches!(levels.dtype(), DataType::List(inner) if matches!(**inner, DataType::Struct(_)))
);
let first = levels.get(0).unwrap().to_string();
assert!(first.contains("-5") && first.contains('7'), "{first}");
assert_eq!(levels.list().unwrap().lst_lengths().get(1), Some(0));
windows_agree(&opened);
}
#[test]
fn fortran_records_have_their_length_on_both_ends() {
let spec = r#"
name = "t.fortran"
[records]
framing = "length_prefixed"
size = "n"
size_adjust = 4
length_suffix = true
fields = [{ name = "n", type = "u4" }, { name = "v", type = "f8", count = "n_values" }]
"#;
assert!(Spec::parse(spec, None).is_err());
let spec = r#"
name = "t.fortran"
[records]
framing = "length_prefixed"
size = "n"
size_adjust = 4
length_suffix = true
fields = [{ name = "n", type = "u4" }, { name = "data", type = "bytes", size = "rest" }]
"#;
let mut bytes = Vec::new();
for payload in [&b"abc"[..], b"hello"] {
bytes.extend((payload.len() as u32).to_le_bytes());
bytes.extend(payload);
bytes.extend((payload.len() as u32).to_le_bytes());
}
let opened = open(spec, bytes.clone());
let df = all(&opened);
assert_eq!(df.height(), 2);
assert_eq!(
df.column("data").unwrap().binary().unwrap().get(1),
Some(&b"hello"[..])
);
let n = bytes.len();
bytes[n - 1] = 9;
let opened = open(spec, bytes);
assert_eq!(opened.records.rows(), 1);
assert!(
opened.notes[0].contains("ends with length"),
"{:?}",
opened.notes
);
}
#[test]
fn tagged_chunks_with_text_types_and_even_alignment() {
let spec = r#"
name = "t.riff"
[records]
framing = "length_prefixed"
size = "size"
size_adjust = 8
align = 2
fields = [{ name = "id", type = "str", size = 4 }, { name = "size", type = "u4" }]
type = "id"
[[variants]]
name = "fmt"
when = "fmt "
fields = [{ name = "channels", type = "u2" }]
[[variants]]
name = "data"
when = "data"
fields = [{ name = "payload", type = "bytes", size = "rest" }]
"#;
let mut bytes = b"fmt ".to_vec();
bytes.extend(2u32.to_le_bytes());
bytes.extend(2u16.to_le_bytes());
bytes.extend(b"data");
bytes.extend(3u32.to_le_bytes());
bytes.extend([1, 2, 3, 0]);
bytes.extend(b"LIST");
bytes.extend(0u32.to_le_bytes());
let opened = open(spec, bytes);
let df = all(&opened);
assert_eq!(df.height(), 3, "{:?}", opened.notes);
assert_eq!(cell(&df, "channels", 0), "2");
assert_eq!(
df.column("payload").unwrap().binary().unwrap().get(1),
Some(&[1u8, 2, 3][..])
);
assert_eq!(cell(&df, "type", 2), "?LIST");
windows_agree(&opened);
}
#[test]
fn strz_heap_strings_varints_and_deltas() {
let spec = r#"
name = "t.mixed"
[header]
fields = [{ name = "str_off", type = "u4" }, { name = "str_size", type = "u4" }, { name = "n", type = "u4" }]
[sections.strings]
offset = "header.str_off"
size = "header.str_size"
[records]
count = "header.n"
fields = [
{ name = "name", type = "strz" },
{ name = "label", type = "u2", string_at = "strings" },
{ name = "ts", type = "vu", delta = true, time = "ms" },
{ name = "dv", type = "vs" },
]
"#;
let heap = b"alpha\0beta\0";
let n = 2500u32;
let mut records = Vec::new();
for i in 0..n {
records.extend(format!("r{i}\0").as_bytes());
records.extend(if i % 2 == 0 { 0u16 } else { 6u16 }.to_le_bytes());
records.extend([0xe8, 0x07]);
records.push(5);
}
let off = 12 + records.len() as u32;
let mut bytes = off.to_le_bytes().to_vec();
bytes.extend((heap.len() as u32).to_le_bytes());
bytes.extend(n.to_le_bytes());
bytes.extend(&records);
bytes.extend(heap);
let opened = open(spec, bytes);
let df = all(&opened);
assert_eq!(df.height(), 2500);
assert_eq!(cell(&df, "name", 7), "r7");
assert_eq!(cell(&df, "label", 0), "alpha");
assert_eq!(cell(&df, "label", 1), "beta");
assert_eq!(cell(&df, "dv", 0), "-3");
assert_eq!(
df.column("ts").unwrap().dtype(),
&DataType::Datetime(TimeUnit::Milliseconds, None)
);
let ts = df.column("ts").unwrap().cast(&DataType::Int64).unwrap();
assert_eq!(ts.get(2048).unwrap(), AnyValue::Int64(2_049_000));
windows_agree(&opened);
}
#[test]
fn the_byte_order_comes_from_the_magic() {
let spec = r#"
name = "t.auto"
endian = "auto"
match = { magic = [0xd4, 0xc3, 0xb2, 0xa1] }
[header]
fields = [{ type = "pad", size = 4 }]
[records]
fields = [{ name = "v", type = "u4" }]
"#;
let le = [vec![0xd4, 0xc3, 0xb2, 0xa1], 7u32.to_le_bytes().to_vec()].concat();
let be = [vec![0xa1, 0xb2, 0xc3, 0xd4], 7u32.to_be_bytes().to_vec()].concat();
let parsed = Spec::parse(spec, None).unwrap();
assert!(parsed.magic_matches(&be));
for bytes in [le, be] {
let df = all(&open(spec, bytes));
assert_eq!(cell(&df, "v", 0), "7");
}
}
#[test]
fn a_footer_counts_the_records_and_checks_the_file() {
let spec = r#"
name = "t.foot"
[records]
count = "footer.n"
fields = [{ name = "v", type = "u2" }]
[footer]
fields = [{ name = "n", type = "u4" }, { name = "crc", type = "u4" }]
checksum = { algo = "crc32", field = "crc" }
"#;
let mut bytes: Vec<u8> = (0..5u16).flat_map(|v| v.to_le_bytes()).collect();
let crc = crate::formats::ChecksumAlgo::Crc32.compute(&bytes) as u32;
bytes.extend(4u32.to_le_bytes());
let mut good = bytes.clone();
good.extend(crc.to_le_bytes());
let opened = open(spec, good);
assert_eq!(opened.records.rows(), 4);
assert!(
!opened.notes.iter().any(|n| n.contains("checksum")),
"{:?}",
opened.notes
);
bytes.extend((crc ^ 1).to_le_bytes());
let opened = open(spec, bytes);
assert!(
opened.notes.iter().any(|n| n.contains("the file's is")),
"{:?}",
opened.notes
);
}
fn compress(codec: &str, raw: &[u8]) -> Vec<u8> {
use std::io::Write;
match codec {
"zstd" => zstd::encode_all(raw, 3).unwrap(),
"gzip" => {
let mut e = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::fast());
e.write_all(raw).unwrap();
e.finish().unwrap()
}
"zlib" => {
let mut e =
flate2::write::ZlibEncoder::new(Vec::new(), flate2::Compression::fast());
e.write_all(raw).unwrap();
e.finish().unwrap()
}
"deflate" => {
let mut e =
flate2::write::DeflateEncoder::new(Vec::new(), flate2::Compression::fast());
e.write_all(raw).unwrap();
e.finish().unwrap()
}
"lz4" => {
let mut e = lz4::EncoderBuilder::new().build(Vec::new()).unwrap();
e.write_all(raw).unwrap();
let (out, r) = e.finish();
r.unwrap();
out
}
"lz4_block" => lz4::block::compress(raw, None, false).unwrap(),
"snappy" => snap::raw::Encoder::new().compress_vec(raw).unwrap(),
"snappy_framed" => {
let mut e = snap::write::FrameEncoder::new(Vec::new());
e.write_all(raw).unwrap();
e.into_inner().unwrap()
}
"brotli" => {
let mut out = Vec::new();
{
let mut e = brotli::CompressorWriter::new(&mut out, 4096, 5, 22);
e.write_all(raw).unwrap();
}
out
}
"bzip2" => {
let mut e = bzip2::write::BzEncoder::new(Vec::new(), bzip2::Compression::fast());
e.write_all(raw).unwrap();
e.finish().unwrap()
}
"xz" => {
let mut e = xz2::write::XzEncoder::new(Vec::new(), 1);
e.write_all(raw).unwrap();
e.finish().unwrap()
}
_ => raw.to_vec(),
}
}
#[test]
fn a_running_sum_in_one_plain_block_reads_the_same_in_any_window() {
let spec = r#"
name = "t.one"
[blocks]
header = [{ name = "clen", type = "u4" }]
size = "clen"
[records]
fields = [{ name = "v", type = "u4", delta = "block" }]
"#;
let mut bytes = 40u32.to_le_bytes().to_vec();
bytes.extend((0..10u32).flat_map(|_| 1u32.to_le_bytes()));
let opened = open(spec, bytes);
let df = all(&opened);
assert_eq!(cell(&df, "v", 9), "10");
windows_agree(&opened);
}
#[test]
fn every_codec_reads_a_block() {
for codec in [
"none",
"gzip",
"deflate",
"zlib",
"zstd",
"lz4",
"lz4_block",
"snappy",
"snappy_framed",
"brotli",
"bzip2",
"xz",
] {
let spec = format!(
r#"
name = "t.blocks"
[blocks]
header = [{{ name = "clen", type = "u4" }}, {{ name = "rawlen", type = "u4" }}]
size = "clen"
compression = "{codec}"
uncompressed = "rawlen"
[records]
fields = [{{ name = "v", type = "u4", delta = "block" }}]
"#
);
let mut bytes = Vec::new();
for block in 0..3u32 {
let raw: Vec<u8> = (0..1500u32)
.flat_map(|_| (block + 1).to_le_bytes())
.collect();
let packed = compress(codec, &raw);
bytes.extend((packed.len() as u32).to_le_bytes());
bytes.extend((raw.len() as u32).to_le_bytes());
bytes.extend(packed);
}
let opened = open(&spec, bytes);
assert!(opened.notes.is_empty(), "{codec}: {:?}", opened.notes);
let df = all(&opened);
assert_eq!(df.height(), 4500, "{codec}");
assert_eq!(cell(&df, "v", 1499), "1500", "{codec}");
assert_eq!(cell(&df, "v", 1500), "2", "{codec}");
assert_eq!(cell(&df, "v", 4499), "4500", "{codec}");
windows_agree(&opened);
}
}
#[test]
fn a_block_index_in_the_footer_and_a_codec_per_block() {
let spec = r#"
name = "t.indexed"
[blocks]
header = [{ name = "codec", type = "u1" }, { name = "clen", type = "u4" }]
size = "clen"
compression = { field = "codec", values = { 0 = "none", 1 = "zstd" } }
index = { at = "footer.index_off", count = "footer.n_blocks", fields = [{ name = "offset", type = "u8" }, { name = "rows", type = "u4" }] }
[records]
fields = [{ name = "v", type = "u2" }]
[footer]
fields = [{ name = "index_off", type = "u8" }, { name = "n_blocks", type = "u4" }]
"#;
let mut bytes = Vec::new();
let mut index = Vec::new();
for block in 0..4u16 {
let raw: Vec<u8> = (0..100u16)
.flat_map(|i| (block * 100 + i).to_le_bytes())
.collect();
let (codec, body) = if block % 2 == 0 {
(0u8, raw.clone())
} else {
(1u8, zstd::encode_all(&raw[..], 1).unwrap())
};
index.extend((bytes.len() as u64).to_le_bytes());
index.extend(100u32.to_le_bytes());
bytes.push(codec);
bytes.extend((body.len() as u32).to_le_bytes());
bytes.extend(body);
}
let index_off = bytes.len() as u64;
bytes.extend(index);
bytes.extend(index_off.to_le_bytes());
bytes.extend(4u32.to_le_bytes());
let opened = open(spec, bytes);
assert!(opened.notes.is_empty(), "{:?}", opened.notes);
let df = all(&opened);
assert_eq!(df.height(), 400);
assert_eq!(cell(&df, "v", 399), "399");
windows_agree(&opened);
}
#[test]
fn columns_at_offsets_in_one_file() {
let spec = r#"
name = "t.cols"
layout = "columns"
[header]
fields = [{ name = "n", type = "u4" }, { name = "a_off", type = "u4" }, { name = "b_off", type = "u4" }]
[records]
count = "header.n"
fields = [{ name = "a", type = "u2", offset = "header.a_off" }, { name = "b", type = "f8", offset = "header.b_off" }]
"#;
let mut bytes = Vec::new();
bytes.extend(3u32.to_le_bytes());
bytes.extend(12u32.to_le_bytes());
bytes.extend(18u32.to_le_bytes());
for v in [1u16, 2, 3] {
bytes.extend(v.to_le_bytes());
}
for v in [0.5f64, 1.5, 2.5] {
bytes.extend(v.to_le_bytes());
}
let df = all(&open(spec, bytes));
assert_eq!(df.height(), 3);
assert_eq!(cell(&df, "a", 2), "3");
assert_eq!(cell(&df, "b", 1), "1.5");
}
#[test]
fn a_ring_buffer_starts_at_its_oldest_record() {
let spec = r#"
name = "t.ring"
[header]
fields = [{ name = "head", type = "u4" }]
[records]
ring = "header.head"
fields = [{ name = "v", type = "u1" }]
"#;
let mut bytes = 2u32.to_le_bytes().to_vec();
bytes.extend([30, 40, 10, 20]);
let df = all(&open(spec, bytes));
let values: Vec<String> = (0..4).map(|i| cell(&df, "v", i)).collect();
assert_eq!(values, ["10", "20", "30", "40"]);
}
#[test]
fn half_floats_and_text_encodings() {
let spec = r#"
name = "t.enc"
[records]
fields = [
{ name = "h", type = "f2" }, { name = "b", type = "bf2" },
{ name = "latin", type = "str", size = 4, encoding = "latin1" },
{ name = "wide", type = "str", size = 6, encoding = "utf16le" },
]
"#;
let mut bytes = half::f16::from_f32(1.5).to_bits().to_le_bytes().to_vec();
bytes.extend(half::bf16::from_f32(-2.0).to_bits().to_le_bytes());
bytes.extend([b'c', 0xe9, b' ', b' ']);
bytes.extend([b'h', 0, b'i', 0, 0, 0]);
let df = all(&open(spec, bytes));
assert_eq!(cell(&df, "h", 0), "1.5");
assert_eq!(cell(&df, "b", 0), "-2.0");
assert_eq!(cell(&df, "latin", 0), "c\u{e9}");
assert_eq!(cell(&df, "wide", 0), "hi");
}
fn udp_frame(payload: &[u8]) -> Vec<u8> {
let mut f = vec![0u8; 12];
f.extend([0x08, 0x00]);
let total = (20 + 8 + payload.len()) as u16;
f.extend([0x45, 0]);
f.extend(total.to_be_bytes());
f.extend([0, 0, 0x40, 0, 64, 17, 0, 0]);
f.extend([10, 0, 0, 1, 10, 0, 0, 2]);
f.extend([0x30, 0x39, 0x30, 0x39]);
f.extend(((8 + payload.len()) as u16).to_be_bytes());
f.extend([0, 0]);
f.extend(payload);
f
}
#[test]
fn a_capture_s_udp_payloads_hold_the_records() {
let spec = r#"
name = "t.mold"
endian = "be"
[capture]
header = [{ name = "session", type = "str", size = 10 }, { name = "seq", type = "u8" }, { name = "count", type = "u2" }]
count = "count"
time = "captured"
[records]
framing = "length_prefixed"
size = "len"
size_adjust = 2
fields = [{ name = "len", type = "u2" }, { name = "msg", type = "str", size = "rest" }]
"#;
let mut pcap = vec![0xd4, 0xc3, 0xb2, 0xa1, 2, 0, 4, 0];
pcap.extend([0u8; 8]);
pcap.extend(65535u32.to_le_bytes());
pcap.extend(1u32.to_le_bytes());
for (i, msgs) in [vec!["hi", "there"], vec!["x"]].into_iter().enumerate() {
let mut payload = b"SESSION001".to_vec();
payload.extend((i as u64).to_be_bytes());
payload.extend((msgs.len() as u16).to_be_bytes());
for m in &msgs {
payload.extend((m.len() as u16).to_be_bytes());
payload.extend(m.as_bytes());
}
let frame = udp_frame(&payload);
pcap.extend((1_700_000_000 + i as u32).to_le_bytes());
pcap.extend(5u32.to_le_bytes());
pcap.extend((frame.len() as u32).to_le_bytes());
pcap.extend((frame.len() as u32).to_le_bytes());
pcap.extend(frame);
}
let opened = open(spec, pcap);
let df = all(&opened);
assert_eq!(df.height(), 3, "{:?}", opened.notes);
assert_eq!(cell(&df, "msg", 1), "there");
assert_eq!(cell(&df, "msg", 2), "x");
assert_eq!(cell(&df, "captured", 2), "2023-11-14 22:13:21.000005");
}
#[test]
fn a_variant_s_fields_past_its_record_s_end_are_null() {
let mut short = add(2, 2, "B", 2);
short[..2].copy_from_slice(&13u16.to_be_bytes());
short.truncate(15);
let bytes = [add(1, 1, "A", 1), short, exec(3, 3), add(4, 4, "D", 4)].concat();
let opened = open(ITCH, bytes);
let df = all(&opened);
assert_eq!(df.height(), 4, "{:?}", opened.notes);
assert_eq!(cell(&df, "shares", 1), "2");
assert_eq!(cell(&df, "stock", 1), "null");
assert_eq!(cell(&df, "stock", 3), "D");
windows_agree(&opened);
}
#[test]
fn a_files_walk_is_kept_for_its_next_open() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("kept.itch");
let bytes: Vec<u8> = (0..3000u64)
.flat_map(|i| [add(i, i as u32, "S", 1), exec(i, 2)].concat())
.collect();
std::fs::write(&path, bytes).unwrap();
let spec = Spec::parse(ITCH, None).unwrap();
let first = spec.open(&path, "kept.itch").unwrap();
let kept = crate::indexed::peek::<super::KeptWalk>(&path).expect("kept");
assert_eq!(kept.rows, 6000);
let again = spec.open(&path, "kept.itch").unwrap();
let same = crate::indexed::peek::<super::KeptWalk>(&path).unwrap();
assert!(Arc::ptr_eq(&kept, &same), "not walked again");
assert!(all(&first).equals_missing(&all(&again)));
windows_agree(&again);
let exec_only = spec
.with_variant("exec")
.unwrap()
.open(&path, "kept.itch")
.unwrap();
assert_eq!(exec_only.records.rows(), 3000);
let replaced = crate::indexed::peek::<super::KeptWalk>(&path).unwrap();
assert_eq!(replaced.rows, 3000);
windows_agree(&exec_only);
}
#[test]
fn a_record_cut_short_is_left_out_and_said() {
let opened = open(
ITCH,
[add(1, 1, "A", 1), add(2, 2, "B", 2)[..10].to_vec()].concat(),
);
assert_eq!(opened.records.rows(), 1);
assert!(
opened.notes[0].contains("not a whole record"),
"{:?}",
opened.notes
);
}
}