Skip to main content

datui_lib/formats/
fixed_records.rs

1//! Fixed-width values read straight out of a memory map. A column is values `width`
2//! bytes wide, `stride` apart, from `start` into a source; `count` side by side make an
3//! Array cell. Records (stride = record size) and one file per column (stride = width)
4//! share one strided decoder for format specs and NumPy; [`decode`] is public for other
5//! packed readers (audio, CAN). Sources are kept for raw-byte viewers. Frames decode
6//! over a row index ([`crate::formats::row_index`]) on the streaming engine; an untouched view's
7//! window reads via [`FixedRecords::window`] without an index.
8
9use crate::formats::row_index::RowSource;
10use polars::prelude::*;
11use std::collections::BTreeMap;
12use std::path::Path;
13use std::sync::Arc;
14
15/// Nanoseconds in a day.
16const DAY_NS: i64 = 86_400_000_000_000;
17
18/// The bytes a column reads.
19pub enum Bytes {
20    /// A file, mapped, and the file, to ask its length again.
21    Mapped(memmap2::Mmap, std::fs::File),
22    /// Bytes in memory, for tests and the fuzz target.
23    Owned(Vec<u8>),
24}
25
26impl Bytes {
27    /// Map `path` read-only.
28    pub fn map(path: &Path) -> std::io::Result<Self> {
29        let file = std::fs::File::open(path)?;
30        // An empty file cannot be mapped on every platform, and has nothing to read.
31        if file.metadata()?.len() == 0 {
32            return Ok(Self::Owned(Vec::new()));
33        }
34        // SAFETY: the map is read-only and lives as long as the scan. A file truncated
35        // by another process while it is mapped faults (SIGBUS) on access past its new
36        // end, as it does for every reader that maps (Polars' own readers included).
37        // Datui does not write to the files it reads, and `still_whole` refuses a read
38        // once the file is shorter, which leaves only a truncation during a read.
39        let map = unsafe { memmap2::Mmap::map(&file)? };
40        Ok(Self::Mapped(map, file))
41    }
42
43    pub fn as_slice(&self) -> &[u8] {
44        match self {
45            Self::Mapped(map, _) => map,
46            Self::Owned(bytes) => bytes,
47        }
48    }
49
50    /// Fails when the mapped file is now shorter than its map, so reading the map
51    /// would go past the file's end.
52    pub fn still_whole(&self) -> PolarsResult<()> {
53        if let Self::Mapped(map, file) = self {
54            let len = file.metadata()?.len();
55            polars_ensure!(
56                len >= map.len() as u64,
57                ComputeError: "the file is now {len} bytes, shorter than the {} it had when it was opened; open it again",
58                map.len()
59            );
60        }
61        Ok(())
62    }
63
64    pub fn len(&self) -> usize {
65        self.as_slice().len()
66    }
67
68    pub fn is_empty(&self) -> bool {
69        self.len() == 0
70    }
71}
72
73/// How a value's bytes are stored.
74#[derive(Debug, Clone, Copy, PartialEq, Eq)]
75pub enum Physical {
76    /// An unsigned integer of 1 to 8 bytes.
77    Unsigned(u8),
78    /// A two's-complement integer of 1 to 8 bytes.
79    Signed(u8),
80    /// An IEEE float of 2 (half), 4 or 8 bytes.
81    Float(u8),
82    /// A bfloat16: the top half of an `f4`.
83    BFloat16,
84    /// One byte, nonzero for true.
85    Bool,
86    /// Text, its padding (NUL and spaces) trimmed from the right.
87    Text,
88    /// Text in ISO 8859-1, one byte a character, trimmed as `Text` is.
89    Latin1,
90    /// Text in UTF-16, two bytes a unit, trimmed as `Text` is.
91    Utf16 { big_endian: bool },
92    /// Text in UTF-32, four bytes a character, NUL characters trimmed from the right:
93    /// NumPy's `U` strings.
94    Utf32 { big_endian: bool },
95    /// Raw bytes.
96    Raw,
97}
98
99impl Physical {
100    /// Bytes the type itself takes, or `None` for text and raw bytes.
101    pub fn width(self) -> Option<usize> {
102        match self {
103            Self::Unsigned(n) | Self::Signed(n) | Self::Float(n) => Some(n as usize),
104            Self::Bool => Some(1),
105            Self::BFloat16 => Some(2),
106            Self::Text | Self::Latin1 | Self::Utf16 { .. } | Self::Utf32 { .. } | Self::Raw => None,
107        }
108    }
109
110    pub fn is_integer(self) -> bool {
111        matches!(self, Self::Unsigned(_) | Self::Signed(_))
112    }
113
114    fn plain_dtype(self) -> DataType {
115        match self {
116            Self::Unsigned(1) => DataType::UInt8,
117            Self::Unsigned(2) => DataType::UInt16,
118            Self::Unsigned(3 | 4) => DataType::UInt32,
119            Self::Unsigned(_) => DataType::UInt64,
120            Self::Signed(1) => DataType::Int8,
121            Self::Signed(2) => DataType::Int16,
122            Self::Signed(3 | 4) => DataType::Int32,
123            Self::Signed(_) => DataType::Int64,
124            Self::Float(2 | 4) | Self::BFloat16 => DataType::Float32,
125            Self::Float(_) => DataType::Float64,
126            Self::Bool => DataType::Boolean,
127            Self::Text | Self::Latin1 | Self::Utf16 { .. } | Self::Utf32 { .. } => DataType::String,
128            Self::Raw => DataType::Binary,
129        }
130    }
131}
132
133/// A stored value that means "no value".
134#[derive(Debug, Clone, Copy, PartialEq)]
135pub enum Null {
136    /// The type's smallest value: `i32::MIN` for an `s4`, 0 for a `u4`.
137    Min,
138    /// The type's largest value: `u32::MAX` for a `u4` (all ones), `i32::MAX` for an `s4`.
139    Max,
140    /// A float NaN.
141    NaN,
142    /// This value, as stored.
143    Value(i128),
144}
145
146/// What a value means beyond how it is stored.
147#[derive(Debug, Clone, PartialEq)]
148pub enum Logical {
149    /// The value as stored.
150    Plain,
151    /// An integer count of a unit since an epoch, as a datetime: `value * multiplier +
152    /// epoch`, in `unit`.
153    Timestamp {
154        unit: TimeUnit,
155        multiplier: i64,
156        epoch: i64,
157    },
158    /// A float count of a unit since an epoch, as a nanosecond datetime: serial dates.
159    FloatTimestamp { ns_per_unit: f64, epoch_ns: i64 },
160    /// An integer count of `multiplier` of a unit, as a duration in `unit`.
161    Duration { unit: TimeUnit, multiplier: i64 },
162    /// A count of days since `epoch_days` days past 1970-01-01, as a date.
163    Days { epoch_days: i32 },
164    /// An integer written as `YYYYMMDD`, as a date; anything else is null.
165    Yyyymmdd,
166    /// A count of a unit since midnight: a time of day, or a datetime on `date_ns`
167    /// (nanoseconds from the Unix epoch to that day's midnight).
168    TimeOfDay {
169        ns_per_unit: i64,
170        date_ns: Option<i64>,
171    },
172    /// An integer with `scale` implied decimal places.
173    Decimal { scale: usize },
174    /// `value * factor + offset`, as a float.
175    Linear { factor: f64, offset: f64 },
176    /// A code and its label; a code with no label reads as its number.
177    Enum(Arc<BTreeMap<i64, String>>),
178    /// An index into a list of symbols, as a categorical; one past the list is null.
179    Lookup(Arc<Vec<String>>),
180}
181
182/// One column: where its values are and how to read them.
183#[derive(Debug, Clone)]
184pub struct ColumnLayout {
185    pub name: PlSmallStr,
186    /// Index into the reader's sources.
187    pub source: usize,
188    /// Byte offset of row 0's value.
189    pub start: usize,
190    /// Bytes from one row's value to the next.
191    pub stride: usize,
192    /// Bytes one value takes.
193    pub width: usize,
194    /// Values side by side in each row: more than one makes the cell an Array.
195    pub count: usize,
196    pub physical: Physical,
197    pub big_endian: bool,
198    pub null: Option<Null>,
199    pub logical: Logical,
200}
201
202impl ColumnLayout {
203    /// A plain column of one `physical` value a row, little-endian.
204    pub fn new(name: &str, start: usize, stride: usize, physical: Physical, width: usize) -> Self {
205        Self {
206            name: name.into(),
207            source: 0,
208            start,
209            stride,
210            width,
211            count: 1,
212            physical,
213            big_endian: false,
214            null: None,
215            logical: Logical::Plain,
216        }
217    }
218
219    /// The type of one value.
220    pub fn value_dtype(&self) -> DataType {
221        match &self.logical {
222            Logical::Timestamp { unit, .. } => DataType::Datetime(*unit, None),
223            Logical::FloatTimestamp { .. } => DataType::Datetime(TimeUnit::Nanoseconds, None),
224            Logical::Duration { unit, .. } => DataType::Duration(*unit),
225            Logical::Days { .. } | Logical::Yyyymmdd => DataType::Date,
226            Logical::TimeOfDay { date_ns: None, .. } => DataType::Time,
227            Logical::TimeOfDay {
228                date_ns: Some(_), ..
229            } => DataType::Datetime(TimeUnit::Nanoseconds, None),
230            Logical::Decimal { scale } => DataType::Decimal(38, *scale),
231            Logical::Linear { .. } => DataType::Float64,
232            Logical::Enum(_) => DataType::String,
233            Logical::Lookup(_) => DataType::from_categories(Categories::global()),
234            Logical::Plain => self.physical.plain_dtype(),
235        }
236    }
237
238    /// The column's type: the value's, or an Array of `count` of them.
239    pub fn dtype(&self) -> DataType {
240        if self.count > 1 {
241            DataType::Array(Box::new(self.value_dtype()), self.count)
242        } else {
243            self.value_dtype()
244        }
245    }
246
247    /// Bytes one row's cell takes, `None` past `usize`.
248    fn cell_width(&self) -> Option<usize> {
249        self.width.checked_mul(self.count.max(1))
250    }
251
252    /// Rows this column can give from a source of `len` bytes.
253    fn rows_in(&self, len: usize) -> usize {
254        let Some(cell) = self.cell_width() else {
255            return 0;
256        };
257        match len
258            .checked_sub(self.start)
259            .and_then(|room| room.checked_sub(cell))
260        {
261            Some(after_first) if self.stride > 0 => after_first / self.stride + 1,
262            _ => 0,
263        }
264    }
265
266    /// Check the layout reads what it says: widths that fit the type, sensible counts.
267    pub fn validate(&self) -> PolarsResult<()> {
268        polars_ensure!(
269            self.stride > 0 && self.width > 0 && self.count > 0,
270            ComputeError: "column {} has no width", self.name
271        );
272        polars_ensure!(
273            self.cell_width().is_some(),
274            ComputeError: "column {}: {} values of {} bytes is too many", self.name, self.count, self.width
275        );
276        if let Some(width) = self.physical.width() {
277            polars_ensure!(
278                width == self.width,
279                ComputeError: "column {} is {} bytes wide, not {width}", self.name, self.width
280            );
281        }
282        if let Physical::Unsigned(n) | Physical::Signed(n) = self.physical {
283            polars_ensure!(
284                (1..=8).contains(&n),
285                ComputeError: "column {}: an integer is 1 to 8 bytes", self.name
286            );
287        }
288        if let Physical::Float(n) = self.physical {
289            polars_ensure!(
290                n == 2 || n == 4 || n == 8,
291                ComputeError: "column {}: a float is 2, 4 or 8 bytes", self.name
292            );
293        }
294        Ok(())
295    }
296}
297
298/// An integer of `bytes.len()` bytes, unsigned.
299pub fn read_unsigned(bytes: &[u8], big_endian: bool) -> u64 {
300    let mut value = 0u64;
301    if big_endian {
302        for b in bytes {
303            value = (value << 8) | u64::from(*b);
304        }
305    } else {
306        for b in bytes.iter().rev() {
307            value = (value << 8) | u64::from(*b);
308        }
309    }
310    value
311}
312
313/// An integer of `bytes.len()` bytes, two's complement.
314pub fn read_signed(bytes: &[u8], big_endian: bool) -> i64 {
315    let raw = read_unsigned(bytes, big_endian);
316    let bits = bytes.len() * 8;
317    if bits == 0 || bits >= 64 {
318        raw as i64
319    } else {
320        // Sign-extend from the value's own width.
321        ((raw << (64 - bits)) as i64) >> (64 - bits)
322    }
323}
324
325/// The value `null` stands for in a `physical` integer, as stored.
326fn null_integer(null: Null, physical: Physical) -> Option<i128> {
327    let bits = u32::from(match physical {
328        Physical::Unsigned(n) | Physical::Signed(n) => n,
329        Physical::Bool => 1,
330        _ => return None,
331    }) * 8;
332    let signed = matches!(physical, Physical::Signed(_));
333    match null {
334        Null::Min if signed => Some(-(1i128 << (bits - 1))),
335        Null::Min => Some(0),
336        Null::Max if signed => Some((1i128 << (bits - 1)) - 1),
337        Null::Max => Some((1i128 << bits) - 1),
338        Null::Value(v) => Some(v),
339        Null::NaN => None,
340    }
341}
342
343/// The first `rows` values of `column` from `bytes`: one cell a row, an Array of
344/// `column.count` values each when it holds more than one. Asking for more rows than
345/// `bytes` holds is an error.
346pub fn decode(bytes: &[u8], column: &ColumnLayout, rows: usize) -> PolarsResult<Column> {
347    column.validate()?;
348    let fits = column.rows_in(bytes.len());
349    polars_ensure!(
350        rows <= fits,
351        ComputeError: "column {}: {rows} rows asked for, {fits} in {} bytes", column.name, bytes.len()
352    );
353    decode_strided(bytes, column, rows, |row| row)
354}
355
356/// The values of `column` from `bytes` for the rows `index` names, in that order. A
357/// missing row, or one `bytes` does not hold, is an error.
358pub fn decode_rows(bytes: &[u8], column: &ColumnLayout, index: &IdxCa) -> PolarsResult<Column> {
359    column.validate()?;
360    let rows = crate::formats::row_index::checked(index, column.rows_in(bytes.len()))?;
361    decode_strided(bytes, column, rows.len(), |i| rows[i] as usize)
362}
363
364/// The values of `column` from `bytes` for records starting at `records`, in order
365/// (`column.start` into each; `stride` unused): for unevenly spaced records like log
366/// messages found by an index. A cell past the end is an error.
367pub fn decode_at(bytes: &[u8], column: &ColumnLayout, records: &[usize]) -> PolarsResult<Column> {
368    // The stride is not used here, so a layout may leave it 0.
369    ColumnLayout {
370        stride: column.stride.max(1),
371        ..column.clone()
372    }
373    .validate()?;
374    let cell = column.cell_width().unwrap_or(usize::MAX);
375    for &record in records {
376        let end = record
377            .checked_add(column.start)
378            .and_then(|start| start.checked_add(cell));
379        polars_ensure!(
380            end.is_some_and(|end| end <= bytes.len()),
381            OutOfBounds: "column {}: a record at {record} runs past the {} bytes on hand", column.name, bytes.len()
382        );
383    }
384    let start = column.start;
385    decode_cells(bytes, column, records.len(), |i| records[i] + start)
386}
387
388/// `rows` cells of `column`, the `i`th of them its row `row(i)`; every row is one
389/// `bytes` holds.
390fn decode_strided(
391    bytes: &[u8],
392    column: &ColumnLayout,
393    rows: usize,
394    row: impl Fn(usize) -> usize,
395) -> PolarsResult<Column> {
396    let (start, stride) = (column.start, column.stride);
397    decode_cells(bytes, column, rows, move |i| start + row(i) * stride)
398}
399
400/// `rows` cells of `column`, the `i`th of them starting at byte `cell(i)`; every cell
401/// is one `bytes` holds.
402fn decode_cells(
403    bytes: &[u8],
404    column: &ColumnLayout,
405    rows: usize,
406    cell: impl Fn(usize) -> usize,
407) -> PolarsResult<Column> {
408    let count = column.count.max(1);
409    let values = rows
410        .checked_mul(count)
411        .ok_or_else(|| polars_err!(ComputeError: "column {}: too many values", column.name))?;
412    // Each value's bytes, row by row and in a row left to right.
413    let at = move |i: usize| {
414        let start = cell(i / count) + (i % count) * column.width;
415        &bytes[start..start + column.width]
416    };
417    let name = column.name.clone();
418    let flat = decode_values(column, values, at)?.with_name(name.clone());
419    if count == 1 {
420        return Ok(flat.into_column());
421    }
422    flat.reshape_array(&[
423        ReshapeDimension::new(rows as i64),
424        ReshapeDimension::new(count as i64),
425    ])
426    .map(|s| s.with_name(name).into_column())
427}
428
429/// `n` values, the `i`th of them in `at(i)`.
430fn decode_values<'a>(
431    column: &ColumnLayout,
432    n: usize,
433    at: impl Fn(usize) -> &'a [u8],
434) -> PolarsResult<Series> {
435    let big = column.big_endian;
436    let name = PlSmallStr::EMPTY;
437    macro_rules! native {
438        ($ty:ty) => {{
439            const N: usize = std::mem::size_of::<$ty>();
440            (0..n)
441                .map(|i| {
442                    let raw: [u8; N] = at(i).try_into().expect("width checked at build");
443                    if big {
444                        <$ty>::from_be_bytes(raw)
445                    } else {
446                        <$ty>::from_le_bytes(raw)
447                    }
448                })
449                .collect::<Vec<$ty>>()
450        }};
451    }
452    // A native width with a sentinel: compared at that width, not through `i128`.
453    if column.logical == Logical::Plain
454        && let Some(null) = column.null
455    {
456        let sentinel = null_integer(null, column.physical);
457        macro_rules! or_null {
458            ($ty:ty) => {{
459                let sentinel = sentinel.and_then(|s| <$ty>::try_from(s).ok());
460                let values = native!($ty).into_iter();
461                Some(Series::new(
462                    name.clone(),
463                    values
464                        .map(|v| (Some(v) != sentinel).then_some(v))
465                        .collect::<Vec<Option<$ty>>>(),
466                ))
467            }};
468        }
469        let fast = match column.physical {
470            Physical::Unsigned(1) => or_null!(u8),
471            Physical::Unsigned(2) => or_null!(u16),
472            Physical::Unsigned(4) => or_null!(u32),
473            Physical::Unsigned(8) => or_null!(u64),
474            Physical::Signed(1) => or_null!(i8),
475            Physical::Signed(2) => or_null!(i16),
476            Physical::Signed(4) => or_null!(i32),
477            Physical::Signed(8) => or_null!(i64),
478            _ => None,
479        };
480        if let Some(series) = fast {
481            return Ok(series);
482        }
483    }
484    // The common case at memory speed: a native width, as stored, no sentinel.
485    if column.logical == Logical::Plain && column.null.is_none() {
486        let fast = match column.physical {
487            Physical::Unsigned(1) => Some(Series::new(name.clone(), native!(u8))),
488            Physical::Unsigned(2) => Some(Series::new(name.clone(), native!(u16))),
489            Physical::Unsigned(4) => Some(Series::new(name.clone(), native!(u32))),
490            Physical::Unsigned(8) => Some(Series::new(name.clone(), native!(u64))),
491            Physical::Signed(1) => Some(Series::new(name.clone(), native!(i8))),
492            Physical::Signed(2) => Some(Series::new(name.clone(), native!(i16))),
493            Physical::Signed(4) => Some(Series::new(name.clone(), native!(i32))),
494            Physical::Signed(8) => Some(Series::new(name.clone(), native!(i64))),
495            Physical::Float(4) => Some(Series::new(name.clone(), native!(f32))),
496            Physical::Float(8) => Some(Series::new(name.clone(), native!(f64))),
497            // Odd widths straight into the type they widen to: through `i128` and a
498            // cast they took four times as long as a native width (#662).
499            Physical::Unsigned(3) => Some(Series::new(
500                name.clone(),
501                (0..n)
502                    .map(|i| read_unsigned(at(i), big) as u32)
503                    .collect::<Vec<u32>>(),
504            )),
505            Physical::Unsigned(5..=7) => Some(Series::new(
506                name.clone(),
507                (0..n)
508                    .map(|i| read_unsigned(at(i), big))
509                    .collect::<Vec<u64>>(),
510            )),
511            Physical::Signed(3) => Some(Series::new(
512                name.clone(),
513                (0..n)
514                    .map(|i| read_signed(at(i), big) as i32)
515                    .collect::<Vec<i32>>(),
516            )),
517            Physical::Signed(5..=7) => Some(Series::new(
518                name.clone(),
519                (0..n)
520                    .map(|i| read_signed(at(i), big))
521                    .collect::<Vec<i64>>(),
522            )),
523            _ => None,
524        };
525        if let Some(series) = fast {
526            return Ok(series);
527        }
528    }
529    match column.physical {
530        Physical::Unsigned(_) | Physical::Signed(_) | Physical::Bool => {
531            let signed = matches!(column.physical, Physical::Signed(_));
532            let sentinel = column
533                .null
534                .and_then(|null| null_integer(null, column.physical));
535            let ints: Vec<Option<i128>> = (0..n)
536                .map(|i| {
537                    let bytes = at(i);
538                    let v = if signed {
539                        i128::from(read_signed(bytes, big))
540                    } else {
541                        i128::from(read_unsigned(bytes, big))
542                    };
543                    (Some(v) != sentinel).then_some(v)
544                })
545                .collect();
546            integers(column, ints)
547        }
548        Physical::Float(_) | Physical::BFloat16 => {
549            let physical = column.physical;
550            let floats: Vec<Option<f64>> = (0..n)
551                .map(|i| {
552                    let raw = read_unsigned(at(i), big);
553                    let v = match physical {
554                        Physical::Float(2) => f64::from(half::f16::from_bits(raw as u16)),
555                        Physical::BFloat16 => f64::from(half::bf16::from_bits(raw as u16)),
556                        Physical::Float(4) => f64::from(f32::from_bits(raw as u32)),
557                        _ => f64::from_bits(raw),
558                    };
559                    let null = match column.null {
560                        Some(Null::NaN) => v.is_nan(),
561                        Some(Null::Value(sentinel)) => v == sentinel as f64,
562                        _ => false,
563                    };
564                    (!null).then_some(v)
565                })
566                .collect();
567            let width = physical.width().unwrap_or(8) as u8;
568            floats_of(column, width, floats)
569        }
570        Physical::Text => {
571            let values: StringChunked = (0..n).map(|i| Some(text(at(i)))).collect();
572            Ok(values.into_series())
573        }
574        Physical::Latin1 => {
575            let values: StringChunked = (0..n).map(|i| Some(latin1(at(i)))).collect();
576            Ok(values.into_series())
577        }
578        Physical::Utf16 { big_endian } => {
579            let values: StringChunked = (0..n).map(|i| Some(utf16(at(i), big_endian))).collect();
580            Ok(values.into_series())
581        }
582        Physical::Utf32 { big_endian } => {
583            let values: StringChunked = (0..n).map(|i| Some(utf32(at(i), big_endian))).collect();
584            Ok(values.into_series())
585        }
586        Physical::Raw => {
587            let values: BinaryChunked = (0..n).map(|i| Some(at(i))).collect();
588            Ok(values.into_series())
589        }
590    }
591}
592
593/// Integers as their column means them: plain at its width, or as its `logical` says
594/// (a datetime, a duration, a factor and offset, an enum's labels). Public for readers
595/// that find their integers some other way, such as bit fields of CAN signals.
596pub fn integers(column: &ColumnLayout, ints: Vec<Option<i128>>) -> PolarsResult<Series> {
597    let as_i64 = |v: Option<i128>| v.and_then(|v| i64::try_from(v).ok());
598    Ok(match &column.logical {
599        Logical::Plain => match column.physical {
600            Physical::Bool => ints
601                .into_iter()
602                .map(|v| v.map(|v| v != 0))
603                .collect::<BooleanChunked>()
604                .into_series(),
605            physical => {
606                let wide: Int128Chunked = ints.into_iter().collect();
607                // Each value fits the type `plain_dtype` names for its width.
608                wide.into_series().strict_cast(&physical.plain_dtype())?
609            }
610        },
611        Logical::Timestamp {
612            unit,
613            multiplier,
614            epoch,
615        } => ints
616            .into_iter()
617            .map(|v| as_i64(v)?.checked_mul(*multiplier)?.checked_add(*epoch))
618            .collect::<Int64Chunked>()
619            .into_datetime(*unit, None)
620            .into_series(),
621        Logical::Duration { unit, multiplier } => ints
622            .into_iter()
623            .map(|v| as_i64(v)?.checked_mul(*multiplier))
624            .collect::<Int64Chunked>()
625            .into_duration(*unit)
626            .into_series(),
627        Logical::FloatTimestamp {
628            ns_per_unit,
629            epoch_ns,
630        } => float_timestamps(
631            ints.into_iter().map(|v| v.map(|v| v as f64)),
632            *ns_per_unit,
633            *epoch_ns,
634        ),
635        Logical::Days { epoch_days } => ints
636            .into_iter()
637            .map(|v| i32::try_from(as_i64(v)?).ok()?.checked_add(*epoch_days))
638            .collect::<Int32Chunked>()
639            .into_date()
640            .into_series(),
641        Logical::Yyyymmdd => ints
642            .into_iter()
643            .map(|v| yyyymmdd_days(as_i64(v)?))
644            .collect::<Int32Chunked>()
645            .into_date()
646            .into_series(),
647        Logical::TimeOfDay {
648            ns_per_unit,
649            date_ns,
650        } => {
651            let of_day = ints.into_iter().map(|v| {
652                let ns = as_i64(v)?.checked_mul(*ns_per_unit)?;
653                (0..DAY_NS).contains(&ns).then_some(ns)
654            });
655            match date_ns {
656                None => of_day.collect::<Int64Chunked>().into_time().into_series(),
657                Some(day) => of_day
658                    .map(|ns| ns?.checked_add(*day))
659                    .collect::<Int64Chunked>()
660                    .into_datetime(TimeUnit::Nanoseconds, None)
661                    .into_series(),
662            }
663        }
664        Logical::Decimal { scale } => ints
665            .into_iter()
666            .collect::<Int128Chunked>()
667            .into_decimal_unchecked(38, *scale)
668            .into_series(),
669        Logical::Linear { factor, offset } => ints
670            .into_iter()
671            .map(|v| v.map(|v| v as f64 * factor + offset))
672            .collect::<Float64Chunked>()
673            .into_series(),
674        Logical::Lookup(symbols) => ints
675            .into_iter()
676            .map(|v| {
677                let i = usize::try_from(v?).ok()?;
678                symbols.get(i).map(String::as_str)
679            })
680            .collect::<StringChunked>()
681            .into_series()
682            .cast(&DataType::from_categories(Categories::global()))?,
683        Logical::Enum(labels) => ints
684            .into_iter()
685            .map(|v| {
686                v.map(|code| {
687                    i64::try_from(code)
688                        .ok()
689                        .and_then(|code| labels.get(&code).cloned())
690                        .unwrap_or_else(|| code.to_string())
691                })
692            })
693            .collect::<StringChunked>()
694            .into_series(),
695    })
696}
697
698/// Floats as their column means them.
699fn floats_of(column: &ColumnLayout, width: u8, floats: Vec<Option<f64>>) -> PolarsResult<Series> {
700    Ok(match &column.logical {
701        Logical::Linear { factor, offset } => floats
702            .into_iter()
703            .map(|v| v.map(|v| v * factor + offset))
704            .collect::<Float64Chunked>()
705            .into_series(),
706        Logical::FloatTimestamp {
707            ns_per_unit,
708            epoch_ns,
709        } => float_timestamps(floats.into_iter(), *ns_per_unit, *epoch_ns),
710        _ if width <= 4 => floats
711            .into_iter()
712            .map(|v| v.map(|v| v as f32))
713            .collect::<Float32Chunked>()
714            .into_series(),
715        _ => floats.into_iter().collect::<Float64Chunked>().into_series(),
716    })
717}
718
719fn float_timestamps(
720    values: impl Iterator<Item = Option<f64>>,
721    ns_per_unit: f64,
722    epoch_ns: i64,
723) -> Series {
724    values
725        .map(|v| {
726            let ns = (v? * ns_per_unit).round();
727            // Out of range or not a number reads as null rather than wrapping.
728            (ns.is_finite() && ns.abs() < 9.0e18).then(|| (ns as i64).checked_add(epoch_ns))?
729        })
730        .collect::<Int64Chunked>()
731        .into_datetime(TimeUnit::Nanoseconds, None)
732        .into_series()
733}
734
735/// Days since 1970-01-01 of a date written as the integer `YYYYMMDD`.
736fn yyyymmdd_days(v: i64) -> Option<i32> {
737    let (year, month, day) = (v / 10_000, (v / 100) % 100, v % 100);
738    let date = chrono::NaiveDate::from_ymd_opt(
739        i32::try_from(year).ok()?,
740        u32::try_from(month).ok()?,
741        u32::try_from(day).ok()?,
742    )?;
743    let epoch = chrono::NaiveDate::from_ymd_opt(1970, 1, 1)?;
744    i32::try_from((date - epoch).num_days()).ok()
745}
746
747/// Text from a fixed-width field: UTF-8 where it is, padding trimmed from the right.
748pub fn text(bytes: &[u8]) -> String {
749    let end = bytes
750        .iter()
751        .rposition(|&b| b != 0 && b != b' ')
752        .map_or(0, |i| i + 1);
753    String::from_utf8_lossy(&bytes[..end]).into_owned()
754}
755
756/// Text from a fixed-width ISO 8859-1 field, padding trimmed as [`text`] trims it.
757pub fn latin1(bytes: &[u8]) -> String {
758    let end = bytes
759        .iter()
760        .rposition(|&b| b != 0 && b != b' ')
761        .map_or(0, |i| i + 1);
762    bytes[..end].iter().map(|&b| char::from(b)).collect()
763}
764
765/// Text from a fixed-width UTF-16 field, NUL and space units trimmed from the right.
766/// An odd last byte and unpaired surrogates read as the replacement character.
767pub fn utf16(bytes: &[u8], big_endian: bool) -> String {
768    let mut units: Vec<u16> = bytes
769        .as_chunks::<2>()
770        .0
771        .iter()
772        .map(|&pair| {
773            if big_endian {
774                u16::from_be_bytes(pair)
775            } else {
776                u16::from_le_bytes(pair)
777            }
778        })
779        .collect();
780    while units.last().is_some_and(|&u| u == 0 || u == 0x20) {
781        units.pop();
782    }
783    String::from_utf16_lossy(&units)
784}
785
786/// Text from a fixed-width UTF-32 field, NUL characters trimmed from the right as NumPy
787/// trims them. A unit that is no character, and an odd last few bytes, read as the
788/// replacement character.
789pub fn utf32(bytes: &[u8], big_endian: bool) -> String {
790    let (units, rest) = bytes.as_chunks::<4>();
791    let mut chars: Vec<char> = units
792        .iter()
793        .map(|&unit| {
794            let code = if big_endian {
795                u32::from_be_bytes(unit)
796            } else {
797                u32::from_le_bytes(unit)
798            };
799            char::from_u32(code).unwrap_or(char::REPLACEMENT_CHARACTER)
800        })
801        .collect();
802    while chars.last() == Some(&'\0') {
803        chars.pop();
804    }
805    let mut text: String = chars.into_iter().collect();
806    if !rest.is_empty() {
807        text.push(char::REPLACEMENT_CHARACTER);
808    }
809    text
810}
811
812/// Bytes as space-separated hex pairs.
813pub fn hex(bytes: &[u8]) -> String {
814    const DIGITS: &[u8; 16] = b"0123456789abcdef";
815    let mut out = String::with_capacity(bytes.len() * 3);
816    for (i, b) in bytes.iter().enumerate() {
817        if i > 0 {
818            out.push(' ');
819        }
820        out.push(DIGITS[(b >> 4) as usize] as char);
821        out.push(DIGITS[(b & 0xf) as usize] as char);
822    }
823    out
824}
825
826/// Columns of fixed-width values over one or more sources, `rows` long.
827pub struct FixedRecords {
828    sources: Vec<Arc<Bytes>>,
829    columns: Vec<ColumnLayout>,
830    rows: usize,
831    schema: SchemaRef,
832}
833
834impl FixedRecords {
835    /// The columns over `sources`, at most `rows` long (fewer when a source runs out, never
836    /// reading past one, dropping a partial last record), from the sources' current lengths.
837    /// A grown file rebuilds records over a fresh map with `usize::MAX` rows. At most
838    /// [`crate::formats::row_index::MAX_ROWS`] are shown.
839    pub fn new(
840        sources: Vec<Arc<Bytes>>,
841        columns: Vec<ColumnLayout>,
842        rows: usize,
843    ) -> PolarsResult<Self> {
844        let mut rows = rows.min(crate::formats::row_index::MAX_ROWS);
845        for column in &columns {
846            column.validate()?;
847            let source = sources
848                .get(column.source)
849                .ok_or_else(|| polars_err!(ComputeError: "column {} has no source", column.name))?;
850            rows = rows.min(column.rows_in(source.len()));
851        }
852        let schema: Schema = columns
853            .iter()
854            .map(|c| Field::new(c.name.clone(), c.dtype()))
855            .collect();
856        // The frame decodes a column by its place in the schema.
857        polars_ensure!(
858            schema.len() == columns.len(),
859            Duplicate: "two columns have the same name"
860        );
861        Ok(Self {
862            sources,
863            columns,
864            rows,
865            schema: Arc::new(schema),
866        })
867    }
868
869    pub fn rows(&self) -> usize {
870        self.rows
871    }
872
873    pub fn schema(&self) -> SchemaRef {
874        self.schema.clone()
875    }
876
877    /// The bytes the columns read, for a reader of the raw bytes.
878    pub fn sources(&self) -> &[Arc<Bytes>] {
879        &self.sources
880    }
881
882    /// The frame: decoded over a row index, only what a query reaches.
883    pub fn lazy(self: &Arc<Self>) -> LazyFrame {
884        crate::formats::row_index::lazy(self)
885    }
886
887    /// Rows `[start, start + len)`, decoded now: the columns start further into the
888    /// same sources, so nothing before `start` is decoded and no index is built.
889    pub fn window(&self, start: usize, len: usize) -> PolarsResult<DataFrame> {
890        let start = start.min(self.rows);
891        let len = len.min(self.rows - start);
892        let columns = self
893            .columns
894            .iter()
895            .map(|c| ColumnLayout {
896                start: c.start + start * c.stride,
897                ..c.clone()
898            })
899            .collect();
900        Self::new(self.sources.clone(), columns, len)?.collect(len)
901    }
902
903    /// The first `rows` rows of every column, decoded now.
904    pub fn collect(&self, rows: usize) -> PolarsResult<DataFrame> {
905        let rows = rows.min(self.rows);
906        let columns = self
907            .columns
908            .iter()
909            .map(|c| self.decode_column(c, rows))
910            .collect::<PolarsResult<Vec<_>>>()?;
911        DataFrame::new(rows, columns)
912    }
913
914    fn decode_column(&self, column: &ColumnLayout, rows: usize) -> PolarsResult<Column> {
915        let source = &self.sources[column.source];
916        source.still_whole()?;
917        decode(source.as_slice(), column, rows)
918    }
919}
920
921impl RowSource for FixedRecords {
922    fn height(&self) -> usize {
923        self.rows
924    }
925
926    fn schema(&self) -> SchemaRef {
927        self.schema.clone()
928    }
929
930    fn decode(&self, column: usize, index: &IdxCa) -> PolarsResult<Column> {
931        let column = &self.columns[column];
932        let source = &self.sources[column.source];
933        source.still_whole()?;
934        decode_rows(source.as_slice(), column, index)
935    }
936}
937
938impl crate::formats::pushdown::Windowed for FixedRecords {
939    fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
940        Ok(FixedRecords::window(self, start, len)?.lazy())
941    }
942}
943
944#[cfg(test)]
945mod tests;