format: bound and batch every chunk fetch over Storage
Only the full read split its chunk fetches into 64 MiB batches. The selection path, the indexed read and the parallel_read decoders fetched every chunk's stored bytes in one read_ranges call, each extent bounded only by the file length, so a crafted chunk index pointing many chunks at one large extent made File::open_storage hold chunks x extent bytes (3.3 GB from a 16.8 MB file) before the first decode error. - storage::for_each_extent_batch is now the one way raw-data reads fetch chunk bytes: batches of at most RAW_BATCH_BYTES (now pub), each decoded before the next is fetched. Used by the full, cached, indexed, selection and parallel_read paths; the sweep read uses read_extent per chunk. - ExtentReq carries each chunk's claimed extent (bounds-checked as before, same errors) and the prefix actually fetched: filters::stored_chunk_limit — the chunk size if unfiltered, else each applied filter's worst-case growth (n + n/4 + 4096 per codec; unbounded only for an application-registered codec). The in-memory path cuts the slice it decodes the same way, so both paths still agree. - tests/raw_fetch_bounds.rs: a crafted chunked_large.h5 (ten chunks all claiming 20 MiB at one padding blob) read through every path over a storage that records the largest single fetch; and 160 MiB of legitimate unfiltered chunks fetched batch by batch. Before: one 80 MiB fetch (selection) and one 160 MiB fetch; after: within the budget. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
This commit is contained in:
@@ -12,24 +12,21 @@ 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};
|
||||
use crate::storage::{ExtentReq, Storage, for_each_extent_batch};
|
||||
|
||||
/// 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,
|
||||
/// The extents of `chunks`' stored bytes (see
|
||||
/// [`crate::chunked_read::chunk_req`]), fetched batch by batch with
|
||||
/// [`for_each_extent_batch`] when the file is not in memory (each chunk's
|
||||
/// bounds error is reported when that chunk is decoded, as before).
|
||||
fn chunk_reqs(
|
||||
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)
|
||||
pipeline: Option<&FilterPipeline>,
|
||||
chunk_total_bytes: usize,
|
||||
) -> Vec<ExtentReq> {
|
||||
chunks
|
||||
.iter()
|
||||
.map(|c| crate::chunked_read::chunk_req(c, pipeline, chunk_total_bytes, true))
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Threshold: only use parallel decompression when chunk count exceeds this.
|
||||
@@ -238,62 +235,75 @@ pub fn decompress_chunks_lane_partitioned_in<S: Storage + ?Sized>(
|
||||
.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();
|
||||
|
||||
// Each lane processes its assigned chunks and returns results + stats.
|
||||
let lane_results: Result<Vec<(Vec<DecompressedChunk>, LaneStats)>, FormatError> = assignments
|
||||
.into_par_iter()
|
||||
.map(|indices| {
|
||||
let mut results = Vec::with_capacity(indices.len());
|
||||
let mut stats = LaneStats::default();
|
||||
|
||||
for &index in &indices {
|
||||
let chunk_info = &chunks[index];
|
||||
let size = chunk_info.chunk_size as usize;
|
||||
let raw_chunk = raw_bytes.get(index, chunk_info.address, size)?;
|
||||
|
||||
let decompressed = decompress_chunk_exact(
|
||||
raw_chunk,
|
||||
pipeline,
|
||||
chunk_total_bytes,
|
||||
element_size,
|
||||
chunk_info.filter_mask,
|
||||
&chunk_info.offsets,
|
||||
)?;
|
||||
|
||||
stats.chunks_processed += 1;
|
||||
stats.compressed_bytes += size as u64;
|
||||
stats.decompressed_bytes += decompressed.len() as u64;
|
||||
|
||||
results.push(DecompressedChunk {
|
||||
index,
|
||||
data: decompressed,
|
||||
});
|
||||
}
|
||||
|
||||
Ok((results, stats))
|
||||
})
|
||||
.collect();
|
||||
|
||||
let lane_results = lane_results?;
|
||||
|
||||
// Aggregate stats
|
||||
let mut partition_stats = PartitionStats::new(num_lanes);
|
||||
let reqs = chunk_reqs(chunks, Some(pipeline), chunk_total_bytes);
|
||||
let mut ordered: Vec<Vec<u8>> = Vec::with_capacity(chunks.len());
|
||||
let mut partition_stats = PartitionStats::new(0);
|
||||
partition_stats.total_chunks = chunks.len();
|
||||
for (lane_idx, (_, stats)) in lane_results.iter().enumerate() {
|
||||
partition_stats.per_lane[lane_idx] = stats.clone();
|
||||
}
|
||||
// Each batch of fetched chunks is partitioned into lanes and decoded
|
||||
// before the next batch is fetched (with the file in memory there is
|
||||
// one batch: all the chunks).
|
||||
for_each_extent_batch(file_data, &reqs, |batch, raw_bytes| {
|
||||
let assignments = lane_partition::partition_chunks(batch.len(), lanes, seed);
|
||||
// Each lane processes its assigned chunks and returns results + stats.
|
||||
let lane_results: Result<Vec<(Vec<DecompressedChunk>, LaneStats)>, FormatError> =
|
||||
assignments
|
||||
.into_par_iter()
|
||||
.map(|indices| {
|
||||
let mut results = Vec::with_capacity(indices.len());
|
||||
let mut stats = LaneStats::default();
|
||||
|
||||
// Flatten and sort by original index to restore order
|
||||
let mut all_chunks: Vec<DecompressedChunk> = lane_results
|
||||
.into_iter()
|
||||
.flat_map(|(chunks, _)| chunks)
|
||||
.collect();
|
||||
all_chunks.sort_by_key(|dc| dc.index);
|
||||
for &local in &indices {
|
||||
let index = batch.start + local;
|
||||
let chunk_info = &chunks[index];
|
||||
let size = chunk_info.chunk_size as usize;
|
||||
let raw_chunk = raw_bytes.get(index, &reqs[index])?;
|
||||
|
||||
let ordered = all_chunks.into_iter().map(|dc| dc.data).collect();
|
||||
let decompressed = decompress_chunk_exact(
|
||||
raw_chunk,
|
||||
pipeline,
|
||||
chunk_total_bytes,
|
||||
element_size,
|
||||
chunk_info.filter_mask,
|
||||
&chunk_info.offsets,
|
||||
)?;
|
||||
|
||||
stats.chunks_processed += 1;
|
||||
stats.compressed_bytes += size as u64;
|
||||
stats.decompressed_bytes += decompressed.len() as u64;
|
||||
|
||||
results.push(DecompressedChunk {
|
||||
index,
|
||||
data: decompressed,
|
||||
});
|
||||
}
|
||||
|
||||
Ok((results, stats))
|
||||
})
|
||||
.collect();
|
||||
let lane_results = lane_results?;
|
||||
|
||||
// Aggregate stats
|
||||
if partition_stats.per_lane.len() < lane_results.len() {
|
||||
partition_stats
|
||||
.per_lane
|
||||
.resize_with(lane_results.len(), LaneStats::default);
|
||||
partition_stats.num_lanes = lane_results.len();
|
||||
}
|
||||
for (lane, (_, stats)) in partition_stats.per_lane.iter_mut().zip(&lane_results) {
|
||||
lane.chunks_processed += stats.chunks_processed;
|
||||
lane.compressed_bytes += stats.compressed_bytes;
|
||||
lane.decompressed_bytes += stats.decompressed_bytes;
|
||||
}
|
||||
|
||||
// Flatten and sort by original index to restore order
|
||||
let mut all_chunks: Vec<DecompressedChunk> = lane_results
|
||||
.into_iter()
|
||||
.flat_map(|(chunks, _)| chunks)
|
||||
.collect();
|
||||
all_chunks.sort_by_key(|dc| dc.index);
|
||||
ordered.extend(all_chunks.into_iter().map(|dc| dc.data));
|
||||
Ok(())
|
||||
})?;
|
||||
Ok((ordered, partition_stats))
|
||||
}
|
||||
|
||||
@@ -325,33 +335,38 @@ pub fn decompress_chunks_parallel_in<S: Storage + ?Sized>(
|
||||
) -> 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 size = chunk_info.chunk_size as usize;
|
||||
let raw_chunk = raw_bytes.get(index, chunk_info.address, size)?;
|
||||
let reqs = chunk_reqs(chunks, Some(pipeline), chunk_total_bytes);
|
||||
let mut ordered: Vec<Vec<u8>> = Vec::with_capacity(chunks.len());
|
||||
for_each_extent_batch(file_data, &reqs, |batch, raw_bytes| {
|
||||
let results: Result<Vec<DecompressedChunk>, FormatError> = batch
|
||||
.clone()
|
||||
.into_par_iter()
|
||||
.map(|index| {
|
||||
let chunk_info = &chunks[index];
|
||||
let raw_chunk = raw_bytes.get(index, &reqs[index])?;
|
||||
|
||||
let decompressed = decompress_chunk_exact(
|
||||
raw_chunk,
|
||||
pipeline,
|
||||
chunk_total_bytes,
|
||||
element_size,
|
||||
chunk_info.filter_mask,
|
||||
&chunk_info.offsets,
|
||||
)?;
|
||||
let decompressed = decompress_chunk_exact(
|
||||
raw_chunk,
|
||||
pipeline,
|
||||
chunk_total_bytes,
|
||||
element_size,
|
||||
chunk_info.filter_mask,
|
||||
&chunk_info.offsets,
|
||||
)?;
|
||||
|
||||
Ok(DecompressedChunk {
|
||||
index,
|
||||
data: decompressed,
|
||||
Ok(DecompressedChunk {
|
||||
index,
|
||||
data: decompressed,
|
||||
})
|
||||
})
|
||||
})
|
||||
.collect();
|
||||
.collect();
|
||||
|
||||
let mut result_vec = results?;
|
||||
result_vec.sort_by_key(|dc| dc.index);
|
||||
Ok(result_vec.into_iter().map(|dc| dc.data).collect())
|
||||
let mut result_vec = results?;
|
||||
result_vec.sort_by_key(|dc| dc.index);
|
||||
ordered.extend(result_vec.into_iter().map(|dc| dc.data));
|
||||
Ok(())
|
||||
})?;
|
||||
Ok(ordered)
|
||||
}
|
||||
|
||||
/// Decompress chunks sequentially (fallback when parallel is not warranted).
|
||||
@@ -373,26 +388,29 @@ pub fn decompress_chunks_sequential_in<S: Storage + ?Sized>(
|
||||
chunk_total_bytes: usize,
|
||||
element_size: u32,
|
||||
) -> Result<Vec<Vec<u8>>, FormatError> {
|
||||
let raw_bytes = fetch_all(file_data, chunks)?;
|
||||
let reqs = chunk_reqs(chunks, pipeline, chunk_total_bytes);
|
||||
let mut result = Vec::with_capacity(chunks.len());
|
||||
for (i, chunk_info) in chunks.iter().enumerate() {
|
||||
let size = chunk_info.chunk_size as usize;
|
||||
let raw_chunk = raw_bytes.get(i, chunk_info.address, size)?;
|
||||
for_each_extent_batch(file_data, &reqs, |batch, raw_bytes| {
|
||||
for i in batch {
|
||||
let chunk_info = &chunks[i];
|
||||
let raw_chunk = raw_bytes.get(i, &reqs[i])?;
|
||||
|
||||
let decompressed = if let Some(pl) = pipeline {
|
||||
decompress_chunk_exact(
|
||||
raw_chunk,
|
||||
pl,
|
||||
chunk_total_bytes,
|
||||
element_size,
|
||||
chunk_info.filter_mask,
|
||||
&chunk_info.offsets,
|
||||
)?
|
||||
} else {
|
||||
raw_chunk.to_vec()
|
||||
};
|
||||
result.push(decompressed);
|
||||
}
|
||||
let decompressed = if let Some(pl) = pipeline {
|
||||
decompress_chunk_exact(
|
||||
raw_chunk,
|
||||
pl,
|
||||
chunk_total_bytes,
|
||||
element_size,
|
||||
chunk_info.filter_mask,
|
||||
&chunk_info.offsets,
|
||||
)?
|
||||
} else {
|
||||
raw_chunk.to_vec()
|
||||
};
|
||||
result.push(decompressed);
|
||||
}
|
||||
Ok(())
|
||||
})?;
|
||||
Ok(result)
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user