format: raw data, VDS and VL data over Storage
Every raw-data path has a generic *_in core, with the &[u8] functions as thin wrappers: data_read (read_raw_data*, read_raw_data_selection, read_chunked_native), chunked_read (the v1 B-tree chunk index, list_chunks, the full, cached, sweep and indexed reads), parallel_read, partial_read, fill_value (read_full_with_fill, apply_to_unallocated_chunks; and dataset_fill_value_from_storage is now generic), vds (the virtual file through Storage, external sources still through the resolver), vl_data (VlResolver<'a, S = [u8]>, read_vl_strings_in, read_vl_bytes_in), AttributeMessage::read_vl_strings_in and provenance::verify_dataset_in. With the whole file in memory nothing changes: chunks and contiguous data are sliced from it as before. Otherwise a chunked read lists its chunks, fetches their stored bytes with one Storage::read_ranges call per 64 MiB batch (chunks the cache already holds are not fetched), then decodes as today; a selection fetches only the chunks it overlaps, and a contiguous selection only its runs. Each extent's bounds error is the one the slice code gave, reported when that extent is reached, so errors keep their order. Tests: the equivalence harness now reads every dataset's values (whole, fill-aware, cached, indexed, three selections, VDS, VL strings and sequences) through the read_at-only storage and requires the slice results (all 653 corpus files agree); a misbehaving storage (a failing Nth read, short reads) only ever yields errors or the right values; and chunked reads are checked to use one read_ranges call. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
This commit is contained in:
@@ -7,12 +7,30 @@
|
||||
//! The lane assignment is seeded by dataset metadata so repeated reads of
|
||||
//! the same region produce identical partitions (cache-friendly, reproducible).
|
||||
|
||||
use crate::addr::to_usize;
|
||||
use crate::chunked_read::ChunkInfo;
|
||||
use crate::error::FormatError;
|
||||
use crate::filter_pipeline::FilterPipeline;
|
||||
use crate::filters::decompress_chunk_exact;
|
||||
use crate::lane_partition::{self, LaneStats, PartitionStats};
|
||||
use crate::storage::{ExtentBytes, Storage};
|
||||
|
||||
/// The stored bytes of every chunk in `chunks`, fetched in one
|
||||
/// [`Storage::read_ranges`] call when the file is not in memory (each
|
||||
/// chunk's bounds error is reported when that chunk is decoded, as before).
|
||||
fn fetch_all<'a, S: Storage + ?Sized>(
|
||||
file_data: &'a S,
|
||||
chunks: &[ChunkInfo],
|
||||
) -> Result<ExtentBytes<'a>, FormatError> {
|
||||
let extents: Vec<(u64, usize, bool)> = if file_data.as_contiguous().is_some() {
|
||||
Vec::new()
|
||||
} else {
|
||||
chunks
|
||||
.iter()
|
||||
.map(|c| (c.address, c.chunk_size as usize, true))
|
||||
.collect()
|
||||
};
|
||||
ExtentBytes::fetch(file_data, &extents)
|
||||
}
|
||||
|
||||
/// Threshold: only use parallel decompression when chunk count exceeds this.
|
||||
const PARALLEL_THRESHOLD: usize = 4;
|
||||
@@ -190,6 +208,27 @@ pub fn decompress_chunks_lane_partitioned(
|
||||
element_size: u32,
|
||||
seed: u64,
|
||||
num_lanes: Option<usize>,
|
||||
) -> Result<(Vec<Vec<u8>>, PartitionStats), FormatError> {
|
||||
decompress_chunks_lane_partitioned_in(
|
||||
file_data,
|
||||
chunks,
|
||||
pipeline,
|
||||
chunk_total_bytes,
|
||||
element_size,
|
||||
seed,
|
||||
num_lanes,
|
||||
)
|
||||
}
|
||||
|
||||
/// [`decompress_chunks_lane_partitioned`] over any [`Storage`].
|
||||
pub fn decompress_chunks_lane_partitioned_in<S: Storage + ?Sized>(
|
||||
file_data: &S,
|
||||
chunks: &[ChunkInfo],
|
||||
pipeline: &FilterPipeline,
|
||||
chunk_total_bytes: usize,
|
||||
element_size: u32,
|
||||
seed: u64,
|
||||
num_lanes: Option<usize>,
|
||||
) -> Result<(Vec<Vec<u8>>, PartitionStats), FormatError> {
|
||||
use rayon::prelude::*;
|
||||
|
||||
@@ -199,6 +238,7 @@ pub fn decompress_chunks_lane_partitioned(
|
||||
.unwrap_or(1)
|
||||
});
|
||||
|
||||
let raw_bytes = fetch_all(file_data, chunks)?;
|
||||
let assignments = lane_partition::partition_chunks(chunks.len(), lanes, seed);
|
||||
let num_lanes = assignments.len();
|
||||
|
||||
@@ -211,19 +251,8 @@ pub fn decompress_chunks_lane_partitioned(
|
||||
|
||||
for &index in &indices {
|
||||
let chunk_info = &chunks[index];
|
||||
let c_addr = to_usize(chunk_info.address)?;
|
||||
let size = chunk_info.chunk_size as usize;
|
||||
|
||||
if c_addr
|
||||
.checked_add(size)
|
||||
.is_none_or(|end| end > file_data.len())
|
||||
{
|
||||
return Err(FormatError::UnexpectedEof {
|
||||
expected: c_addr.saturating_add(size),
|
||||
available: file_data.len(),
|
||||
});
|
||||
}
|
||||
let raw_chunk = &file_data[c_addr..c_addr + size];
|
||||
let raw_chunk = raw_bytes.get(index, chunk_info.address, size)?;
|
||||
|
||||
let decompressed = decompress_chunk_exact(
|
||||
raw_chunk,
|
||||
@@ -282,25 +311,27 @@ pub fn decompress_chunks_parallel(
|
||||
pipeline: &FilterPipeline,
|
||||
chunk_total_bytes: usize,
|
||||
element_size: u32,
|
||||
) -> Result<Vec<Vec<u8>>, FormatError> {
|
||||
decompress_chunks_parallel_in(file_data, chunks, pipeline, chunk_total_bytes, element_size)
|
||||
}
|
||||
|
||||
/// [`decompress_chunks_parallel`] over any [`Storage`].
|
||||
pub fn decompress_chunks_parallel_in<S: Storage + ?Sized>(
|
||||
file_data: &S,
|
||||
chunks: &[ChunkInfo],
|
||||
pipeline: &FilterPipeline,
|
||||
chunk_total_bytes: usize,
|
||||
element_size: u32,
|
||||
) -> Result<Vec<Vec<u8>>, FormatError> {
|
||||
use rayon::prelude::*;
|
||||
|
||||
let raw_bytes = fetch_all(file_data, chunks)?;
|
||||
let results: Result<Vec<DecompressedChunk>, FormatError> = chunks
|
||||
.par_iter()
|
||||
.enumerate()
|
||||
.map(|(index, chunk_info)| {
|
||||
let c_addr = to_usize(chunk_info.address)?;
|
||||
let size = chunk_info.chunk_size as usize;
|
||||
if c_addr
|
||||
.checked_add(size)
|
||||
.is_none_or(|end| end > file_data.len())
|
||||
{
|
||||
return Err(FormatError::UnexpectedEof {
|
||||
expected: c_addr.saturating_add(size),
|
||||
available: file_data.len(),
|
||||
});
|
||||
}
|
||||
let raw_chunk = &file_data[c_addr..c_addr + size];
|
||||
let raw_chunk = raw_bytes.get(index, chunk_info.address, size)?;
|
||||
|
||||
let decompressed = decompress_chunk_exact(
|
||||
raw_chunk,
|
||||
@@ -331,20 +362,22 @@ pub fn decompress_chunks_sequential(
|
||||
chunk_total_bytes: usize,
|
||||
element_size: u32,
|
||||
) -> Result<Vec<Vec<u8>>, FormatError> {
|
||||
decompress_chunks_sequential_in(file_data, chunks, pipeline, chunk_total_bytes, element_size)
|
||||
}
|
||||
|
||||
/// [`decompress_chunks_sequential`] over any [`Storage`].
|
||||
pub fn decompress_chunks_sequential_in<S: Storage + ?Sized>(
|
||||
file_data: &S,
|
||||
chunks: &[ChunkInfo],
|
||||
pipeline: Option<&FilterPipeline>,
|
||||
chunk_total_bytes: usize,
|
||||
element_size: u32,
|
||||
) -> Result<Vec<Vec<u8>>, FormatError> {
|
||||
let raw_bytes = fetch_all(file_data, chunks)?;
|
||||
let mut result = Vec::with_capacity(chunks.len());
|
||||
for chunk_info in chunks {
|
||||
let c_addr = to_usize(chunk_info.address)?;
|
||||
for (i, chunk_info) in chunks.iter().enumerate() {
|
||||
let size = chunk_info.chunk_size as usize;
|
||||
if c_addr
|
||||
.checked_add(size)
|
||||
.is_none_or(|end| end > file_data.len())
|
||||
{
|
||||
return Err(FormatError::UnexpectedEof {
|
||||
expected: c_addr.saturating_add(size),
|
||||
available: file_data.len(),
|
||||
});
|
||||
}
|
||||
let raw_chunk = &file_data[c_addr..c_addr + size];
|
||||
let raw_chunk = raw_bytes.get(i, chunk_info.address, size)?;
|
||||
|
||||
let decompressed = if let Some(pl) = pipeline {
|
||||
decompress_chunk_exact(
|
||||
|
||||
Reference in New Issue
Block a user