use core::num::NonZeroUsize;
use crate::checksum::jenkins_lookup3;
use crate::chunked_write::{ea_compute_stats, split_into_chunks, write_ea_addr};
use crate::convert::TryToUsize;
use crate::data_layout::DataLayout;
use crate::dataspace::Dataspace;
use crate::datatype::Datatype;
use crate::edit::{LOSSY_TAIL_REFUSAL, pipeline_lossless};
use crate::error::{Error, FormatError};
use crate::extensible_array::{DataBlockGeom, EaGeometry, ExtensibleArrayHeader, SuperBlockGeom};
use crate::fill_value::FillPattern;
use crate::filter_pipeline::FilterPipeline;
use crate::filters::{ChunkContext, FilterScratch, compress_chunk_with, decompress_chunk};
use crate::message_type::MessageType;
use crate::source::Source;
#[cfg(test)]
pub(crate) mod alloc_probe {
use core::cell::Cell;
thread_local! {
static DATA_BLOCKS: Cell<usize> = const { Cell::new(0) };
static SUPER_BLOCKS: Cell<usize> = const { Cell::new(0) };
}
pub(crate) fn note_data_block() {
DATA_BLOCKS.with(|c| c.set(c.get() + 1));
}
pub(crate) fn note_super_block() {
SUPER_BLOCKS.with(|c| c.set(c.get() + 1));
}
pub(crate) fn take() -> (usize, usize) {
(
DATA_BLOCKS.with(|c| c.replace(0)),
SUPER_BLOCKS.with(|c| c.replace(0)),
)
}
}
pub(crate) fn undef_addr(offset_size: u8) -> u64 {
match offset_size {
4 => 0xFFFF_FFFF,
_ => u64::MAX,
}
}
pub(crate) fn is_undef(addr: u64, offset_size: u8) -> bool {
addr == undef_addr(offset_size)
}
fn push_undef_element(buf: &mut Vec<u8>, offset_size: u8, ea_elem_size: usize) {
write_ea_addr(buf, undef_addr(offset_size), offset_size);
for _ in offset_size as usize..ea_elem_size {
buf.push(0);
}
}
#[derive(Clone, Copy, Debug)]
pub(crate) struct ElemRecord {
pub addr: u64,
pub stored_size: u64,
pub filter_mask: u32,
}
impl ElemRecord {
pub(crate) fn addr_only(addr: u64) -> Self {
Self {
addr,
stored_size: 0,
filter_mask: 0,
}
}
}
pub(crate) trait Store: Source {
fn offset_size(&self) -> u8;
fn length_size(&self) -> u8;
fn alloc_raw(&mut self, bytes: &[u8]) -> Result<u64, Error>;
fn write_at(&mut self, offset: u64, bytes: &[u8]) -> Result<(), Error>;
fn patch_superblock_eof(&mut self) -> Result<(), Error>;
fn sync(&mut self) -> Result<(), Error>;
fn read_addr_at(&self, offset: u64) -> Result<u64, Error> {
let mut buf = [0u8; 8];
let width = if self.offset_size() == 4 { 4 } else { 8 };
self.read_at(offset, &mut buf[..width])?;
Ok(u64::from_le_bytes(buf))
}
fn publish_checksummed(
&mut self,
start: u64,
cks_off: u64,
at: u64,
value: &[u8],
) -> Result<(), Error> {
let span = (cks_off + 4 - start).to_usize()?;
let from = at
.checked_sub(start)
.and_then(|d| d.to_usize().ok())
.filter(|&f| f.checked_add(value.len()).is_some_and(|e| e <= span - 4))
.ok_or(Error::AppendUnsupported(
"a checksummed structure was patched outside itself",
))?;
let mut bytes = self.read_exact_at(start, span)?;
bytes[from..from + value.len()].copy_from_slice(value);
let (body, cks_field) = bytes.split_at_mut(span - 4);
cks_field.copy_from_slice(&jenkins_lookup3(body).to_le_bytes());
self.write_at(at, &bytes[from..])
}
}
const MAX_EA_ELEM: usize = 8 + 8 + 4;
pub(crate) struct MessageSpans {
pub datatype: (u64, usize),
pub filter: Option<(u64, usize)>,
pub fill: Option<(MessageType, u64, usize)>,
}
pub(crate) struct LocateResult {
pub located: Located,
pub spans: MessageSpans,
}
pub(crate) struct Located {
pub dim0_off: u64,
pub current_dim: u64,
pub ohdr_chunk_start: u64,
pub ohdr_chunk_msg_end: u64,
pub chunk_elems: u64,
pub elem_bytes: NonZeroUsize,
pub chunk_bytes: usize,
pub client_id: u8,
pub ea_addr: u64,
pub geom: EaGeometry,
pub idx_blk_elmts: u64,
pub ea_elem_size: usize,
pub page_nelmts: u64,
pub blk_off_size: usize,
pub index_block_addr: u64,
pub num_chunks: u64,
}
impl Located {
pub(crate) fn locate_at<F: Store>(
file: &F,
oh_addr: u64,
unsupported: fn(&'static str) -> Error,
) -> Result<LocateResult, Error> {
let os = file.offset_size();
let ls = file.length_size();
let walk = walk_v2_object_header(file, oh_addr, os, ls)?;
let dataspace_msg = walk
.messages
.iter()
.find(|m| m.msg_type == MessageType::Dataspace)
.ok_or(Error::MissingMessage(MessageType::Dataspace))?;
let layout_msg = walk
.messages
.iter()
.find(|m| m.msg_type == MessageType::DataLayout)
.ok_or(Error::MissingMessage(MessageType::DataLayout))?;
let datatype_msg = walk
.messages
.iter()
.find(|m| m.msg_type == MessageType::Datatype)
.ok_or(Error::MissingMessage(MessageType::Datatype))?;
let filter_msg = walk
.messages
.iter()
.find(|m| m.msg_type == MessageType::FilterPipeline);
let fill_msg = walk
.messages
.iter()
.find(|m| m.msg_type == MessageType::FillValue)
.or_else(|| {
walk.messages
.iter()
.find(|m| m.msg_type == MessageType::FillValueOld)
});
for msg in [Some(dataspace_msg), Some(datatype_msg), filter_msg]
.into_iter()
.flatten()
{
if crate::shared_message::is_shared(msg.flags) {
return Err(unsupported(
"dataset has a committed (shared) datatype, dataspace, or filter pipeline",
));
}
}
let ds_bytes = file.read_metadata_at(dataspace_msg.data_off, dataspace_msg.size)?;
let dataspace = Dataspace::parse(&ds_bytes, ls)?;
if dataspace.rank != 1 {
return Err(unsupported("only rank-1 datasets are supported"));
}
match &dataspace.max_dimensions {
Some(maxs) if maxs.first() == Some(&u64::MAX) => {}
_ => {
return Err(unsupported("dataset has no unlimited (maxshape) dimension"));
}
}
let current_dim = dataspace.dimensions[0];
let dim0_off = dataspace_msg.data_off + 4;
let layout_bytes = file.read_metadata_at(layout_msg.data_off, layout_msg.size)?;
let layout = DataLayout::parse(&layout_bytes, os, ls)?;
let (ea_addr, chunk_dims) = match layout {
DataLayout::Chunked {
chunk_index_type: Some(4),
btree_address: Some(addr),
chunk_dimensions,
..
} => (addr, chunk_dimensions),
DataLayout::Chunked {
chunk_index_type: Some(4),
btree_address: None,
..
} => {
return Err(unsupported(
"the dataset's extensible-array index is not allocated yet (an empty \
dataset with no chunks); write initial data at creation or make the \
first append with Dataset::append_staged",
));
}
_ => {
return Err(unsupported(
"only Extensible-Array-indexed chunked datasets are supported",
));
}
};
if chunk_dims.len() != 2 {
return Err(unsupported(
"unexpected chunk dimensionality (expected a rank-1 chunked layout)",
));
}
let chunk_elems = chunk_dims[0] as u64;
let Some(elem_bytes) = NonZeroUsize::new(chunk_dims[1] as usize) else {
return Err(unsupported("dataset has a zero-sized element"));
};
let chunk_bytes = chunk_elems.to_usize()? * elem_bytes.get();
let ea_header = ExtensibleArrayHeader::parse_from_source(file, ea_addr, os, ls)?;
let has_filters = filter_msg.is_some();
if (ea_header.client_id == 1) != has_filters {
return Err(unsupported(
"dataset filter metadata is inconsistent (chunk-index client id \
disagrees with the filter pipeline)",
));
}
let elem_w = ea_header.element_size as usize;
if ea_header.client_id == 1 && !(os as usize + 5..=os as usize + 12).contains(&elem_w) {
return Err(unsupported(
"malformed filtered extensible-array element width",
));
}
let geom = EaGeometry::from_header(&ea_header);
let page_nelmts = 1u64 << ea_header.max_dblk_nelmts_bits;
let blk_off_size = (ea_header.max_nelmts_bits as usize).div_ceil(8);
let index_block_addr = ea_header.index_block_address;
let num_chunks = if chunk_elems == 0 {
ea_header.num_elements
} else {
current_dim.div_ceil(chunk_elems)
};
Ok(LocateResult {
located: Located {
dim0_off,
current_dim,
ohdr_chunk_start: dataspace_msg.chunk_start,
ohdr_chunk_msg_end: dataspace_msg.chunk_msg_end,
chunk_elems,
elem_bytes,
chunk_bytes,
client_id: ea_header.client_id,
ea_addr,
geom,
idx_blk_elmts: ea_header.idx_blk_elmts as u64,
ea_elem_size: ea_header.element_size as usize,
page_nelmts,
blk_off_size,
index_block_addr,
num_chunks,
},
spans: MessageSpans {
datatype: (datatype_msg.data_off, datatype_msg.size),
filter: filter_msg.map(|m| (m.data_off, m.size)),
fill: fill_msg.map(|m| (m.msg_type, m.data_off, m.size)),
},
})
}
fn read_element_at<F: Store>(&self, file: &F, off: u64) -> Result<ElemRecord, Error> {
let os = file.offset_size() as usize;
if self.client_id == 0 {
return Ok(ElemRecord::addr_only(file.read_addr_at(off)?));
}
let elem = file.read_exact_at(off, self.ea_elem_size)?;
let addr_w = if file.offset_size() == 4 { 4 } else { 8 };
let mut a = [0u8; 8];
a[..addr_w].copy_from_slice(&elem[..addr_w]);
let addr = u64::from_le_bytes(a);
let csz = self.ea_elem_size - os - 4;
let mut stored_size = 0u64;
for i in 0..csz {
stored_size |= (elem[os + i] as u64) << (8 * i);
}
let fm = os + csz;
let filter_mask = u32::from_le_bytes([elem[fm], elem[fm + 1], elem[fm + 2], elem[fm + 3]]);
Ok(ElemRecord {
addr,
stored_size,
filter_mask,
})
}
fn elem_slot_off<F: Store>(&self, file: &F, e: u64) -> Result<Option<u64>, Error> {
let os = file.offset_size() as usize;
let elem_size = self.ea_elem_size as u64;
let idx = self.idx_blk_elmts;
let blk_off = self.blk_off_size;
let page_nelmts = self.page_nelmts;
if e < idx {
let ib_prefix = (4 + 1 + 1 + os) as u64;
return Ok(Some(self.index_block_addr + ib_prefix + e * elem_size));
}
let region = locate_data_block(&self.geom, idx, e);
if region.ndblks == 0 {
return Err(Error::AppendUnsupported(
"chunk index geometry does not cover the appended element",
));
}
let is_paged = region.blocks(page_nelmts).is_paged();
let slot = e - region.db_start;
let sblk_addr = match region.parent {
Parent::Super { sblk_j, .. } => {
let a = self.super_block_addr(file, sblk_j)?;
if is_undef(a, file.offset_size()) {
return Ok(None);
}
Some(a)
}
Parent::IndexDirect { .. } => None,
};
let dblk_ptr_off = self.dblk_ptr_off(sblk_addr, ®ion, os, blk_off)?;
let dblk_addr = file.read_addr_at(dblk_ptr_off)?;
if is_undef(dblk_addr, file.offset_size()) {
return Ok(None);
}
let off = if !is_paged {
let db_prefix = (4 + 1 + 1 + os + blk_off) as u64;
dblk_addr + db_prefix + slot * elem_size
} else {
let header_size = (4 + 1 + 1 + os + blk_off + 4) as u64;
let page = slot / page_nelmts;
let slot_in_page = slot % page_nelmts;
let page_bytes = page_nelmts * elem_size + 4;
let page_off = dblk_addr + header_size + page * page_bytes;
page_off + slot_in_page * elem_size
};
Ok(Some(off))
}
fn dblk_ptr_off(
&self,
sblk_addr: Option<u64>,
region: &DataBlockLoc,
os: usize,
blk_off: usize,
) -> Result<u64, Error> {
match region.parent {
Parent::IndexDirect { ordinal } => {
let ib_prefix = (4 + 1 + 1 + os) as u64;
Ok(self.index_block_addr
+ ib_prefix
+ self.idx_blk_elmts * self.ea_elem_size as u64
+ (ordinal * os) as u64)
}
Parent::Super { dblk_local, .. } => {
let sblk_addr = sblk_addr.expect("super-block address resolved for a Super parent");
Ok(sb_dblk_slot_off(
os,
sblk_addr,
dblk_local,
region.super_block(self.page_nelmts),
blk_off,
))
}
}
}
pub(crate) fn read_element<F: Store>(
&self,
file: &F,
e: u64,
) -> Result<Option<ElemRecord>, Error> {
match self.elem_slot_off(file, e)? {
None => Ok(None),
Some(off) => {
let rec = self.read_element_at(file, off)?;
if is_undef(rec.addr, file.offset_size()) {
Ok(None)
} else {
Ok(Some(rec))
}
}
}
}
pub(crate) fn ea_insert<F: Store>(
&self,
file: &mut F,
e: u64,
rec: ElemRecord,
) -> Result<(), Error> {
let os = file.offset_size() as usize;
let elem_size = self.ea_elem_size as u64;
let idx = self.idx_blk_elmts;
let blk_off = self.blk_off_size;
if e < idx {
let ib_prefix = (4 + 1 + 1 + os) as u64;
let slot_off = self.index_block_addr + ib_prefix + e * elem_size;
let mut buf = [0u8; MAX_EA_ELEM];
let n = self.element_bytes(&mut buf, os, rec)?;
return self.publish_index_block(file, slot_off, &buf[..n]);
}
let region = locate_data_block(&self.geom, idx, e);
if region.ndblks == 0 {
return Err(Error::AppendUnsupported(
"chunk index geometry does not cover the appended element",
));
}
let dblk_nelmts = region.dblk_nelmts;
let sb_geom = region.super_block(self.page_nelmts);
let is_paged = sb_geom.blocks.is_paged();
let slot = e - region.db_start;
let block_offset_rel = region.db_start - idx;
let sblk_addr = match region.parent {
Parent::Super { sblk_j, .. } => {
Some(self.ensure_super_block(file, sblk_j, region.sb_block_offset, sb_geom)?)
}
Parent::IndexDirect { .. } => None,
};
let dblk_ptr_off = self.dblk_ptr_off(sblk_addr, ®ion, os, blk_off)?;
let existing = file.read_addr_at(dblk_ptr_off)?;
let dblk_addr = if is_undef(existing, file.offset_size()) {
let new_addr = if is_paged {
self.alloc_undef_paged_data_block(file, dblk_nelmts, block_offset_rel)?
} else {
self.alloc_undef_data_block(file, dblk_nelmts, block_offset_rel)?
};
#[cfg(test)]
alloc_probe::note_data_block();
file.sync()?;
match region.parent {
Parent::IndexDirect { .. } => {
self.publish_index_block(file, dblk_ptr_off, &new_addr.to_le_bytes()[..os])?;
}
Parent::Super { .. } => self.publish_super_block(
file,
sblk_addr.unwrap(),
sb_geom,
blk_off,
dblk_ptr_off,
&new_addr.to_le_bytes()[..os],
)?,
}
new_addr
} else {
existing
};
if !is_paged {
let db_prefix = (4 + 1 + 1 + os + blk_off) as u64;
let elem_off = dblk_addr + db_prefix + slot * elem_size;
let mut buf = [0u8; MAX_EA_ELEM];
let n = self.element_bytes(&mut buf, os, rec)?;
let cks_off = dblk_addr + db_prefix + dblk_nelmts * elem_size;
file.publish_checksummed(dblk_addr, cks_off, elem_off, &buf[..n])?;
} else {
let page_nelmts = self.page_nelmts;
let header_size = (4 + 1 + 1 + os + blk_off + 4) as u64;
let page = slot / page_nelmts;
let slot_in_page = slot % page_nelmts;
let page_bytes = page_nelmts * elem_size + 4;
let page_off = dblk_addr + header_size + page * page_bytes;
let elem_off = page_off + slot_in_page * elem_size;
let mut buf = [0u8; MAX_EA_ELEM];
let n = self.element_bytes(&mut buf, os, rec)?;
let page_cks_off = page_off + page_nelmts * elem_size;
file.publish_checksummed(page_off, page_cks_off, elem_off, &buf[..n])?;
if slot_in_page == 0 {
let sblk_addr = sblk_addr.unwrap();
let npages = dblk_nelmts / self.page_nelmts;
if let Parent::Super { dblk_local, .. } = region.parent {
let global_page = dblk_local as u64 * npages + page;
let (byte, set) = sb_page_bit(file, sblk_addr, blk_off, global_page)?;
self.publish_super_block(file, sblk_addr, sb_geom, blk_off, byte, &[set])?;
}
}
}
Ok(())
}
fn super_block_addr<F: Store>(&self, file: &F, sblk_j: usize) -> Result<u64, Error> {
let os = file.offset_size() as usize;
let ib_prefix = (4 + 1 + 1 + os) as u64;
let ndblk_addrs = self.geom.direct_dblk_nelmts.len();
let slot_off = self.index_block_addr
+ ib_prefix
+ self.idx_blk_elmts * self.ea_elem_size as u64
+ ((ndblk_addrs + sblk_j) * os) as u64;
file.read_addr_at(slot_off)
}
fn ensure_super_block<F: Store>(
&self,
file: &mut F,
sblk_j: usize,
sb_block_offset: u64,
sb: SuperBlockGeom,
) -> Result<u64, Error> {
let existing = self.super_block_addr(file, sblk_j)?;
if !is_undef(existing, file.offset_size()) {
return Ok(existing);
}
let os = file.offset_size() as usize;
let ib_prefix = (4 + 1 + 1 + os) as u64;
let ndblk_addrs = self.geom.direct_dblk_nelmts.len();
let bitmap = vec![0u8; sb.bitmap_size().to_usize()?];
let undef = vec![undef_addr(file.offset_size()); sb.ndblks.to_usize()?];
let aesb = crate::chunked_write::build_aesb(
self.ea_addr,
sb_block_offset,
&bitmap,
&undef,
file.offset_size(),
self.blk_off_size,
self.client_id,
);
let new_addr = file.alloc_raw(&aesb)?;
#[cfg(test)]
alloc_probe::note_super_block();
file.sync()?;
let slot_off = self.index_block_addr
+ ib_prefix
+ self.idx_blk_elmts * self.ea_elem_size as u64
+ ((ndblk_addrs + sblk_j) * os) as u64;
self.publish_index_block(file, slot_off, &new_addr.to_le_bytes()[..os])?;
Ok(new_addr)
}
fn alloc_undef_data_block<F: Store>(
&self,
file: &mut F,
dblk_nelmts: u64,
block_offset_rel: u64,
) -> Result<u64, Error> {
let os = file.offset_size();
let mut buf = Vec::new();
buf.extend_from_slice(b"EADB");
buf.push(0); buf.push(self.client_id);
write_ea_addr(&mut buf, self.ea_addr, os);
buf.extend_from_slice(&block_offset_rel.to_le_bytes()[..self.blk_off_size]);
for _ in 0..dblk_nelmts {
push_undef_element(&mut buf, os, self.ea_elem_size);
}
let cks = jenkins_lookup3(&buf);
buf.extend_from_slice(&cks.to_le_bytes());
file.alloc_raw(&buf)
}
fn alloc_undef_paged_data_block<F: Store>(
&self,
file: &mut F,
dblk_nelmts: u64,
block_offset_rel: u64,
) -> Result<u64, Error> {
let os = file.offset_size();
let page_nelmts = self.page_nelmts;
let mut buf = Vec::new();
buf.extend_from_slice(b"EADB");
buf.push(0); buf.push(self.client_id);
write_ea_addr(&mut buf, self.ea_addr, os);
buf.extend_from_slice(&block_offset_rel.to_le_bytes()[..self.blk_off_size]);
let header_cks = jenkins_lookup3(&buf);
buf.extend_from_slice(&header_cks.to_le_bytes());
let npages = (dblk_nelmts / page_nelmts).to_usize()?;
for _ in 0..npages {
let mut page = Vec::with_capacity(page_nelmts.to_usize()? * self.ea_elem_size + 4);
for _ in 0..page_nelmts {
push_undef_element(&mut page, os, self.ea_elem_size);
}
let page_cks = jenkins_lookup3(&page);
page.extend_from_slice(&page_cks.to_le_bytes());
buf.extend_from_slice(&page);
}
file.alloc_raw(&buf)
}
fn element_size_width(&self, os: usize) -> usize {
if self.client_id == 0 {
return 0;
}
debug_assert!(
(os + 5..=os + 12).contains(&self.ea_elem_size),
"locate admitted an element width of {} for a {os}-byte address",
self.ea_elem_size
);
self.ea_elem_size - os - 4
}
fn check_stored_size(&self, os: usize, stored_size: u64) -> Result<(), Error> {
let csz = self.element_size_width(os);
if self.client_id != 0 && csz < 8 && stored_size >= (1u64 << (8 * csz)) {
return Err(Error::AppendUnsupported(
"recompressed chunk size exceeds the dataset's extensible-array element width",
));
}
Ok(())
}
fn element_bytes(
&self,
buf: &mut [u8; MAX_EA_ELEM],
os: usize,
rec: ElemRecord,
) -> Result<usize, Error> {
buf[..os].copy_from_slice(&rec.addr.to_le_bytes()[..os]);
if self.client_id == 0 {
return Ok(os);
}
self.check_stored_size(os, rec.stored_size)?;
let csz = self.element_size_width(os);
buf[os..os + csz].copy_from_slice(&rec.stored_size.to_le_bytes()[..csz]);
buf[os + csz..os + csz + 4].copy_from_slice(&rec.filter_mask.to_le_bytes());
Ok(os + csz + 4)
}
fn publish_index_block<F: Store>(
&self,
file: &mut F,
at: u64,
value: &[u8],
) -> Result<(), Error> {
let os = file.offset_size() as usize;
let ib_prefix = (4 + 1 + 1 + os) as u64;
let ndblk_addrs = self.geom.direct_dblk_nelmts.len();
let nsblk_addrs = self.geom.nsblk_addrs;
let cks_off = self.index_block_addr
+ ib_prefix
+ self.idx_blk_elmts * self.ea_elem_size as u64
+ ((ndblk_addrs + nsblk_addrs) * os) as u64;
file.publish_checksummed(self.index_block_addr, cks_off, at, value)
}
fn publish_super_block<F: Store>(
&self,
file: &mut F,
sblk_addr: u64,
sb: SuperBlockGeom,
blk_off: usize,
at: u64,
value: &[u8],
) -> Result<(), Error> {
let os = file.offset_size() as usize;
let prefix = (4 + 1 + 1 + os + blk_off) as u64;
let cks_off = sblk_addr + prefix + sb.bitmap_size() + sb.ndblks * os as u64;
file.publish_checksummed(sblk_addr, cks_off, at, value)
}
pub(crate) fn update_ea_header<F: Store>(
&self,
file: &mut F,
num_chunks: u64,
) -> Result<(), Error> {
let stats = ea_compute_stats(
&self.geom,
self.idx_blk_elmts,
self.ea_elem_size,
self.page_nelmts,
file.offset_size(),
self.blk_off_size,
num_chunks,
crate::chunked_write::SlotOccupancy::Dense(num_chunks),
);
let ls = file.length_size() as usize;
let ea_addr = self.ea_addr;
let mut stat_block = [0u8; 6 * 8];
let mut at = 0;
for value in [
stats.nsuper_blks,
stats.super_blk_size,
stats.ndata_blks,
stats.data_blk_size,
stats.max_idx_set,
stats.nelmts,
] {
stat_block[at..at + ls].copy_from_slice(&value.to_le_bytes()[..ls]);
at += ls;
}
let aehd_size =
ExtensibleArrayHeader::serialized_size(file.offset_size(), file.length_size()) as u64;
let cks_off = ea_addr + aehd_size - 4;
file.publish_checksummed(ea_addr, cks_off, ea_addr + 12, &stat_block[..at])
}
pub(crate) fn patch_dimension<F: Store>(
&self,
file: &mut F,
new_dim: u64,
) -> Result<(), Error> {
let ls = file.length_size() as usize;
file.publish_checksummed(
self.ohdr_chunk_start,
self.ohdr_chunk_msg_end,
self.dim0_off,
&new_dim.to_le_bytes()[..ls],
)
}
}
pub(crate) struct AppendPlan {
pub new_chunk_bytes: Vec<Vec<u8>>,
pub n_full: u64,
pub new_dim: u64,
pub new_num_chunks: u64,
}
pub(crate) fn plan_ea_append<F: Store>(
file: &F,
loc: &Located,
datatype: &Datatype,
spatial: &[u64],
element_size: NonZeroUsize,
pipeline: Option<&FilterPipeline>,
grow_visible_tail: bool,
raw: &[u8],
new_elems: u64,
fill: FillPattern<'_>,
) -> Result<AppendPlan, Error> {
let chunk_elems = loc.chunk_elems;
let current_dim = loc.current_dim;
let new_dim = current_dim
.checked_add(new_elems)
.ok_or(Error::AppendUnsupported(
"append would overflow the dataset dimension",
))?;
let n_full = current_dim / chunk_elems;
let has_partial = current_dim % chunk_elems != 0;
if has_partial && let Some(pl) = pipeline {
if !grow_visible_tail {
return Err(Error::AppendUnsupported(
"a filtered dataset whose length is not a whole multiple of the chunk length \
cannot be appended in place by this writer: growing its trailing partial chunk \
would repoint an index element a concurrent reader can already see. Use \
Dataset::append_staged, or append whole chunks",
));
}
if !pipeline_lossless(pl) {
return Err(Error::AppendUnsupported(LOSSY_TAIL_REFUSAL));
}
}
let mut tail_raw: Vec<u8> = Vec::new();
if has_partial {
let rec = loc
.read_element(file, n_full)?
.ok_or(Error::AppendUnsupported(
"trailing partial chunk is missing from the index",
))?;
let stored_len = if pipeline.is_some() {
usize::try_from(rec.stored_size)
.map_err(|_| Error::AppendUnsupported("chunk size exceeds this platform"))?
} else {
chunk_elems.to_usize()? * element_size.get()
};
rec.addr
.checked_add(stored_len as u64)
.filter(|&e| e <= file.len())
.ok_or(Error::AppendUnsupported(
"trailing chunk extends past end-of-file",
))?;
let stored = file.read_exact_at(rec.addr, stored_len)?;
let full = if let Some(pl) = pipeline {
let ctx = ChunkContext::from_datatype(spatial, datatype)?;
decompress_chunk(&stored, pl, ctx, rec.filter_mask).map_err(Error::Format)?
} else {
stored
};
let live_elems = usize::try_from(current_dim % chunk_elems)
.map_err(|_| Error::AppendUnsupported("chunk length exceeds this platform"))?;
let live_bytes = live_elems * element_size.get();
if full.len() < live_bytes {
return Err(Error::AppendUnsupported(
"trailing chunk decoded shorter than its live element count",
));
}
tail_raw.extend_from_slice(&full[..live_bytes]);
}
tail_raw.extend_from_slice(raw);
let tail_len_elems = new_dim - n_full * chunk_elems;
let split = split_into_chunks(&tail_raw, &[tail_len_elems], spatial, element_size, fill)
.map_err(Error::Format)?;
let new_chunk_bytes: Vec<Vec<u8>> = if let Some(pl) = pipeline {
let ctx = ChunkContext::from_datatype(spatial, datatype)?;
let mut out = Vec::with_capacity(split.len());
let mut scratch = FilterScratch::new();
for buf in &split {
out.push(compress_chunk_with(&mut scratch, buf, pl, ctx).map_err(Error::Format)?);
}
out
} else {
split
};
let os = file.offset_size() as usize;
for blob in &new_chunk_bytes {
loc.check_stored_size(os, blob.len() as u64)?;
}
let new_num_chunks = n_full + new_chunk_bytes.len() as u64;
Ok(AppendPlan {
new_chunk_bytes,
n_full,
new_dim,
new_num_chunks,
})
}
pub(crate) fn apply_ea_append<F: Store>(
file: &mut F,
loc: &mut Located,
plan: &AppendPlan,
max_phase: u8,
) -> Result<(), Error> {
let mut recorded_eof = file.len();
let mut chunk_addrs = Vec::with_capacity(plan.new_chunk_bytes.len());
for blob in &plan.new_chunk_bytes {
chunk_addrs.push((file.alloc_raw(blob)?, blob.len() as u64));
}
file.sync()?;
if file.len() != recorded_eof {
file.patch_superblock_eof()?;
file.sync()?;
recorded_eof = file.len();
}
if max_phase < 2 {
return Ok(());
}
for (k, &(addr, stored_size)) in chunk_addrs.iter().enumerate() {
let e = plan.n_full + k as u64;
let rec = ElemRecord {
addr,
stored_size,
filter_mask: 0,
};
loc.ea_insert(file, e, rec)?;
}
file.sync()?;
if max_phase < 3 {
return Ok(());
}
if file.len() != recorded_eof {
file.patch_superblock_eof()?;
}
loc.update_ea_header(file, plan.new_num_chunks)?;
file.sync()?;
if max_phase < 4 {
return Ok(());
}
loc.patch_dimension(file, plan.new_dim)?;
file.sync()?;
loc.current_dim = plan.new_dim;
loc.num_chunks = plan.new_num_chunks;
Ok(())
}
fn sb_page_bit<F: Store>(
file: &F,
sblk_addr: u64,
blk_off: usize,
global_page: u64,
) -> Result<(u64, u8), Error> {
let os = file.offset_size() as usize;
let bitmap_start = sblk_addr + (4 + 1 + 1 + os + blk_off) as u64;
let byte = bitmap_start + global_page / 8;
let mut v = [0u8; 1];
file.read_at(byte, &mut v)?;
Ok((byte, v[0] | (0x80u8 >> (global_page % 8))))
}
fn sb_dblk_slot_off(
os: usize,
sblk_addr: u64,
dblk_local: usize,
sb: SuperBlockGeom,
blk_off: usize,
) -> u64 {
let prefix = (4 + 1 + 1 + os + blk_off) as u64;
sblk_addr + prefix + sb.bitmap_size() + (dblk_local * os) as u64
}
enum Parent {
IndexDirect { ordinal: usize },
Super { sblk_j: usize, dblk_local: usize },
}
struct DataBlockLoc {
db_start: u64,
dblk_nelmts: u64,
ndblks: u64,
sb_block_offset: u64,
parent: Parent,
}
impl DataBlockLoc {
fn super_block(&self, page_nelmts: u64) -> SuperBlockGeom {
SuperBlockGeom {
ndblks: self.ndblks,
blocks: self.blocks(page_nelmts),
}
}
fn blocks(&self, page_nelmts: u64) -> DataBlockGeom {
DataBlockGeom {
dblk_nelmts: self.dblk_nelmts,
page_nelmts,
}
}
}
fn locate_data_block(geom: &EaGeometry, idx_blk_elmts: u64, e: u64) -> DataBlockLoc {
let mut elem = idx_blk_elmts;
for (ordinal, &dn) in geom.direct_dblk_nelmts.iter().enumerate() {
if e < elem + dn {
return DataBlockLoc {
db_start: elem,
dblk_nelmts: dn,
ndblks: 1,
sb_block_offset: 0,
parent: Parent::IndexDirect { ordinal },
};
}
elem += dn;
}
for j in 0..geom.nsblk_addrs {
let (ndblks, dn) = geom.sblks[geom.first_indirect_sblk + j];
let span = ndblks * dn;
if e < elem + span {
let sb_block_offset = elem - idx_blk_elmts;
let within = e - elem;
#[expect(
clippy::cast_possible_truncation,
reason = "within/dn is a data-block index bounded by ndblks (small)"
)]
let dblk_local = (within / dn) as usize;
let db_start = elem + dblk_local as u64 * dn;
return DataBlockLoc {
db_start,
dblk_nelmts: dn,
ndblks,
sb_block_offset,
parent: Parent::Super {
sblk_j: j,
dblk_local,
},
};
}
elem += span;
}
DataBlockLoc {
db_start: e,
dblk_nelmts: u64::MAX,
ndblks: 0,
sb_block_offset: 0,
parent: Parent::IndexDirect { ordinal: 0 },
}
}
struct WalkedMessage {
msg_type: MessageType,
flags: u8,
data_off: u64,
size: usize,
chunk_start: u64,
chunk_msg_end: u64,
}
struct Walk {
messages: Vec<WalkedMessage>,
}
fn walk_v2_object_header<S: Source + ?Sized>(
source: &S,
offset: u64,
offset_size: u8,
length_size: u8,
) -> Result<Walk, Error> {
let head = match source.read_metadata_at(offset, 6) {
Ok(head) => head,
Err(FormatError::UnexpectedEof { .. }) => {
return Err(Error::Format(FormatError::InvalidObjectHeaderSignature));
}
Err(e) => return Err(Error::Format(e)),
};
if &head[..4] != b"OHDR" {
return Err(Error::Format(FormatError::InvalidObjectHeaderSignature));
}
let flags = head[5];
let mut pos = offset + 6;
if flags & 0x20 != 0 {
pos += 16; }
if flags & 0x10 != 0 {
pos += 4; }
let chunk_size_width = 1usize << (flags & 0x03);
let size_buf = source.read_metadata_at(pos, chunk_size_width)?;
let chunk0_size = read_uint(&size_buf, 0, chunk_size_width)?.to_usize()?;
pos += chunk_size_width as u64;
let chunk0_start = offset;
let chunk0_msg_start = pos;
let chunk0_msg_end = chunk0_msg_start + chunk0_size as u64;
let has_creation_order = flags & 0x04 != 0;
let mut messages = Vec::new();
let mut continuations: Vec<(u64, usize)> = Vec::new();
let chunk0 = source.read_metadata_at(chunk0_msg_start, chunk0_size)?;
walk_messages(
&chunk0,
chunk0_msg_start,
chunk0_start,
chunk0_msg_end,
has_creation_order,
offset_size,
length_size,
&mut messages,
&mut continuations,
)?;
let mut guard = 256;
while let Some((cont_off, cont_len)) = continuations.pop() {
guard -= 1;
if guard == 0 {
return Err(Error::Format(FormatError::NestingDepthExceeded));
}
if cont_len < 8 {
return Err(Error::Format(FormatError::InvalidObjectHeaderSignature));
}
let chunk = source.read_metadata_at(cont_off, cont_len)?;
if &chunk[..4] != b"OCHK" {
return Err(Error::Format(FormatError::InvalidObjectHeaderSignature));
}
let msg_start = cont_off + 4;
let msg_end = cont_off + (cont_len - 4) as u64; walk_messages(
&chunk[4..cont_len - 4],
msg_start,
cont_off,
msg_end,
has_creation_order,
offset_size,
length_size,
&mut messages,
&mut continuations,
)?;
}
Ok(Walk { messages })
}
#[allow(clippy::too_many_arguments)]
fn walk_messages(
chunk: &[u8],
base: u64,
chunk_start: u64,
chunk_msg_end: u64,
has_creation_order: bool,
offset_size: u8,
length_size: u8,
messages: &mut Vec<WalkedMessage>,
continuations: &mut Vec<(u64, usize)>,
) -> Result<(), Error> {
let msg_header_size = if has_creation_order { 6 } else { 4 };
let end = chunk.len();
let mut pos = 0usize;
while pos + msg_header_size <= end {
let msg_type_raw = chunk[pos] as u16;
let msg_data_size = u16::from_le_bytes([chunk[pos + 1], chunk[pos + 2]]) as usize;
let msg_flags = chunk[pos + 3];
pos += msg_header_size;
if pos + msg_data_size > end {
break; }
let msg_type = MessageType::from_u16(msg_type_raw);
if msg_type == MessageType::ObjectHeaderContinuation {
let cont_off = read_uint(chunk, pos, offset_size as usize)?;
let cont_len =
read_uint(chunk, pos + offset_size as usize, length_size as usize)?.to_usize()?;
continuations.push((cont_off, cont_len));
} else {
messages.push(WalkedMessage {
msg_type,
flags: msg_flags,
data_off: base + pos as u64,
size: msg_data_size,
chunk_start,
chunk_msg_end,
});
}
pos += msg_data_size;
}
Ok(())
}
fn read_uint(data: &[u8], pos: usize, size: usize) -> Result<u64, Error> {
if pos + size > data.len() {
return Err(Error::Format(FormatError::UnexpectedEof {
expected: pos + size,
available: data.len(),
}));
}
let mut v = 0u64;
for i in 0..size {
v |= (data[pos + i] as u64) << (8 * i);
}
Ok(v)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::group_v2;
use crate::signature;
use crate::source::BytesSource;
use crate::superblock::Superblock;
use crate::writer::FileBuilder;
use std::cell::{Cell, RefCell};
struct WindowProbeStore {
data: Vec<u8>,
superblock: Superblock,
sb_sig_off: usize,
max_read: Cell<usize>,
total_read: Cell<usize>,
reads: RefCell<Vec<(u64, usize)>>,
superblock_patches: Cell<usize>,
syncs: Cell<usize>,
appends: Cell<usize>,
free: Option<(u64, u64)>,
reuses: Cell<usize>,
}
impl WindowProbeStore {
fn open(data: Vec<u8>) -> Self {
let sb_sig_off = signature::find_signature(&data).unwrap();
let superblock = Superblock::parse(&data, sb_sig_off).unwrap();
Self {
data,
superblock,
sb_sig_off,
max_read: Cell::new(0),
total_read: Cell::new(0),
reads: RefCell::new(Vec::new()),
superblock_patches: Cell::new(0),
syncs: Cell::new(0),
appends: Cell::new(0),
free: None,
reuses: Cell::new(0),
}
}
fn with_free_tail(mut self, len: u64) -> Self {
let addr = self.data.len() as u64;
self.data
.resize(self.data.len() + len.to_usize().unwrap(), 0);
self.patch_superblock_eof().unwrap();
self.free = Some((addr, len));
self
}
fn append_bytes(&mut self, bytes: &[u8]) -> Result<u64, Error> {
self.appends.set(self.appends.get() + 1);
let addr = self.data.len() as u64;
self.data.extend_from_slice(bytes);
Ok(addr)
}
fn reset_counters(&self) {
self.max_read.set(0);
self.total_read.set(0);
self.reads.borrow_mut().clear();
self.superblock_patches.set(0);
self.syncs.set(0);
self.appends.set(0);
self.reuses.set(0);
}
}
impl Source for WindowProbeStore {
fn len(&self) -> u64 {
self.data.len() as u64
}
fn read_at(&self, offset: u64, buf: &mut [u8]) -> Result<(), FormatError> {
self.max_read.set(self.max_read.get().max(buf.len()));
self.total_read.set(self.total_read.get() + buf.len());
self.reads.borrow_mut().push((offset, buf.len()));
BytesSource::new(&self.data).read_at(offset, buf)
}
}
impl Store for WindowProbeStore {
fn offset_size(&self) -> u8 {
self.superblock.offset_size
}
fn length_size(&self) -> u8 {
self.superblock.length_size
}
fn alloc_raw(&mut self, bytes: &[u8]) -> Result<u64, Error> {
let len = bytes.len() as u64;
if let Some((addr, avail)) = self.free.filter(|&(_, avail)| avail >= len) {
self.reuses.set(self.reuses.get() + 1);
self.free = (avail > len).then(|| (addr + len, avail - len));
self.write_at(addr, bytes)?;
return Ok(addr);
}
self.append_bytes(bytes)
}
fn write_at(&mut self, offset: u64, bytes: &[u8]) -> Result<(), Error> {
let offset = offset.to_usize()?;
self.data[offset..offset + bytes.len()].copy_from_slice(bytes);
Ok(())
}
fn patch_superblock_eof(&mut self) -> Result<(), Error> {
self.superblock_patches
.set(self.superblock_patches.get() + 1);
self.superblock.eof_address = self.data.len() as u64;
let bytes = self.superblock.serialize();
let off = self.sb_sig_off;
self.data[off..off + bytes.len()].copy_from_slice(&bytes);
Ok(())
}
fn sync(&mut self) -> Result<(), Error> {
self.syncs.set(self.syncs.get() + 1);
Ok(())
}
}
fn build_unlimited(n: i32, chunk: u64) -> Vec<u8> {
let data: Vec<i32> = (0..n).collect();
let mut b = FileBuilder::new();
b.create_dataset("d")
.with_i32_data(&data)
.with_shape(&[n as u64])
.with_maxshape(&[u64::MAX])
.with_chunks(&[chunk]);
b.finish().unwrap()
}
fn locate(store: &WindowProbeStore) -> (Located, crate::datatype::Datatype) {
let oh_addr = group_v2::resolve_path_any(&store.data, &store.superblock, "d").unwrap();
let result = Located::locate_at(store, oh_addr, Error::AppendUnsupported).unwrap();
let (dt_off, dt_size) = result.spans.datatype;
let dt_bytes = store.read_metadata_at(dt_off, dt_size).unwrap();
let (datatype, _) = crate::datatype::Datatype::parse(&dt_bytes).unwrap();
(result.located, datatype)
}
fn read_i32s(store: &WindowProbeStore, loc: &Located) -> Vec<i32> {
let width = loc.chunk_elems.to_usize().unwrap() * 4;
let mut out = Vec::new();
for e in 0..loc.num_chunks {
let rec = loc.read_element(store, e).unwrap().expect("indexed chunk");
let bytes = store.read_exact_at(rec.addr, width).unwrap();
out.extend(
bytes
.as_chunks::<4>()
.0
.iter()
.copied()
.map(i32::from_le_bytes),
);
}
out
}
fn append_i32s(
store: &mut WindowProbeStore,
loc: &mut Located,
datatype: &crate::datatype::Datatype,
values: std::ops::Range<i32>,
) {
let raw: Vec<u8> = values.clone().flat_map(|v| v.to_le_bytes()).collect();
let new_elems = (values.end - values.start) as u64;
let spatial = vec![loc.chunk_elems];
let plan = plan_ea_append(
store,
loc,
datatype,
&spatial,
loc.elem_bytes,
None,
true,
&raw,
new_elems,
FillPattern::ZERO,
)
.unwrap();
apply_ea_append(store, loc, &plan, 4).unwrap();
}
#[test]
fn allocating_an_index_block_barriers_before_the_pointer_that_names_it() {
let mut store = WindowProbeStore::open(build_unlimited(4, 1));
let (mut loc, datatype) = locate(&store);
let mut rounds: Vec<(usize, usize)> = Vec::new();
for i in 0..320i32 {
store.reset_counters();
append_i32s(&mut store, &mut loc, &datatype, (4 + i)..(5 + i));
rounds.push((store.syncs.get(), store.appends.get()));
}
let base = rounds
.iter()
.find(|&&(_, appends)| appends == 1)
.map(|&(syncs, _)| syncs)
.expect("some append allocates no index block");
let blocks = |appends: usize| appends - 1;
assert!(
rounds.iter().any(|&(_, appends)| blocks(appends) == 1),
"the run must allocate a data block somewhere: {rounds:?}"
);
assert!(
rounds.iter().any(|&(_, appends)| blocks(appends) >= 2),
"the run must reach a super block, which allocates two blocks in one \
append, or the second barrier site is never exercised: {rounds:?}"
);
for (round, &(syncs, appends)) in rounds.iter().enumerate() {
assert_eq!(
syncs,
base + blocks(appends),
"append {round} allocated {} index block(s) and took {syncs} \
barriers against the {base} a plain append takes: each fresh \
block owes one, between its bytes and the pointer naming it",
blocks(appends)
);
}
}
#[test]
fn phase_three_patches_the_superblock_only_when_the_index_grew() {
let mut store = WindowProbeStore::open(build_unlimited(4, 1));
let (mut loc, datatype) = locate(&store);
let mut patches = Vec::new();
for i in 0..40i32 {
store.reset_counters();
append_i32s(&mut store, &mut loc, &datatype, (4 + i)..(5 + i));
patches.push(store.superblock_patches.get());
}
assert!(
patches.iter().all(|&p| p >= 1),
"every append extends the file, so every one must record the new \
end-of-file at least once: {patches:?}"
);
assert!(
patches.contains(&1),
"an append that allocates no index block must not rewrite the \
superblock a second time: {patches:?}"
);
assert!(
patches.contains(&2),
"an append that does allocate one must, or the file would advertise \
an end-of-file short of its own index: {patches:?}"
);
assert!(
patches.iter().all(|&p| p <= 2),
"the sequence defines exactly two such points: {patches:?}"
);
}
#[test]
fn phase_one_patches_the_superblock_only_when_the_file_grew() {
let mut store = WindowProbeStore::open(build_unlimited(4, 4)).with_free_tail(64);
let (mut loc, datatype) = locate(&store);
let mut rounds = Vec::new();
for i in 0..6i32 {
store.reset_counters();
let base = 4 + i * 4;
append_i32s(&mut store, &mut loc, &datatype, base..(base + 4));
rounds.push((
store.reuses.get(),
store.appends.get(),
store.superblock_patches.get(),
));
}
for &(reuses, appends, patches) in &rounds {
assert_eq!(
patches == 0,
appends == 0,
"the superblock is rewritten exactly when the file grew: {rounds:?}"
);
assert!(reuses + appends >= 1, "every round allocates: {rounds:?}");
}
assert!(
rounds.iter().any(|&(_, appends, _)| appends == 0),
"the first rounds fit the hole, so they must allocate nothing at \
end-of-file: {rounds:?}"
);
assert!(
rounds.iter().any(|&(_, appends, _)| appends > 0),
"the hole is spent well before the last round: {rounds:?}"
);
assert_eq!(read_i32s(&store, &loc), (0..28).collect::<Vec<i32>>());
}
#[test]
fn a_hole_serves_the_index_blocks_as_well_as_the_chunks() {
let mut store = WindowProbeStore::open(build_unlimited(4, 1)).with_free_tail(16 * 1024);
let (mut loc, datatype) = locate(&store);
let mut rounds = Vec::new();
for i in 0..40i32 {
store.reset_counters();
append_i32s(&mut store, &mut loc, &datatype, (4 + i)..(5 + i));
rounds.push((store.reuses.get(), store.appends.get()));
}
assert!(
rounds.iter().all(|&(_, appends)| appends == 0),
"the hole is larger than the whole run, so nothing may reach \
end-of-file: {rounds:?}"
);
assert!(
rounds.iter().any(|&(reuses, _)| reuses >= 2),
"some append must allocate an index block beside its chunk, or this \
says nothing about the blocks: {rounds:?}"
);
assert_eq!(
read_i32s(&store, &loc),
(0..44).collect::<Vec<i32>>(),
"the reused regions must hold the dataset's elements"
);
}
#[test]
fn a_relocated_partial_tail_may_land_in_freed_space() {
const HOLE: u64 = 1024;
let mut store = WindowProbeStore::open(build_unlimited(4, 4)).with_free_tail(HOLE);
let hole = (store.len() - HOLE)..store.len();
let (mut loc, datatype) = locate(&store);
let mut tails = Vec::new();
for i in 4..12i32 {
append_i32s(&mut store, &mut loc, &datatype, i..(i + 1));
let tail = loc.num_chunks - 1;
let addr = loc.read_element(&store, tail).unwrap().unwrap().addr;
tails.push((loc.num_chunks, addr));
assert_eq!(
read_i32s(&store, &loc)[..(i as usize + 1)],
(0..=i).collect::<Vec<i32>>()[..],
"after appending {i} the live elements must all read back"
);
}
assert!(
tails
.windows(2)
.any(|w| w[0].0 == w[1].0 && w[0].1 != w[1].1),
"a partial trailing chunk must have been rewritten to a new address \
at least once: {tails:?}"
);
assert!(
tails.iter().all(|&(_, addr)| hole.contains(&addr)),
"every tail chunk must have landed in the freed region, or this says \
nothing about reuse: {tails:?}"
);
}
#[test]
fn append_reads_stay_bounded_windows() {
let n = 100_000i32;
let mut store = WindowProbeStore::open(build_unlimited(n, 256));
let file_len = store.data.len();
assert!(
file_len > 300_000,
"test file unexpectedly small: {file_len}"
);
let (mut loc, datatype) = locate(&store);
const WINDOW: usize = 16 * 1024;
assert!(
store.max_read.get() <= WINDOW,
"locate read a {}-byte window (> {WINDOW})",
store.max_read.get()
);
store.reset_counters();
append_i32s(&mut store, &mut loc, &datatype, n..n + 10);
assert!(
store.max_read.get() <= WINDOW,
"append read a {}-byte window (> {WINDOW})",
store.max_read.get()
);
assert!(
store.total_read.get() <= 64 * 1024,
"append read {} bytes total (> 64 KiB) on a {file_len}-byte file",
store.total_read.get()
);
assert!(
store
.reads
.borrow()
.iter()
.all(|&(_, len)| len < file_len / 2),
"an append read scaled with file size"
);
let mut worst = 0usize;
let mut next = n + 10;
for _ in 0..20 {
store.reset_counters();
append_i32s(&mut store, &mut loc, &datatype, next..next + 300);
worst = worst.max(store.total_read.get());
next += 300;
}
assert!(
worst <= 64 * 1024,
"per-append read volume grew to {worst} bytes"
);
let file = crate::File::from_bytes(store.data).unwrap();
let ds = file.dataset("d").unwrap();
let got = ds.read_i32().unwrap();
let expected: Vec<i32> = (0..next).collect();
assert_eq!(got.len(), expected.len());
assert_eq!(got, expected);
}
#[test]
fn partial_tail_append_reads_one_chunk_window() {
let n = 50_000i32;
let chunk = 256u64;
let mut store = WindowProbeStore::open(build_unlimited(n + 7, chunk));
let (mut loc, datatype) = locate(&store);
store.reset_counters();
append_i32s(&mut store, &mut loc, &datatype, n + 7..n + 7 + 13);
let chunk_bytes = (chunk as usize) * 4;
assert!(
store.max_read.get() <= chunk_bytes.max(8 * 1024),
"partial-tail append read a {}-byte window",
store.max_read.get()
);
let file = crate::File::from_bytes(store.data).unwrap();
let got = file.dataset("d").unwrap().read_i32().unwrap();
let expected: Vec<i32> = (0..n + 7 + 13).collect();
assert_eq!(got, expected);
}
#[test]
fn an_element_width_wider_than_a_u64_is_refused_not_a_panic() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("wide_elem.h5");
let mut b = crate::writer::FileBuilder::new();
b.create_dataset("d")
.with_i32_data(&(0..4096).collect::<Vec<_>>())
.with_shape(&[4096])
.with_maxshape(&[u64::MAX])
.with_chunks(&[4])
.with_deflate(1);
b.write(&path).unwrap();
let mut bytes = std::fs::read(&path).unwrap();
let at = bytes
.windows(4)
.position(|w| w == b"EAHD")
.expect("a filtered unlimited dataset is indexed by an extensible array");
assert_eq!(bytes[at + 6], 14, "the writer's own element width moved");
bytes[at + 6] = 30;
crate::checksum::stamp_trailing(&mut bytes, at, 12 + 6 * 8 + 8 + 4);
std::fs::write(&path, &bytes).unwrap();
let f = crate::reader::File::open_rw(&path).unwrap();
let err = f
.dataset("d")
.and_then(|mut d| d.append(&[1i32, 2, 3, 4]))
.expect_err("an element width wider than a u64 must be refused");
assert!(
std::format!("{err:?}").contains("element width"),
"the refusal should name the element width, but said {err:?}"
);
}
}