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::error::{Error, FormatError};
use crate::extensible_array::{EaGeometry, ExtensibleArrayHeader};
use crate::filter_pipeline::FilterPipeline;
use crate::filters::{ChunkContext, compress_chunk, decompress_chunk};
use crate::message_type::MessageType;
use crate::source::Source;
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 append_bytes(&mut self, bytes: &[u8]) -> Result<u64, Error>;
fn append_raw(&mut self, bytes: &[u8]) -> Result<u64, Error> {
self.append_bytes(bytes)
}
fn append_meta(&mut self, bytes: &[u8]) -> Result<u64, Error> {
self.append_bytes(bytes)
}
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 write_addr_at(&mut self, offset: u64, addr: u64) -> Result<(), Error> {
#[expect(
clippy::cast_possible_truncation,
reason = "the 4-byte arm is taken only when this file's offset_size is 4 bytes"
)]
match self.offset_size() {
4 => self.write_at(offset, &(addr as u32).to_le_bytes()),
_ => self.write_at(offset, &addr.to_le_bytes()),
}
}
fn write_length_at(&mut self, offset: u64, val: u64) -> Result<(), Error> {
#[expect(
clippy::cast_possible_truncation,
reason = "the 4-byte arm is taken only when this file's length_size is 4 bytes"
)]
match self.length_size() {
4 => self.write_at(offset, &(val as u32).to_le_bytes()),
_ => self.write_at(offset, &val.to_le_bytes()),
}
}
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 rechecksum_range(&mut self, start: u64, cks_off: u64) -> Result<(), Error> {
let bytes = self.read_exact_at(start, (cks_off - start).to_usize()?)?;
let cks = jenkins_lookup3(&bytes);
self.write_at(cks_off, &cks.to_le_bytes())
}
}
pub(crate) struct MessageSpans {
pub datatype: (u64, usize),
pub filter: Option<(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: usize,
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 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 elem_bytes = chunk_dims[1] as usize;
if elem_bytes == 0 {
return Err(unsupported("dataset has a zero-sized element"));
}
let chunk_bytes = chunk_elems.to_usize()? * elem_bytes;
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)",
));
}
if ea_header.client_id == 1 && (ea_header.element_size as usize) < os as usize + 5 {
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)),
},
})
}
fn write_element_at<F: Store>(
&self,
file: &mut F,
off: u64,
rec: ElemRecord,
) -> Result<(), Error> {
if self.client_id == 0 {
return file.write_addr_at(off, rec.addr);
}
let os = file.offset_size() as usize;
let csz = self.ea_elem_size - os - 4;
if csz < 8 && rec.stored_size >= (1u64 << (8 * csz)) {
return Err(Error::AppendUnsupported(
"recompressed chunk size exceeds the dataset's extensible-array element width",
));
}
file.write_addr_at(off, rec.addr)?;
file.write_at(off + os as u64, &rec.stored_size.to_le_bytes()[..csz])?;
file.write_at(off + (os + csz) as u64, &rec.filter_mask.to_le_bytes())
}
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.dblk_nelmts > page_nelmts;
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");
sb_dblk_slot_off(
os,
sblk_addr,
dblk_local,
region.ndblks,
region.dblk_nelmts,
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;
self.write_element_at(file, slot_off, rec)?;
self.rechecksum_index_block(file)?;
return Ok(());
}
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 is_paged = dblk_nelmts > self.page_nelmts;
let slot = e - region.db_start;
let block_offset_rel = region.db_start - idx;
let ndblks = region.ndblks;
let sblk_addr = match region.parent {
Parent::Super { sblk_j, .. } => Some(self.ensure_super_block(
file,
sblk_j,
region.sb_block_offset,
ndblks,
dblk_nelmts,
)?),
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)?
};
file.write_addr_at(dblk_ptr_off, new_addr)?;
match region.parent {
Parent::IndexDirect { .. } => self.rechecksum_index_block(file)?,
Parent::Super { .. } => self.rechecksum_super_block(
file,
sblk_addr.unwrap(),
ndblks,
dblk_nelmts,
self.page_nelmts,
blk_off,
)?,
}
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;
self.write_element_at(file, elem_off, rec)?;
let cks_off = dblk_addr + db_prefix + dblk_nelmts * elem_size;
file.rechecksum_range(dblk_addr, cks_off)?;
} 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;
self.write_element_at(file, page_off + slot_in_page * elem_size, rec)?;
let page_cks_off = page_off + page_nelmts * elem_size;
file.rechecksum_range(page_off, page_cks_off)?;
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;
self.set_sb_page_bit(file, sblk_addr, blk_off, global_page)?;
self.rechecksum_super_block(
file,
sblk_addr,
ndblks,
dblk_nelmts,
self.page_nelmts,
blk_off,
)?;
}
}
}
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,
ndblks: u64,
dblk_nelmts: u64,
) -> 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(ndblks, dblk_nelmts, self.page_nelmts)?];
let undef = vec![undef_addr(file.offset_size()); 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.append_meta(&aesb)?;
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.write_addr_at(slot_off, new_addr)?;
self.rechecksum_index_block(file)?;
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.append_meta(&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.append_meta(&buf)
}
fn set_sb_page_bit<F: Store>(
&self,
file: &mut F,
sblk_addr: u64,
blk_off: usize,
global_page: u64,
) -> Result<(), 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 mask = 0x80u8 >> (global_page % 8);
let mut v = [0u8; 1];
file.read_at(byte, &mut v)?;
file.write_at(byte, &[v[0] | mask])
}
fn rechecksum_index_block<F: Store>(&self, file: &mut F) -> 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.rechecksum_range(self.index_block_addr, cks_off)
}
fn rechecksum_super_block<F: Store>(
&self,
file: &mut F,
sblk_addr: u64,
ndblks: u64,
dblk_nelmts: u64,
page_nelmts: u64,
blk_off: usize,
) -> Result<(), Error> {
let os = file.offset_size() as usize;
let prefix = (4 + 1 + 1 + os + blk_off) as u64;
let bitmap = sb_bitmap_size(ndblks, dblk_nelmts, page_nelmts)? as u64;
let cks_off = sblk_addr + prefix + bitmap + ndblks * os as u64;
file.rechecksum_range(sblk_addr, cks_off)
}
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,
);
let ls = file.length_size() as u64;
let ea_addr = self.ea_addr;
file.write_length_at(ea_addr + 12, stats.nsuper_blks)?;
file.write_length_at(ea_addr + 12 + ls, stats.super_blk_size)?;
file.write_length_at(ea_addr + 12 + 2 * ls, stats.ndata_blks)?;
file.write_length_at(ea_addr + 12 + 3 * ls, stats.data_blk_size)?;
file.write_length_at(ea_addr + 12 + 4 * ls, stats.max_idx_set)?;
file.write_length_at(ea_addr + 12 + 5 * ls, stats.nelmts)?;
let aehd_size =
ExtensibleArrayHeader::serialized_size(file.offset_size(), file.length_size()) as u64;
let cks_off = ea_addr + aehd_size - 4;
file.rechecksum_range(ea_addr, cks_off)
}
pub(crate) fn patch_dimension<F: Store>(
&self,
file: &mut F,
new_dim: u64,
) -> Result<(), Error> {
file.write_length_at(self.dim0_off, new_dim)?;
file.rechecksum_range(self.ohdr_chunk_start, self.ohdr_chunk_msg_end)
}
}
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: usize,
pipeline: Option<&FilterPipeline>,
raw: &[u8],
new_elems: u64,
) -> 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 pipeline.is_some() && (has_partial || new_elems % chunk_elems != 0) {
return Err(Error::AppendUnsupported(
"a filtered dataset can only be appended in place in whole chunks (the current \
length and the appended length must both be multiples of the chunk length); \
use Dataset::append_staged for a non-chunk-aligned filtered append",
));
}
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
};
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;
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);
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());
for (_, buf) in &split {
out.push(compress_chunk(buf, pl, ctx).map_err(Error::Format)?);
}
out
} else {
split.into_iter().map(|(_, buf)| buf).collect()
};
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 chunk_addrs = Vec::with_capacity(plan.new_chunk_bytes.len());
for blob in &plan.new_chunk_bytes {
chunk_addrs.push((file.append_raw(blob)?, blob.len() as u64));
}
file.sync()?;
file.patch_superblock_eof()?;
file.sync()?;
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(());
}
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_bitmap_size(ndblks: u64, dblk_nelmts: u64, page_nelmts: u64) -> Result<usize, Error> {
if dblk_nelmts > page_nelmts {
let npages = (dblk_nelmts / page_nelmts).to_usize()?;
Ok(ndblks.to_usize()? * npages.div_ceil(8))
} else {
Ok(0)
}
}
#[allow(clippy::too_many_arguments)]
fn sb_dblk_slot_off(
os: usize,
sblk_addr: u64,
dblk_local: usize,
ndblks: u64,
dblk_nelmts: u64,
page_nelmts: u64,
blk_off: usize,
) -> Result<u64, Error> {
let prefix = (4 + 1 + 1 + os + blk_off) as u64;
let bitmap = sb_bitmap_size(ndblks, dblk_nelmts, page_nelmts)? as u64;
Ok(sblk_addr + prefix + bitmap + (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,
}
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,
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;
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,
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)>>,
}
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()),
}
}
fn reset_counters(&self) {
self.max_read.set(0);
self.total_read.set(0);
self.reads.borrow_mut().clear();
}
}
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 append_bytes(&mut self, bytes: &[u8]) -> Result<u64, Error> {
let addr = self.data.len() as u64;
self.data.extend_from_slice(bytes);
Ok(addr)
}
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.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> {
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 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,
&raw,
new_elems,
)
.unwrap();
apply_ea_append(store, loc, &plan, 4).unwrap();
}
#[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);
}
}