use std::collections::{BTreeMap, HashMap, HashSet};
use std::fs;
use std::io::{Read, Seek, SeekFrom};
use std::path::Path;
use crate::checksum::jenkins_lookup3;
use crate::chunk_index_inplace::{Located, Store, apply_ea_append, plan_ea_append};
use crate::chunked_read::{
chunk_index_spans_from_source, enumerate_chunks_from_source, plan_dense_grid,
};
use crate::chunked_write::{
ChunkMeta, ChunkOptions, ChunkProvider, WrittenChunk, build_chunked_data_at_ext,
build_extensible_array_at, emit_chunked_data_verbatim, plan_chunked_data_verbatim,
serialize_v4_extensible_array, split_into_chunks,
};
use crate::convert::TryToUsize;
use crate::data_layout::DataLayout;
use crate::dataspace::{Dataspace, DataspaceType};
use crate::datatype::{Datatype, DatatypeByteOrder};
use crate::error::{Error, FormatError, OBJECT_HEADER_MESSAGE_MAX};
use crate::extensible_array::ExtensibleArrayHeader;
use crate::file_create_properties::FileCreateProperties;
use crate::file_lock::{self, FileLocking};
use crate::file_space_info::{FileSpaceInfo, FileSpaceStrategy, NUM_FILE_FSM_MANAGERS};
use crate::file_writer::{
LENGTH_SIZE, OFFSET_SIZE, build_chunked_dataset_oh, build_dataset_oh, make_link,
};
use crate::filter_pipeline::{
FILTER_DEFLATE, FILTER_FLETCHER32, FILTER_LZF, FILTER_SCALEOFFSET, FILTER_SHUFFLE,
FilterPipeline,
};
use crate::filters::{ChunkContext, compress_chunk, decompress_chunk};
use crate::free_space::FreeList;
use crate::free_space_manager::{
self, FreeSection, FsmHeader, PageType, SECT_CLASS_SIMPLE, align_up, free_sections, fshd_len,
plan_paged_managers, serialize_file_fsm,
};
use crate::group_v2::resolve_group_entries_from_source;
use crate::image::{FileImage, HandleImage, MirrorImage};
use crate::link_message::{LinkMessage, LinkTarget};
use crate::message_type::MessageType;
use crate::object_header::ObjectHeader;
use crate::reader::FileAccessProperties;
use crate::signature;
use crate::source::{BaseOffsetSource, BytesSource, MetadataCacheConfig, Source};
use crate::superblock::Superblock;
use crate::type_builders::{
AttrValue, DatasetBuilder, ObjectRefPatch, ObjectRefTarget, VlStringStaging,
build_attr_message, build_global_heap_collections, make_f32_type, make_f64_type, make_i8_type,
make_i16_type, make_i32_type, make_i64_type, make_u8_type, make_u16_type, make_u32_type,
make_u64_type, patch_vl_refs, patch_vl_refs_masked, write_reference_address,
};
const UNDEF: u64 = u64::MAX;
const MAX_COMPACT_ATTRS: usize = 8;
const MAX_COPY_DEPTH: u32 = 1000;
const MAX_LINK_GRAPH_NODES: u32 = 1 << 24;
const MAX_OH_CHUNKS: usize = 256;
const OH_PREFIX_MAX: usize = 34;
type PathKey = Vec<String>;
type PendingVlAttrs = Vec<(crate::attribute::AttributeMessage, Vec<Vec<u8>>)>;
pub struct AppendBuilder {
raw: Vec<u8>,
elem_dt: Option<Datatype>,
dt_conflict: bool,
}
impl AppendBuilder {
pub(crate) fn new() -> Self {
Self {
raw: Vec::new(),
elem_dt: None,
dt_conflict: false,
}
}
pub(crate) fn raw(&self) -> &[u8] {
&self.raw
}
pub(crate) fn elem_dt(&self) -> Option<&Datatype> {
self.elem_dt.as_ref()
}
pub(crate) fn dt_conflict(&self) -> bool {
self.dt_conflict
}
fn set_dt(&mut self, dt: Datatype) {
match &self.elem_dt {
Some(prev) if *prev != dt => self.dt_conflict = true,
Some(_) => {}
None => self.elem_dt = Some(dt),
}
}
pub fn append_raw(&mut self, bytes: &[u8]) -> &mut Self {
self.raw.extend_from_slice(bytes);
self
}
pub fn append<T: crate::element::H5Element>(&mut self, data: &[T]) -> &mut Self {
T::append_into(self, data);
self
}
}
macro_rules! append_typed {
($($method:ident, $ty:ty, $make:ident;)*) => {
impl AppendBuilder {
$(
#[doc = concat!("Append `", stringify!($ty), "` values to the dataset.")]
pub fn $method(&mut self, data: &[$ty]) -> &mut Self {
self.set_dt($make());
for &v in data {
self.raw.extend_from_slice(&v.to_le_bytes());
}
self
}
)*
}
};
}
append_typed! {
append_f64, f64, make_f64_type;
append_f32, f32, make_f32_type;
append_i8, i8, make_i8_type;
append_i16, i16, make_i16_type;
append_i32, i32, make_i32_type;
append_i64, i64, make_i64_type;
append_u8, u8, make_u8_type;
append_u16, u16, make_u16_type;
append_u32, u32, make_u32_type;
append_u64, u64, make_u64_type;
}
pub(crate) struct WriteEngine {
image: Box<dyn FileImage>,
sb_sig_off: usize,
superblock: Superblock,
pending_datasets: Vec<(PathKey, DatasetBuilder)>,
pending_writes: Vec<(PathKey, DatasetBuilder)>,
pending_appends: Vec<(PathKey, AppendBuilder)>,
pending_groups: Vec<PathKey>,
pending_group_attrs: Vec<(PathKey, AttrOp)>,
pending_dataset_attrs: Vec<(PathKey, AttrOp)>,
pending_deletes: Vec<PathKey>,
pending_copies: Vec<(PathKey, PathKey)>,
pending_cross_copies: Vec<(PathKey, CopyTree)>,
free: FreeList,
persist: Option<PersistState>,
located: HashMap<u64, LocatedState>,
swmr_mode: bool,
paged: Option<PagedEdit>,
committed: bool,
resolved: HashMap<String, u64>,
batched_appends: bool,
bounded: bool,
fsm_len: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum MemoryStrategy {
Bounded,
Auto,
Mirrored,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum EditBacking {
Bounded,
Mirrored,
}
impl From<EditBacking> for MemoryStrategy {
fn from(backing: EditBacking) -> Self {
match backing {
EditBacking::Bounded => Self::Bounded,
EditBacking::Mirrored => Self::Mirrored,
}
}
}
fn bounded_only_limitation(session: &WriteEngine) -> Option<&'static str> {
if session.superblock.version < 2 {
return Some(
"bounded read-write access requires a latest-format file (v2/v3 superblock); \
leave MemoryStrategy unset, or pass MemoryStrategy::Auto, to fall back to \
the whole-file mirror here",
);
}
if session.superblock.base_address != 0 {
return Some(
"bounded read-write access does not support a file with a userblock \
(non-zero base address); leave MemoryStrategy unset, or pass \
MemoryStrategy::Auto, to fall back to the whole-file mirror here",
);
}
None
}
pub(crate) fn create_would_refuse_reopen(
create: &FileCreateProperties,
access: &FileAccessProperties,
) -> Option<&'static str> {
if let Some((FileSpaceStrategy::Page, false, _)) = create.file_space_strategy() {
return Some(
"a paged file (FileSpaceStrategy::Page) with persist = false cannot be reopened \
read-write, so creating one this way would write the file and then fail to open \
it; pass persist = true to with_file_space_strategy, or build the file with \
FileBuilder if it is only ever going to be read",
);
}
if create.userblock() != 0 && access.memory_strategy() == Some(MemoryStrategy::Bounded) {
return Some(
"a userblock cannot be combined with MemoryStrategy::Bounded: the bounded engine \
cannot edit a file with a non-zero base address, so creating one this way would \
write the file and then refuse to open it; drop the userblock, or leave \
MemoryStrategy unset to mirror this file",
);
}
None
}
struct PagedEdit {
page_size: u64,
meta: FreeList,
raw_small: FreeList,
raw_large: FreeList,
last: Option<PageType>,
meta_pad: Vec<(u64, u64)>,
raw_pad: Vec<(u64, u64)>,
}
impl PagedEdit {
fn begin(&mut self, image: &mut dyn FileImage, ty: PageType) -> Result<(), Error> {
let len = image.len();
if len % self.page_size != 0 {
let pad_len = self.page_size - len % self.page_size;
let pad = match self.last {
Some(prev) if prev != ty => Some(Some(prev)),
Some(_) => None, None => Some(None),
};
if let Some(prev) = pad {
let pad_at = len;
image.append(&vec![0u8; pad_len.to_usize()?])?;
match prev {
Some(PageType::Meta) => self.meta_pad.push((pad_at, pad_len)),
Some(PageType::Raw) => self.raw_pad.push((pad_at, pad_len)),
None => {} }
}
}
self.last = Some(ty);
Ok(())
}
fn new(page_size: u64) -> Self {
PagedEdit {
page_size,
meta: FreeList::new(),
raw_small: FreeList::new(),
raw_large: FreeList::new(),
last: None,
meta_pad: Vec::new(),
raw_pad: Vec::new(),
}
}
fn route_free(
meta: &mut FreeList,
raw_small: &mut FreeList,
raw_large: &mut FreeList,
page_size: u64,
addr: u64,
size: u64,
ty: PageType,
) {
match ty {
PageType::Meta => meta.free(addr, size),
PageType::Raw if size >= page_size => raw_large.free(addr, size),
PageType::Raw => raw_small.free(addr, size),
}
}
fn all_sections(&self) -> Vec<(u64, u64)> {
let mut out = self.meta.sections();
out.extend(self.raw_small.sections());
out.extend(self.raw_large.sections());
out.sort_by_key(|&(addr, _)| addr);
out
}
}
#[derive(Clone, Copy)]
pub(crate) enum AppendTarget<'a> {
Path(&'a str),
Header(u64),
}
const APPEND_BATCH_BYTES: u64 = 1 << 20;
pub(crate) struct AppendGeometry {
pub(crate) chunk_elems: u64,
pub(crate) element_size: usize,
pub(crate) current_dim: u64,
pub(crate) filtered: bool,
pub(crate) full_batch_elems: u64,
}
const SWMR_WRITE_FLAGS: u32 = file_lock::WRITE_ACCESS | file_lock::SWMR_WRITE_ACCESS;
pub(crate) struct LocatedState {
pub(crate) loc: Located,
pub(crate) datatype: Datatype,
pub(crate) spatial: Vec<u64>,
pub(crate) element_size: usize,
pub(crate) pipeline: Option<FilterPipeline>,
}
struct PersistState {
strategy: FileSpaceStrategy,
threshold: u64,
page_size: u64,
old_blocks: Vec<(u64, u64)>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct SpaceAccounting {
pub logical_size: u64,
pub reusable_free_bytes: u64,
pub reusable_free_space: Vec<(u64, u64)>,
}
impl WriteEngine {
pub fn open_with_locking<P: AsRef<Path>>(path: P, locking: FileLocking) -> Result<Self, Error> {
Self::open_inner(path.as_ref(), Some(locking))
}
#[cfg(test)]
pub(crate) fn open_source_only(path: &Path) -> Result<Self, Error> {
Self::open_imaged(path, Some(FileLocking::Enabled), |handle, _len| {
Ok(Box::new(crate::image::SourceOnlyImage::new(
Self::read_mirror(handle)?,
)))
})
}
#[cfg(test)]
pub(crate) fn open_bounded_counting(
path: &Path,
read_bytes: std::sync::Arc<std::sync::atomic::AtomicU64>,
) -> Result<Self, Error> {
let mut session = Self::open_imaged(path, Some(FileLocking::Enabled), |handle, len| {
Ok(Box::new(crate::image::CountingImage::new(
Box::new(HandleImage::new(
handle,
len,
crate::source::MetadataCacheConfig::disabled(),
)),
read_bytes,
)))
})?;
session.batched_appends = true;
Ok(session)
}
pub(crate) fn open_rw_with_strategy(
path: &Path,
cache: MetadataCacheConfig,
locking: FileLocking,
strategy: MemoryStrategy,
) -> Result<Self, Error> {
if strategy == MemoryStrategy::Mirrored {
return Self::open_with_locking(path, locking);
}
let mut session = Self::open_imaged(path, Some(locking), |handle, len| {
Ok(Box::new(HandleImage::new(handle, len, cache)))
})?;
session.batched_appends = true;
session.bounded = true;
if session.paged.is_some() && session.persist.is_none() {
return Err(Error::EditUnsupported(
"read-write access to a paged file (H5F_FSPACE_STRATEGY_PAGE) requires \
persisted free space; recreate the file with \
with_file_space_strategy(FileSpaceStrategy::Page, true, ..) to grow it in place",
));
}
if let Some(reason) = bounded_only_limitation(&session) {
if strategy == MemoryStrategy::Bounded {
return Err(Error::EditUnsupported(reason));
}
drop(session);
return Self::open_with_locking(path, locking);
}
Ok(session)
}
pub(crate) fn open_swmr_writer<P: AsRef<Path>>(path: P) -> Result<Self, Error> {
let mut session = Self::open_inner(path.as_ref(), None)?;
if session.superblock.version < 3
|| session.superblock.base_address != 0
|| session.persist.is_some()
{
return Err(Error::SwmrAppendUnsupported(
"SWMR writing requires a latest-format file (v3 superblock) with no userblock \
and no persisted free-space",
));
}
session.swmr_mode = true;
session.set_consistency_flags(SWMR_WRITE_FLAGS)?;
Ok(session)
}
pub(crate) fn set_consistency_flags(&mut self, flags: u32) -> Result<(), Error> {
self.superblock.consistency_flags = flags;
let bytes = self.superblock.serialize();
self.write_at(self.sb_sig_off, &bytes)?;
self.image.sync_data()?;
Ok(())
}
fn open_inner(path: &Path, lock: Option<FileLocking>) -> Result<Self, Error> {
Self::open_imaged(path, lock, |handle, _len| {
Ok(Box::new(Self::read_mirror(handle)?))
})
}
fn read_mirror(mut handle: fs::File) -> Result<MirrorImage, Error> {
handle.seek(SeekFrom::Start(0)).map_err(Error::Io)?;
let mut data = Vec::new();
handle.read_to_end(&mut data).map_err(Error::Io)?;
Ok(MirrorImage::new(handle, data))
}
fn open_imaged(
path: &Path,
lock: Option<FileLocking>,
build: impl FnOnce(fs::File, u64) -> Result<Box<dyn FileImage>, Error>,
) -> Result<Self, Error> {
let handle = fs::OpenOptions::new()
.read(true)
.write(true)
.open(path)
.map_err(Error::Io)?;
if let Some(policy) = lock {
file_lock::acquire_exclusive(&handle, policy, path)?;
}
let len = handle.metadata().map_err(Error::Io)?.len();
let probe = crate::image::BorrowedHandle::new(&handle, len);
let sb_sig_off = signature::find_signature_in(&probe)?.to_usize()?;
let mut superblock = Superblock::parse_from_source(&probe, sb_sig_off as u64)?;
if superblock.version > 3 {
return Err(Error::EditUnsupported("unsupported superblock version"));
}
file_lock::check_status_flags(&superblock, file_lock::OpenIntent::Write, path)?;
if superblock.offset_size != OFFSET_SIZE || superblock.length_size != LENGTH_SIZE {
return Err(Error::EditUnsupported(
"only 8-byte offsets and lengths are supported for in-place editing",
));
}
if superblock.base_address != sb_sig_off as u64 {
return Err(Error::EditUnsupported(
"a file whose superblock is not located at its base address is not editable in place",
));
}
superblock.root_group_address = superblock
.root_group_address
.checked_add(superblock.base_address)
.ok_or(FormatError::OffsetOverflow {
offset: superblock.root_group_address,
length: superblock.base_address,
})?;
let image = build(handle, len)?;
let mut session = Self {
image,
sb_sig_off,
superblock,
pending_datasets: Vec::new(),
pending_writes: Vec::new(),
pending_appends: Vec::new(),
pending_groups: Vec::new(),
pending_group_attrs: Vec::new(),
pending_dataset_attrs: Vec::new(),
pending_deletes: Vec::new(),
pending_copies: Vec::new(),
pending_cross_copies: Vec::new(),
free: FreeList::new(),
persist: None,
located: HashMap::new(),
swmr_mode: false,
paged: None,
committed: false,
resolved: HashMap::new(),
batched_appends: false,
bounded: false,
fsm_len: len,
};
session.load_persisted_free_space();
Ok(session)
}
fn load_persisted_free_space(&mut self) {
if self.superblock.version < 2 {
return; }
let Some(ext_rel) = self.superblock.superblock_extension_address else {
return;
};
if ext_rel == UNDEF {
return;
}
let Ok(ext_addr) = ext_rel
.checked_add(self.superblock.base_address)
.ok_or(())
.and_then(|a| usize::try_from(a).map_err(|_| ()))
else {
return;
};
let Some(info) = self.extension_fsinfo(ext_addr) else {
return;
};
if self.superblock.base_address != 0 {
if info.strategy == FileSpaceStrategy::Page && info.page_size > 0 {
self.paged = Some(PagedEdit::new(info.page_size));
}
return;
}
let paged = info.strategy == FileSpaceStrategy::Page && info.page_size > 0;
if paged {
self.paged = Some(PagedEdit::new(info.page_size));
}
if !info.persist {
return;
}
let os = self.superblock.offset_size;
let file_len = self.image.len();
if paged {
let mut tagged: Vec<(FreeSection, PageType, bool)> = Vec::new();
for (slot, &m) in info.manager_addrs.iter().enumerate() {
if m == UNDEF {
continue;
}
let Ok(sections) =
free_space_manager::read_persisted_sections_source(&self.image(), &[m], 0, os)
.map(|(sections, _)| sections)
else {
continue;
};
let (ty, large) = match slot {
2 => (PageType::Raw, false),
6 => (PageType::Raw, true),
_ => (PageType::Meta, false),
};
for s in sections {
tagged.push((s, ty, large));
}
}
tagged.sort_by_key(|(s, _, _)| s.addr);
let mut prev_end = 0u64;
for (s, ty, large) in tagged {
let Some(end) = s.addr.checked_add(s.size) else {
continue;
};
if s.size == 0 || end > file_len || s.addr < prev_end {
continue;
}
prev_end = end;
let pg = self
.paged
.as_mut()
.expect("the paged state was just installed");
if large {
pg.raw_large.free(s.addr, s.size);
} else {
match ty {
PageType::Meta => pg.meta.free(s.addr, s.size),
PageType::Raw => pg.raw_small.free(s.addr, s.size),
}
}
}
} else if let Ok(mut sections) = free_space_manager::read_persisted_sections_source(
&self.image(),
&info.manager_addrs,
0,
os,
)
.map(|(sections, _)| sections)
{
sections.sort_by_key(|s| s.addr);
let mut prev_end = 0u64;
for s in sections {
let Some(end) = s.addr.checked_add(s.size) else {
continue;
};
if s.size == 0 || end > file_len || s.addr < prev_end {
continue;
}
prev_end = end;
self.free.free(s.addr, s.size);
}
}
let mut old_blocks = Vec::new();
if let Ok(spans) = self.oh_chunk_spans(ext_addr) {
old_blocks.extend(spans);
}
for &m in &info.manager_addrs {
if m == UNDEF {
continue;
}
let Ok(hdr_len) = fshd_len(os).to_usize() else {
continue;
};
let Ok(fshd) = self.image().read_metadata_at(m, hdr_len) else {
continue;
};
if let Ok(h) = FsmHeader::parse(&fshd, os) {
old_blocks.push((m, fshd_len(os)));
if h.fsse_addr != UNDEF
&& h.fsse_addr
.checked_add(h.fsse_used)
.is_some_and(|end| end <= file_len)
{
old_blocks.push((h.fsse_addr, h.fsse_used));
}
}
}
self.persist = Some(PersistState {
strategy: info.strategy,
threshold: info.threshold,
page_size: info.page_size,
old_blocks,
});
}
fn extension_fsinfo(&self, ext_addr: usize) -> Option<FileSpaceInfo> {
let os = self.superblock.offset_size;
let ls = self.superblock.length_size;
let base = self.superblock.base_address;
let oh =
ObjectHeader::parse_from_source(&self.image(), ext_addr as u64, os, ls, base).ok()?;
let msg = oh
.messages
.iter()
.find(|m| m.msg_type == MessageType::FileSpaceInfo)?;
FileSpaceInfo::parse(&msg.data, os, ls).ok()
}
pub(crate) fn stage_created_dataset(&mut self, path: &str, mut builder: DatasetBuilder) {
let mut comps = split_path(path);
builder.name = comps.pop().unwrap_or_default();
self.pending_datasets.push((comps, builder));
}
pub(crate) fn stage_dataset_write(&mut self, path: &str, mut builder: DatasetBuilder) {
let comps = split_path(path);
builder.name = comps.last().cloned().unwrap_or_default();
self.pending_writes.push((comps, builder));
}
pub(crate) fn stage_dataset_append(&mut self, path: &str, builder: AppendBuilder) {
self.pending_appends.push((split_path(path), builder));
}
pub fn has_staged_edits(&self) -> bool {
!self.pending_datasets.is_empty()
|| !self.pending_writes.is_empty()
|| !self.pending_appends.is_empty()
|| !self.pending_groups.is_empty()
|| !self.pending_group_attrs.is_empty()
|| !self.pending_dataset_attrs.is_empty()
|| !self.pending_deletes.is_empty()
|| !self.pending_copies.is_empty()
|| !self.pending_cross_copies.is_empty()
}
pub(crate) fn image_slice(&self) -> Option<&[u8]> {
self.image.as_slice()
}
pub(crate) fn image(&self) -> &dyn Source {
self.image.as_ref()
}
pub(crate) fn superblock(&self) -> &Superblock {
&self.superblock
}
pub(crate) fn edit_backing(&self) -> EditBacking {
if self.bounded {
EditBacking::Bounded
} else {
EditBacking::Mirrored
}
}
#[must_use]
pub fn space_accounting(&self) -> SpaceAccounting {
let reusable_free_space = match &self.paged {
Some(pg) => pg.all_sections(),
None => self.free.sections(),
};
let reusable_free_bytes = reusable_free_space.iter().map(|(_, len)| len).sum();
SpaceAccounting {
logical_size: self.image.len(),
reusable_free_bytes,
reusable_free_space,
}
}
fn store(&mut self) -> EditStore<'_> {
EditStore {
image: self.image.as_mut(),
superblock: &mut self.superblock,
sb_sig_off: self.sb_sig_off,
paged: self.paged.as_mut(),
}
}
pub(crate) fn sync(&mut self) -> Result<(), Error> {
self.image.sync_all()
}
pub(crate) fn finalize_persist(&mut self) -> Result<(), Error> {
if self.persist.is_none() || self.image.len() == self.fsm_len {
return Ok(());
}
self.commit_persisting(self.superblock.root_group_address, Vec::new())
}
fn append_prepare(&mut self, target: AppendTarget<'_>) -> Result<u64, Error> {
if self.superblock.base_address != 0 {
return Err(Error::AppendInPlaceUnsupported(
"in-place append does not support a file with a userblock (non-zero base \
address); use Dataset::append_staged",
));
}
if self.superblock.version < 2 {
return Err(Error::AppendInPlaceUnsupported(
"in-place append requires a latest-format file (v2/v3 superblock); use \
Dataset::append_staged",
));
}
if self.paged.is_some() && self.persist.is_none() {
return Err(Error::AppendInPlaceUnsupported(
"in-place append is not supported on a paged file \
(H5F_FSPACE_STRATEGY_PAGE) without persisted free space; recreate the \
file with with_file_space_strategy(FileSpaceStrategy::Page, true, ..)",
));
}
match target {
AppendTarget::Path(dataset) => {
if self.append_conflicts_with_pending(&split_path(dataset)) {
return Err(Error::AppendInPlaceUnsupported(
"the dataset or an ancestor has a staged edit pending in this session; \
commit the staged edits before appending in place, or use \
Dataset::append_staged",
));
}
}
AppendTarget::Header(_) if self.has_staged_edits() || self.committed => {
return Err(Error::AppendInPlaceUnsupported(
"this append target was reached by object reference, so it names a dataset \
by object-header address, and this session has staged or committed edits \
that can move that header; re-open the dataset by path to append to it",
));
}
AppendTarget::Header(_) => {}
}
let oh_addr = match target {
AppendTarget::Path(dataset) => match self.resolved.get(dataset) {
Some(&addr) => addr,
None => {
let addr = crate::group_v2::resolve_path_any_from_source(
&self.image(),
&self.superblock,
dataset,
)
.map_err(|_| {
Error::AppendInPlaceUnsupported("nothing to append to at the given path")
})?;
self.resolved.insert(dataset.to_string(), addr);
addr
}
},
AppendTarget::Header(addr) => addr,
};
if !self.located.contains_key(&oh_addr) {
let store = self.store();
let state = locate_dataset_state(&store, oh_addr)?;
self.located.insert(oh_addr, state);
}
Ok(oh_addr)
}
pub(crate) fn append_geometry(
&mut self,
target: AppendTarget<'_>,
) -> Result<AppendGeometry, Error> {
let oh_addr = self.append_prepare(target)?;
let st = &self.located[&oh_addr];
let chunk_elems = st.loc.chunk_elems.max(1);
Ok(AppendGeometry {
chunk_elems,
element_size: st.element_size,
current_dim: st.loc.current_dim,
filtered: st.pipeline.is_some(),
full_batch_elems: self.batch_elems(st.loc.chunk_bytes, chunk_elems),
})
}
fn batch_elems(&self, chunk_bytes: usize, chunk_elems: u64) -> u64 {
if !self.batched_appends {
return u64::MAX;
}
(APPEND_BATCH_BYTES / (chunk_bytes.max(1) as u64)).max(1) * chunk_elems
}
pub(crate) fn append_inplace_gathered(
&mut self,
target: AppendTarget<'_>,
b: &AppendBuilder,
max_phase: u8,
) -> Result<(), Error> {
if b.dt_conflict() {
return Err(Error::AppendInPlaceUnsupported(
"append mixes element types in one call; use one element type per append",
));
}
let oh_addr = self.append_prepare(target)?;
let raw = b.raw();
let new_elems = validate_gathered_append(&self.located[&oh_addr], b)?;
if new_elems == 0 {
return Ok(());
}
if self.swmr_mode {
let st = &self.located[&oh_addr];
if st.pipeline.is_some() {
return Err(Error::SwmrAppendUnsupported(
"filtered datasets are not supported for SWMR append",
));
}
let chunk_elems = st.loc.chunk_elems;
if chunk_elems == 0
|| st.loc.current_dim % chunk_elems != 0
|| new_elems % chunk_elems != 0
{
return Err(Error::SwmrAppendUnsupported(
"SWMR append must be chunk-aligned: the current length and the appended \
length must both be whole multiples of the chunk length",
));
}
}
let (chunk_elems, elem_bytes, full_batch_elems, filtered, current_dim) = {
let st = &self.located[&oh_addr];
(
st.loc.chunk_elems.max(1),
st.element_size as u64,
self.batch_elems(st.loc.chunk_bytes, st.loc.chunk_elems.max(1)),
st.pipeline.is_some(),
st.loc.current_dim,
)
};
if filtered && (current_dim % chunk_elems != 0 || new_elems % chunk_elems != 0) {
return Err(Error::AppendInPlaceUnsupported(
"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 done = 0u64;
while done < new_elems {
let current_dim = self.located[&oh_addr].loc.current_dim;
let to_boundary = (chunk_elems - current_dim % chunk_elems) % chunk_elems;
let take = (new_elems - done).min(to_boundary.saturating_add(full_batch_elems));
let batch =
&raw[(done * elem_bytes).to_usize()?..((done + take) * elem_bytes).to_usize()?];
let plan_result = {
let Self {
image,
superblock,
sb_sig_off,
paged,
located,
..
} = self;
let st = &located[&oh_addr];
let store = EditStore {
image: image.as_mut(),
superblock,
sb_sig_off: *sb_sig_off,
paged: paged.as_mut(),
};
plan_ea_append(
&store,
&st.loc,
&st.datatype,
&st.spatial,
st.element_size,
st.pipeline.as_ref(),
batch,
take,
)
};
let plan = plan_result.map_err(as_inplace_error)?;
{
let Self {
image,
superblock,
sb_sig_off,
paged,
located,
..
} = self;
let st = located.get_mut(&oh_addr).expect("dataset located above");
let mut store = EditStore {
image: image.as_mut(),
superblock,
sb_sig_off: *sb_sig_off,
paged: paged.as_mut(),
};
apply_ea_append(&mut store, &mut st.loc, &plan, max_phase)
.map_err(as_inplace_error)?;
}
if max_phase < 4 {
return Ok(());
}
done += take;
}
Ok(())
}
#[cfg(test)]
fn append_inplace_i32_phased(
&mut self,
dataset: &str,
values: &[i32],
max_phase: u8,
) -> Result<(), Error> {
let mut b = AppendBuilder::new();
b.append_i32(values);
self.append_inplace_gathered(AppendTarget::Path(dataset), &b, max_phase)
}
fn append_conflicts_with_pending(&self, target: &[String]) -> bool {
let hits = |p: &[String]| paths_overlap(target, p);
self.pending_writes.iter().any(|(p, _)| hits(p))
|| self.pending_appends.iter().any(|(p, _)| hits(p))
|| self.pending_deletes.iter().any(|p| hits(p))
|| self.pending_copies.iter().any(|(_, dst)| hits(dst))
|| self.pending_cross_copies.iter().any(|(dst, _)| hits(dst))
|| self.pending_dataset_attrs.iter().any(|(p, _)| hits(p))
|| self.pending_datasets.iter().any(|(parent, db)| {
let mut full = parent.clone();
full.push(db.name.clone());
paths_overlap(target, &full)
})
}
pub fn create_group(&mut self, path: &str) {
self.pending_groups.push(split_path(path));
}
pub fn set_group_attr(&mut self, path: &str, name: &str, value: AttrValue) -> &mut Self {
self.pending_group_attrs.push((
split_path(path),
AttrOp::Set {
name: name.to_string(),
value,
},
));
self
}
pub fn remove_group_attr(&mut self, path: &str, name: &str) -> &mut Self {
self.pending_group_attrs.push((
split_path(path),
AttrOp::Remove {
name: name.to_string(),
},
));
self
}
pub fn set_dataset_attr(&mut self, path: &str, name: &str, value: AttrValue) -> &mut Self {
self.pending_dataset_attrs.push((
split_path(path),
AttrOp::Set {
name: name.to_string(),
value,
},
));
self
}
pub fn remove_dataset_attr(&mut self, path: &str, name: &str) -> &mut Self {
self.pending_dataset_attrs.push((
split_path(path),
AttrOp::Remove {
name: name.to_string(),
},
));
self
}
pub fn delete(&mut self, path: &str) {
self.pending_deletes.push(split_path(path));
}
pub fn copy(&mut self, src: &str, dst: &str) {
self.pending_copies.push((split_path(src), split_path(dst)));
}
pub fn copy_from(
&mut self,
source: &crate::reader::File,
src: &str,
dst: &str,
) -> Result<(), Error> {
let src_data = source.in_memory_image().ok_or(Error::EditUnsupported(
"cross-file copy requires a buffered source file (File::open or File::from_bytes), not a streaming one",
))?;
let src_sb = source.superblock();
if src_sb.offset_size != OFFSET_SIZE || src_sb.length_size != LENGTH_SIZE {
return Err(Error::EditUnsupported(
"cross-file copy requires the source file to use 8-byte offsets and lengths",
));
}
if source.base_address() != 0 {
return Err(Error::EditUnsupported(
"cross-file copy requires the source file to have no userblock (base address 0)",
));
}
let src = split_path(src);
if src.is_empty() {
return Err(Error::EditUnsupported("cannot copy the root group"));
}
let dst = split_path(dst);
if dst.is_empty() {
return Err(Error::EditUnsupported("copy destination path is empty"));
}
let src_addr = crate::group_v2::resolve_path_any(src_data, src_sb, &src.join("/"))
.map_err(|_| Error::EditUnsupported("copy source does not exist in the source file"))?;
let tree = Self::read_copy_subtree(&BytesSource::new(src_data), src_addr, 0, true, 0)?;
self.pending_cross_copies.push((dst, tree));
Ok(())
}
pub fn commit(&mut self) -> Result<(), Error> {
if self.pending_datasets.is_empty()
&& self.pending_writes.is_empty()
&& self.pending_appends.is_empty()
&& self.pending_groups.is_empty()
&& self.pending_group_attrs.is_empty()
&& self.pending_dataset_attrs.is_empty()
&& self.pending_deletes.is_empty()
&& self.pending_copies.is_empty()
&& self.pending_cross_copies.is_empty()
{
return Ok(());
}
if self.paged.is_some() && self.persist.is_none() {
return Err(Error::EditUnsupported(
"committing an edit to a paged file (H5F_FSPACE_STRATEGY_PAGE) requires \
persisted free space; recreate the file with \
with_file_space_strategy(FileSpaceStrategy::Page, true, ..) to edit it in place",
));
}
self.located.clear();
self.resolved.clear();
self.committed = true;
let base = self.superblock.base_address;
let writes = std::mem::take(&mut self.pending_writes);
let mut inplace_writes: Vec<(usize, Vec<u8>)> = Vec::new();
let mut moving_writes: Vec<(PathKey, String, MovingWrite)> = Vec::new();
let mut write_targets: Vec<PathKey> = Vec::new();
let mut incoming_links: Option<Option<HashMap<u64, u32>>> = None;
for (full, db) in writes {
if full.is_empty() {
return Err(Error::EditUnsupported("cannot overwrite the root group"));
}
if write_targets.contains(&full) {
return Err(Error::EditUnsupported(
"the same dataset is overwritten twice in one commit; use separate commits",
));
}
let path_str = full.join("/");
let addr = crate::group_v2::resolve_path_any_from_source(
&self.image(),
&self.superblock,
&path_str,
)
.map_err(|_| Error::EditUnsupported("nothing to overwrite at the given path"))?;
let addr = usize::try_from(addr)
.map_err(|_| Error::EditUnsupported("dataset address exceeds this platform"))?;
let fd = flatten_dataset(db)?;
match Self::prepare_write(&self.image(), addr as u64, &fd, base)? {
WritePlan::InPlace { data_addr, raw } => inplace_writes.push((data_addr, raw)),
WritePlan::InPlaceChunks { writes } => inplace_writes.extend(writes),
WritePlan::Moving(mw) => {
let counts = incoming_links
.get_or_insert_with(|| self.count_incoming_hard_links())
.as_ref();
match counts.and_then(|c| c.get(&(addr as u64))) {
Some(&1) => {}
_ => {
return Err(Error::EditUnsupported(
"overwriting a dataset that resizes or relocates its header is \
only supported when it has a single hard link",
));
}
}
let leaf = full.last().unwrap().clone();
let parent = full[..full.len() - 1].to_vec();
moving_writes.push((parent, leaf, mw));
}
}
write_targets.push(full);
}
let appends = std::mem::take(&mut self.pending_appends);
for (full, ab) in appends {
if full.is_empty() {
return Err(Error::AppendUnsupported("cannot append to the root group"));
}
if ab.raw.is_empty() {
continue; }
if write_targets.contains(&full) {
return Err(Error::AppendUnsupported(
"the same dataset is edited more than once in one commit; use separate commits",
));
}
let path_str = full.join("/");
let addr = crate::group_v2::resolve_path_any_from_source(
&self.image(),
&self.superblock,
&path_str,
)
.map_err(|_| Error::AppendUnsupported("nothing to append to at the given path"))?;
let addr = usize::try_from(addr)
.map_err(|_| Error::AppendUnsupported("dataset address exceeds this platform"))?;
let mw = Self::prepare_append(&self.image(), addr as u64, &ab, base)?;
let counts = incoming_links
.get_or_insert_with(|| self.count_incoming_hard_links())
.as_ref();
match counts.and_then(|c| c.get(&(addr as u64))) {
Some(&1) => {}
_ => {
return Err(Error::AppendUnsupported(
"appending relocates the dataset header; only supported when it \
has a single hard link",
));
}
}
let leaf = full.last().unwrap().clone();
let parent = full[..full.len() - 1].to_vec();
moving_writes.push((parent, leaf, mw));
write_targets.push(full);
}
let dataset_attrs = std::mem::take(&mut self.pending_dataset_attrs);
if !dataset_attrs.is_empty() {
let mut order: Vec<PathKey> = Vec::new();
let mut ops_by_path: HashMap<PathKey, Vec<AttrOp>> = HashMap::new();
for (path, op) in dataset_attrs {
if !ops_by_path.contains_key(&path) {
order.push(path.clone());
}
ops_by_path.entry(path).or_default().push(op);
}
for full in order {
let ops = ops_by_path.remove(&full).unwrap();
if full.is_empty() {
return Err(Error::EditUnsupported(
"cannot set a dataset attribute on the root group; use set_group_attr",
));
}
if write_targets.contains(&full) {
return Err(Error::EditUnsupported(
"the same dataset is edited more than once in one commit (an attribute \
edit plus another edit); use separate commits",
));
}
let path_str = full.join("/");
let addr = crate::group_v2::resolve_path_any_from_source(
&self.image(),
&self.superblock,
&path_str,
)
.map_err(|_| {
Error::EditUnsupported("nothing to set an attribute on at the given path")
})?;
let addr = usize::try_from(addr)
.map_err(|_| Error::EditUnsupported("dataset address exceeds this platform"))?;
let counts = incoming_links
.get_or_insert_with(|| self.count_incoming_hard_links())
.as_ref();
match counts.and_then(|c| c.get(&(addr as u64))) {
Some(&1) => {}
_ => {
return Err(Error::EditUnsupported(
"editing a dataset attribute relocates its header; only supported \
when it has a single hard link",
));
}
}
let region = Self::gather_oh_messages(&self.image(), addr as u64, base)?;
let (region, pending_vl_attrs) = apply_group_attr_ops(®ion, &ops)?;
let leaf = full.last().unwrap().clone();
let parent = full[..full.len() - 1].to_vec();
moving_writes.push((
parent,
leaf,
MovingWrite::AttrEdit {
region,
pending_vl_attrs,
},
));
write_targets.push(full);
}
}
if moving_writes.is_empty()
&& self.pending_datasets.is_empty()
&& self.pending_groups.is_empty()
&& self.pending_group_attrs.is_empty()
&& self.pending_deletes.is_empty()
&& self.pending_copies.is_empty()
&& self.pending_cross_copies.is_empty()
{
for (data_addr, raw) in &inplace_writes {
self.write_at(*data_addr, raw)?;
}
self.image.sync_all()?;
return Ok(());
}
let mut nodes: BTreeMap<PathKey, Node> = BTreeMap::new();
nodes.entry(PathKey::new()).or_default(); let mut add_targets: Vec<PathKey> = Vec::new();
let mut attr_targets: Vec<PathKey> = Vec::new();
for path in std::mem::take(&mut self.pending_groups) {
if path.is_empty() {
return Err(Error::EditUnsupported("cannot create the root group"));
}
ensure_ancestors(&mut nodes, &path);
nodes.entry(path.clone()).or_default().is_new = true;
add_targets.push(path);
}
for (parent, db) in std::mem::take(&mut self.pending_datasets) {
let mut full = parent.clone();
full.push(db.name.clone());
add_targets.push(full);
ensure_ancestors(&mut nodes, &parent);
nodes.entry(parent).or_default().datasets.push(db);
}
for (parent, leaf, mw) in moving_writes {
ensure_ancestors(&mut nodes, &parent);
nodes.entry(parent).or_default().writes.push((leaf, mw));
}
for (path, op) in std::mem::take(&mut self.pending_group_attrs) {
ensure_ancestors(&mut nodes, &path);
nodes.entry(path.clone()).or_default().attr_ops.push(op);
attr_targets.push(path);
}
for (src, dst) in std::mem::take(&mut self.pending_copies) {
if src.is_empty() {
return Err(Error::EditUnsupported("cannot copy the root group"));
}
if dst.is_empty() {
return Err(Error::EditUnsupported("copy destination path is empty"));
}
if is_prefix(&src, &dst) {
return Err(Error::EditUnsupported(
"cannot copy an object into itself or its own subtree",
));
}
let src_str = src.join("/");
let src_addr = crate::group_v2::resolve_path_any_from_source(
&self.image(),
&self.superblock,
&src_str,
)
.map_err(|_| Error::EditUnsupported("copy source does not exist"))?;
let src_addr = usize::try_from(src_addr)
.map_err(|_| Error::EditUnsupported("source address exceeds this platform"))?;
let tree = Self::read_copy_subtree(&self.image(), src_addr as u64, 0, false, base)?;
add_targets.push(dst.clone());
let leaf = dst.last().unwrap().clone();
let parent = dst[..dst.len() - 1].to_vec();
ensure_ancestors(&mut nodes, &parent);
nodes.entry(parent).or_default().copies.push((leaf, tree));
}
for (dst, tree) in std::mem::take(&mut self.pending_cross_copies) {
if dst.is_empty() {
return Err(Error::EditUnsupported("copy destination path is empty"));
}
add_targets.push(dst.clone());
let leaf = dst.last().unwrap().clone();
let parent = dst[..dst.len() - 1].to_vec();
ensure_ancestors(&mut nodes, &parent);
nodes.entry(parent).or_default().copies.push((leaf, tree));
}
let delete_targets = std::mem::take(&mut self.pending_deletes);
let mut deleted_addrs: Vec<usize> = Vec::new();
for (i, d) in delete_targets.iter().enumerate() {
if d.is_empty() {
return Err(Error::EditUnsupported("cannot delete the root group"));
}
let path_str = d.join("/");
let del_addr = crate::group_v2::resolve_path_any_from_source(
&self.image(),
&self.superblock,
&path_str,
)
.map_err(|_| Error::EditUnsupported("nothing to delete at the given path"))?;
if let Ok(a) = usize::try_from(del_addr) {
deleted_addrs.push(a);
}
for t in &add_targets {
if is_prefix(d, t) || is_prefix(t, d) {
return Err(Error::EditUnsupported(
"a deletion overlaps an addition in the same commit; use separate commits",
));
}
}
for t in &attr_targets {
if is_prefix(d, t) {
return Err(Error::EditUnsupported(
"a deletion overlaps a group-attribute edit in the same commit; use separate commits",
));
}
}
for t in &write_targets {
if is_prefix(d, t) {
return Err(Error::EditUnsupported(
"a deletion overlaps a value overwrite in the same commit; use separate commits",
));
}
}
for (j, d2) in delete_targets.iter().enumerate() {
if i != j && is_prefix(d, d2) {
return Err(Error::EditUnsupported(
"overlapping deletions in one commit; delete the common parent only",
));
}
}
let parent = d[..d.len() - 1].to_vec();
ensure_ancestors(&mut nodes, &parent);
nodes
.entry(parent)
.or_default()
.deletes
.push(d.last().unwrap().clone());
}
let keys: Vec<PathKey> = nodes.keys().cloned().collect();
let mut superseded_addrs: Vec<usize> = Vec::new();
for key in &keys {
let is_new = nodes[key].is_new;
if is_new {
nodes.get_mut(key).unwrap().base_region = fresh_group_region();
} else {
let path_str = key.join("/");
let addr = crate::group_v2::resolve_path_any_from_source(
&self.image(),
&self.superblock,
&path_str,
)
.map_err(|_| {
Error::EditUnsupported(
"a target group does not exist; create it first in this session",
)
})?;
let addr = usize::try_from(addr)
.map_err(|_| Error::EditUnsupported("group address exceeds this platform"))?;
let info = self.inspect_group(addr)?;
superseded_addrs.push(addr);
let node = nodes.get_mut(key).unwrap();
node.base_region = info.region;
node.existing_links = info.link_names;
}
}
for key in &keys {
let node = nodes.get_mut(key).unwrap();
let ops = std::mem::take(&mut node.attr_ops);
if !ops.is_empty() {
let region = std::mem::take(&mut node.base_region);
let (region, pending_vl_attrs) = apply_group_attr_ops(®ion, &ops)?;
node.base_region = region;
node.pending_vl_attrs = pending_vl_attrs;
}
}
let mut children: BTreeMap<PathKey, Vec<PathKey>> = BTreeMap::new();
for key in &keys {
if !key.is_empty() {
let parent = key[..key.len() - 1].to_vec();
children.entry(parent).or_default().push(key.clone());
}
}
for key in &keys {
let node = &nodes[key];
let mut adding: Vec<&str> = Vec::new();
for db in &node.datasets {
adding.push(&db.name);
}
for child in children.get(key).into_iter().flatten() {
if nodes[child].is_new {
adding.push(child.last().unwrap());
}
}
for (leaf, _) in &node.copies {
adding.push(leaf);
}
for (i, name) in adding.iter().enumerate() {
if node.existing_links.iter().any(|n| n == name) || adding[..i].contains(name) {
return Err(Error::EditUnsupported(
"a link with this name already exists in the target group",
));
}
}
}
let mut flat: BTreeMap<PathKey, Vec<FlatDataset>> = BTreeMap::new();
for key in &keys {
let dbs = std::mem::take(&mut nodes.get_mut(key).unwrap().datasets);
let mut v = Vec::with_capacity(dbs.len());
for db in dbs {
v.push(flatten_dataset(db)?);
}
flat.insert(key.clone(), v);
}
Self::preflight_reference_targets(
&keys,
&flat,
&nodes,
&add_targets,
&write_targets,
&delete_targets,
&self.image(),
&self.superblock,
)?;
let mut to_free: Vec<(u64, u64, PageType)> = Vec::new();
deleted_addrs.sort_unstable();
deleted_addrs.dedup();
if !deleted_addrs.is_empty() {
if let Some(incoming) = self.count_incoming_hard_links() {
for &a in &deleted_addrs {
self.collect_free_spans(a, 0, &incoming, &mut to_free);
}
}
}
for &a in &superseded_addrs {
if let Ok(spans) = self.oh_chunk_spans(a) {
to_free.extend(spans.into_iter().map(|(a, l)| (a, l, PageType::Meta)));
}
}
for key in &keys {
for (leaf, mw) in &nodes[key].writes {
match mw {
MovingWrite::Contiguous {
old_extent: Some(extent),
..
} => to_free.push((extent.0, extent.1, PageType::Raw)),
MovingWrite::Chunked { old_addr, .. } => {
if let Ok(a) = usize::try_from(*old_addr) {
if let Some(spans) = self.chunked_storage_spans(a) {
to_free.extend(spans);
}
}
}
MovingWrite::AppendedChunks {
old_addr,
old_tail_extent,
..
} => {
if let Ok(a) = usize::try_from(*old_addr) {
if let Some(spans) = self.chunked_index_spans(a) {
to_free
.extend(spans.into_iter().map(|(a, l)| (a, l, PageType::Raw)));
}
}
if let Some(ext) = old_tail_extent {
to_free.push((ext.0, ext.1, PageType::Raw));
}
}
_ => {}
}
let mut full = key.clone();
full.push(leaf.clone());
let path_str = full.join("/");
if let Ok(addr) = crate::group_v2::resolve_path_any_from_source(
&self.image(),
&self.superblock,
&path_str,
) {
if let Ok(a) = usize::try_from(addr) {
if let Ok(spans) = self.oh_chunk_spans(a) {
to_free.extend(spans.into_iter().map(|(a, l)| (a, l, PageType::Meta)));
}
}
}
}
}
retain_disjoint_in_bounds(&mut to_free, self.image.len());
let mut path_addr: BTreeMap<PathKey, u64> = BTreeMap::new();
let mut by_depth = keys.clone();
by_depth.sort_by_key(|k| std::cmp::Reverse(k.len())); for key in &by_depth {
let (mut region, deletes, copies, writes, pending_vl_attrs) = {
let node = nodes.get_mut(key).unwrap();
(
std::mem::take(&mut node.base_region),
std::mem::take(&mut node.deletes),
std::mem::take(&mut node.copies),
std::mem::take(&mut node.writes),
std::mem::take(&mut node.pending_vl_attrs),
)
};
for name in &deletes {
region = remove_link_from_region(®ion, name)?;
}
for (leaf, tree) in copies {
let root = self.write_copy_subtree(&tree)?;
region.extend_from_slice(&encode_link_message(&leaf, root - base));
}
let mut group_datasets: Vec<FlatDataset> =
flat.remove(key).into_iter().flatten().collect();
group_datasets.sort_by_key(|fd| fd.reference_targets.is_some());
for mut fd in group_datasets {
for (idx, collections) in std::mem::take(&mut fd.vl_attrs) {
let addrs = self.place_vl_collections(&collections)?;
patch_vl_refs(&mut fd.attrs[idx].raw_data, &addrs);
}
if let Some(patches) = fd.reference_targets.take() {
for patch in &patches {
let addr = Self::resolve_reference_target(
&patch.target,
&path_addr,
&nodes,
&add_targets,
&write_targets,
&delete_targets,
&self.image(),
&self.superblock,
)?;
write_reference_address(&mut fd.raw, patch.byte_offset, addr);
}
}
let oh = if fd.chunk_options.is_chunked() || fd.maxshape.is_some() {
self.build_chunked_dataset(&fd)?
} else {
if let Some(staging) = fd.vl_string_staging.take() {
if !staging.collections.is_empty() {
let addrs = self.place_vl_collections(&staging.collections)?;
patch_vl_refs_masked(&mut fd.raw, &staging.patch_offsets, &addrs);
}
}
let data_addr = if fd.raw.is_empty() {
u64::MAX
} else {
self.alloc_or_append_typed(&fd.raw, PageType::Raw)? - base
};
build_dataset_oh(
&fd.dt,
&fd.ds,
data_addr,
fd.raw.len() as u64,
&fd.attrs,
None,
fd.fill.as_deref(),
)?
};
let oh_addr = self.alloc_or_append_typed(&oh, PageType::Meta)?;
region.extend_from_slice(&encode_link_message(&fd.name, oh_addr - base));
let mut full = key.clone();
full.push(fd.name.clone());
path_addr.insert(full, oh_addr);
}
for (leaf, mw) in &writes {
let new_oh = self.write_moving(mw)?;
patch_link_target(&mut region, leaf, new_oh - base)?;
}
for child in children.get(key).into_iter().flatten() {
let child_name = child.last().unwrap();
let child_addr = path_addr[child] - base;
if nodes[child].is_new {
region.extend_from_slice(&encode_link_message(child_name, child_addr));
} else {
patch_link_target(&mut region, child_name, child_addr)?;
}
}
for (mut msg, collections) in pending_vl_attrs {
let addrs = self.place_vl_collections(&collections)?;
patch_vl_refs(&mut msg.raw_data, &addrs);
region.extend_from_slice(®ion_message(
MessageType::Attribute,
&msg.serialize(LENGTH_SIZE),
));
}
let oh = build_v2_object_header(®ion);
let addr = self.alloc_or_append_typed(&oh, PageType::Meta)?;
path_addr.insert(key.clone(), addr);
}
for (data_addr, raw) in &inplace_writes {
self.write_at(*data_addr, raw)?;
}
let new_root = path_addr[&PathKey::new()];
if self.persist.is_some() {
return self.commit_persisting(new_root, to_free);
}
for (a, l, _) in to_free.drain(..) {
self.free.free(a, l);
}
let cur_eof = self.image.len();
let trunc_to = self.free.take_trailing(cur_eof);
let new_eof = trunc_to.unwrap_or(cur_eof);
self.image.sync_all()?;
if self.superblock.version >= 2 {
let mut new_sb = self.superblock.clone();
new_sb.root_group_address = new_root - base;
new_sb.eof_address = new_eof;
new_sb.consistency_flags = 0;
let sb_bytes = new_sb.serialize();
self.write_at(self.sb_sig_off, &sb_bytes)?;
self.image.sync_all()?;
new_sb.root_group_address = new_root;
self.superblock = new_sb;
} else {
self.repoint_v0v1_root(new_root - base, new_eof)?;
self.image.sync_all()?;
self.superblock.root_group_address = new_root;
self.superblock.eof_address = new_eof;
}
if let Some(cut) = trunc_to {
self.image.truncate(cut)?;
self.image.sync_all()?;
}
Ok(())
}
fn commit_persisting(
&mut self,
new_root: u64,
to_free: Vec<(u64, u64, PageType)>,
) -> Result<(), Error> {
if self.paged.is_some() {
return self.commit_persisting_paged(new_root, to_free);
}
let os = self.superblock.offset_size;
let (strategy, threshold, page_size, old_blocks) = {
let ps = self
.persist
.as_ref()
.expect("commit_persisting is only called when persistence is armed");
(
ps.strategy,
ps.threshold,
ps.page_size,
ps.old_blocks.clone(),
)
};
let mut post = self.free.clone();
for &(a, l, _) in &to_free {
post.free(a, l);
}
for &(a, l) in &old_blocks {
post.free(a, l);
}
let sections: Vec<FreeSection> = post
.sections()
.into_iter()
.map(|(addr, size)| FreeSection { addr, size })
.collect();
let old_ext_rel = self
.superblock
.superblock_extension_address
.filter(|&a| a != UNDEF)
.ok_or(Error::EditUnsupported(
"a persisting file has no superblock extension to update",
))?;
let old_ext_addr = usize::try_from(old_ext_rel)
.map_err(|_| Error::EditUnsupported("extension address exceeds this platform"))?;
let placeholder =
FileSpaceInfo::persistent_single_manager(strategy, threshold, page_size, 0, 0);
let ext_len =
build_v2_object_header(&self.rewrite_extension_region(old_ext_addr, &placeholder)?)
.len() as u64;
let ext_addr = self.image.len();
let fshd_addr = ext_addr + ext_len;
let (ext_oh, fsm_blocks, final_eof) = if sections.is_empty() {
let info = FileSpaceInfo::persistent_empty(strategy, threshold, page_size);
let ext_oh =
build_v2_object_header(&self.rewrite_extension_region(old_ext_addr, &info)?);
let final_eof = ext_addr + ext_oh.len() as u64;
(ext_oh, None, final_eof)
} else {
let fsse_addr = fshd_addr + fshd_len(os);
let eoa_pre_fsm = fshd_addr;
let info = FileSpaceInfo::persistent_single_manager(
strategy,
threshold,
page_size,
fshd_addr,
eoa_pre_fsm,
);
let ext_oh =
build_v2_object_header(&self.rewrite_extension_region(old_ext_addr, &info)?);
debug_assert_eq!(
ext_oh.len() as u64,
ext_len,
"extension length must be stable across the placeholder and real messages"
);
let (fshd, fsse) =
serialize_file_fsm(§ions, fshd_addr, fsse_addr, os, SECT_CLASS_SIMPLE);
let final_eof = fsse_addr + fsse.len() as u64;
(ext_oh, Some((fshd, fsse)), final_eof)
};
let written_ext = self.append(&ext_oh)?;
debug_assert_eq!(written_ext, ext_addr);
let mut new_old_blocks = vec![(ext_addr, ext_oh.len() as u64)];
if let Some((fshd, fsse)) = fsm_blocks {
let wf = self.append(&fshd)?;
debug_assert_eq!(wf, fshd_addr);
new_old_blocks.push((fshd_addr, fshd.len() as u64));
let ws = self.append(&fsse)?;
new_old_blocks.push((ws, fsse.len() as u64));
}
self.image.sync_all()?;
let mut new_sb = self.superblock.clone();
new_sb.root_group_address = new_root;
new_sb.eof_address = final_eof;
new_sb.superblock_extension_address = Some(ext_addr);
new_sb.consistency_flags = 0;
let sb_bytes = new_sb.serialize();
self.write_at(self.sb_sig_off, &sb_bytes)?;
self.image.sync_all()?;
self.superblock = new_sb;
self.free = post;
self.persist = Some(PersistState {
strategy,
threshold,
page_size,
old_blocks: new_old_blocks,
});
self.fsm_len = self.image.len();
Ok(())
}
fn commit_persisting_paged(
&mut self,
new_root: u64,
to_free: Vec<(u64, u64, PageType)>,
) -> Result<(), Error> {
let os = self.superblock.offset_size;
let (strategy, threshold, page_size, old_blocks) = {
let ps = self
.persist
.as_ref()
.expect("commit_persisting is only called when persistence is armed");
(
ps.strategy,
ps.threshold,
ps.page_size,
ps.old_blocks.clone(),
)
};
self.pad_to_page()?;
let (post_meta, post_raw_small, post_raw_large) = {
let pg = self
.paged
.as_ref()
.expect("commit_persisting_paged is only called on a paged file");
let (mut meta, mut raw_small, mut raw_large) =
(pg.meta.clone(), pg.raw_small.clone(), pg.raw_large.clone());
for &(a, l) in &pg.meta_pad {
meta.free(a, l);
}
for &(a, l) in &pg.raw_pad {
raw_small.free(a, l);
}
for &(a, l, ty) in &to_free {
PagedEdit::route_free(
&mut meta,
&mut raw_small,
&mut raw_large,
page_size,
a,
l,
ty,
);
}
for &(a, l) in &old_blocks {
meta.free(a, l);
}
(meta, raw_small, raw_large)
};
let old_ext_rel = self
.superblock
.superblock_extension_address
.filter(|&a| a != UNDEF)
.ok_or(Error::EditUnsupported(
"a persisting file has no superblock extension to update",
))?;
let old_ext_addr = usize::try_from(old_ext_rel)
.map_err(|_| Error::EditUnsupported("extension address exceeds this platform"))?;
let placeholder = FileSpaceInfo::persistent_managers(
strategy,
threshold,
page_size,
[UNDEF; NUM_FILE_FSM_MANAGERS],
0,
);
let ext_len =
build_v2_object_header(&self.rewrite_extension_region(old_ext_addr, &placeholder)?)
.len() as u64;
let ext_addr = self.image.len();
debug_assert_eq!(
ext_addr % page_size,
0,
"the extension begins on a page boundary"
);
let plan = plan_paged_managers(
&free_sections(&post_meta),
&free_sections(&post_raw_small),
&free_sections(&post_raw_large),
page_size,
ext_addr + ext_len,
os,
);
let (ext_oh, final_eof) = if plan.is_empty() {
let info = FileSpaceInfo::persistent_empty(strategy, threshold, page_size);
let ext_oh =
build_v2_object_header(&self.rewrite_extension_region(old_ext_addr, &info)?);
let final_eof = align_up(ext_addr + ext_oh.len() as u64, page_size);
(ext_oh, final_eof)
} else {
let final_eof = align_up(plan.end_of_managers, page_size);
let info = FileSpaceInfo::persistent_managers(
strategy, threshold, page_size, plan.slots, final_eof,
);
let ext_oh =
build_v2_object_header(&self.rewrite_extension_region(old_ext_addr, &info)?);
debug_assert_eq!(
ext_oh.len() as u64,
ext_len,
"extension length must be stable across the placeholder and real messages"
);
(ext_oh, final_eof)
};
let written_ext = self.append(&ext_oh)?;
debug_assert_eq!(written_ext, ext_addr);
let mut new_old_blocks = vec![(ext_addr, ext_oh.len() as u64)];
for b in &plan.blocks {
let (fshd, fsse) =
serialize_file_fsm(&b.sections, b.fshd_addr, b.fsse_addr, os, b.class);
let wf = self.append(&fshd)?;
debug_assert_eq!(wf, b.fshd_addr);
new_old_blocks.push((b.fshd_addr, fshd.len() as u64));
let ws = self.append(&fsse)?;
debug_assert_eq!(ws, b.fsse_addr);
new_old_blocks.push((ws, fsse.len() as u64));
}
self.pad_zeros_to(final_eof)?;
self.image.sync_all()?;
let mut new_sb = self.superblock.clone();
new_sb.root_group_address = new_root;
new_sb.eof_address = final_eof;
new_sb.superblock_extension_address = Some(ext_addr);
new_sb.consistency_flags = 0;
let sb_bytes = new_sb.serialize();
self.write_at(self.sb_sig_off, &sb_bytes)?;
self.image.sync_all()?;
self.superblock = new_sb;
if let Some(pg) = self.paged.as_mut() {
pg.meta = post_meta;
pg.raw_small = post_raw_small;
pg.raw_large = post_raw_large;
pg.meta_pad.clear();
pg.raw_pad.clear();
pg.last = Some(PageType::Meta);
}
self.persist = Some(PersistState {
strategy,
threshold,
page_size,
old_blocks: new_old_blocks,
});
self.fsm_len = self.image.len();
Ok(())
}
fn pad_to_page(&mut self) -> Result<(), Error> {
let len = self.image.len();
let pad = match &self.paged {
Some(pg) if len % pg.page_size != 0 => {
Some((pg.last, pg.page_size - len % pg.page_size))
}
_ => None,
};
if let Some((last, pad_len)) = pad {
let pad_at = len;
self.append(&vec![0u8; pad_len.to_usize()?])?;
if let Some(pg) = self.paged.as_mut() {
match last {
Some(PageType::Meta) => pg.meta_pad.push((pad_at, pad_len)),
_ => pg.raw_pad.push((pad_at, pad_len)),
}
}
}
Ok(())
}
fn pad_zeros_to(&mut self, target: u64) -> Result<(), Error> {
let len = self.image.len();
if target > len {
let pad = (target - len).to_usize()?;
self.append(&vec![0u8; pad])?;
}
debug_assert_eq!(self.image.len(), target);
Ok(())
}
fn rewrite_extension_region(
&self,
ext_addr: usize,
info: &FileSpaceInfo,
) -> Result<Vec<u8>, Error> {
let region =
Self::gather_oh_messages(&self.image(), ext_addr as u64, self.superblock.base_address)?;
rewrite_extension_region_bytes(®ion, info)
}
fn repoint_v0v1_root(&mut self, new_root: u64, new_eof: u64) -> Result<(), Error> {
let os = self.superblock.offset_size as usize;
let var_start = if self.superblock.version == 0 { 24 } else { 28 };
let base = self.sb_sig_off + var_start;
let eof_off = base + 2 * os;
let ste = base + 4 * os;
let oh_addr_off = ste + os;
let cache_off = ste + 2 * os;
self.write_at(eof_off, &new_eof.to_le_bytes()[..os])?;
self.write_at(cache_off, &[0u8; 4])?; self.write_at(cache_off + 8, &[0u8; 16])?; self.write_at(oh_addr_off, &new_root.to_le_bytes()[..os])?;
Ok(())
}
fn gather_oh_messages<S: Source + ?Sized>(
src: &S,
addr: u64,
base: u64,
) -> Result<Vec<u8>, Error> {
let mut out = Vec::new();
for chunk in read_oh_chunks(src, addr, base)? {
let (region, mut p) = chunk.message_region();
while let Some((msg_type, _body, body_end)) = next_message(region, p)? {
if msg_type != MessageType::ObjectHeaderContinuation {
out.extend_from_slice(®ion[p..body_end]);
}
p = body_end;
}
}
Ok(out)
}
fn reconstruct_v1_group(&self, addr: usize) -> Result<GroupInfo, Error> {
let os = self.superblock.offset_size;
let ls = self.superblock.length_size;
let base = self.superblock.base_address;
let oh = ObjectHeader::parse_from_source(&self.image(), addr as u64, os, ls, base)?;
if oh
.messages
.iter()
.any(|m| m.msg_type == MessageType::DataLayout)
{
return Err(Error::EditUnsupported(
"a target path names a dataset, not a group",
));
}
let entries = resolve_group_entries_from_source(&self.image(), &oh, os, ls, base)?;
let mut region = fresh_group_region();
let mut link_names = Vec::with_capacity(entries.len());
for e in &entries {
region.extend_from_slice(&encode_link_message(&e.name, e.object_header_address));
link_names.push(e.name.clone());
}
for m in &oh.messages {
if m.msg_type == MessageType::Attribute {
if m.flags != 0 {
return Err(Error::EditUnsupported(
"a v0/v1 group has a shared attribute message (not convertible in place yet)",
));
}
if m.data.len() > OBJECT_HEADER_MESSAGE_MAX {
return Err(Error::EditUnsupported(
"a v0/v1 group attribute is too large to convert in place",
));
}
#[expect(
clippy::cast_possible_truncation,
reason = "message type ids are a small enum that fits the 1-byte v2 type field"
)]
region.push(MessageType::Attribute.to_u16() as u8);
#[expect(
clippy::cast_possible_truncation,
reason = "attribute body length fits the 2-byte message-size field (oversized \
bodies are rejected above)"
)]
region.extend_from_slice(&(m.data.len() as u16).to_le_bytes());
region.push(0); region.extend_from_slice(&m.data);
}
}
Ok(GroupInfo { region, link_names })
}
fn inspect_group(&self, addr: usize) -> Result<GroupInfo, Error> {
let sig = self.image().read_metadata_at(addr as u64, 4);
if sig.as_deref() != Ok(&b"OHDR"[..]) {
return self.reconstruct_v1_group(addr);
}
let mut region =
Self::gather_oh_messages(&self.image(), addr as u64, self.superblock.base_address)?;
let mut p = 0;
let mut has_link_info = false;
let mut link_names = Vec::new();
while let Some((msg_type, body, body_end)) = next_message(®ion, p)? {
match msg_type {
MessageType::LinkInfo => {
has_link_info = true;
let mut q = body + 2;
if body_end - body >= 2 && region[body + 1] & 0x01 != 0 {
q += 8;
}
if q + 8 <= body_end {
let heap_addr = u64::from_le_bytes(region[q..q + 8].try_into().unwrap());
if heap_addr != u64::MAX {
return Err(Error::EditUnsupported(
"a target group uses dense (fractal-heap) link storage (not supported in place yet)",
));
}
}
}
MessageType::Link => {
if let Ok(link) = LinkMessage::parse(®ion[body..body_end], OFFSET_SIZE) {
link_names.push(link.name);
}
}
MessageType::DataLayout => {
return Err(Error::EditUnsupported(
"a target path names a dataset, not a group",
));
}
_ => {}
}
p = body_end;
}
if !has_link_info {
return Err(Error::EditUnsupported(
"a target group's object header has no link-info message",
));
}
ensure_group_info(&mut region)?;
Ok(GroupInfo { region, link_names })
}
fn prepare_write<S: Source + ?Sized>(
src: &S,
addr: u64,
fd: &FlatDataset,
base: u64,
) -> Result<WritePlan, Error> {
if fd.chunk_options.is_chunked() || fd.maxshape.is_some() {
return Err(Error::EditUnsupported(
"write_dataset overwrites values only; it cannot make a dataset \
chunked, filtered, or extensible",
));
}
if !fd.attrs.is_empty() {
return Err(Error::EditUnsupported(
"write_dataset overwrites values only; it cannot set attributes \
(set them with a separate edit)",
));
}
if fd.fill.is_some() {
return Err(Error::EditUnsupported(
"write_dataset overwrites values only; it cannot change the fill \
value (set it when the dataset is created)",
));
}
if fd.vl_string_staging.is_some() {
return Err(Error::EditUnsupported(
"write_dataset cannot overwrite a variable-length-string dataset's \
data in place yet",
));
}
let region = Self::gather_oh_messages(src, addr, base)?;
let mut datatype: Option<(usize, usize)> = None;
let mut dataspace: Option<(usize, usize)> = None;
let mut layout: Option<(usize, usize)> = None;
let mut filter: Option<(usize, usize)> = None;
let mut has_link = false;
let mut p = 0;
while let Some((msg_type, body, body_end)) = next_message(®ion, p)? {
match msg_type {
MessageType::Datatype => datatype = Some((body, body_end)),
MessageType::Dataspace => dataspace = Some((body, body_end)),
MessageType::DataLayout => layout = Some((body, body_end)),
MessageType::FilterPipeline => filter = Some((body, body_end)),
MessageType::Link | MessageType::LinkInfo | MessageType::SymbolTable => {
has_link = true;
}
_ => {}
}
p = body_end;
}
if has_link {
return Err(Error::EditUnsupported(
"write_dataset target is a group, not a dataset",
));
}
let (dt_b, dt_e) =
datatype.ok_or(Error::EditUnsupported("dataset header has no datatype"))?;
let (ds_b, ds_e) =
dataspace.ok_or(Error::EditUnsupported("dataset header has no dataspace"))?;
let (lb, le) = layout.ok_or(Error::EditUnsupported("dataset header has no data layout"))?;
let (disk_dt, _) = crate::datatype::Datatype::parse(®ion[dt_b..dt_e])
.map_err(|_| Error::EditUnsupported("dataset header datatype could not be parsed"))?;
if disk_dt != fd.dt {
return Err(Error::EditUnsupported(
"write_dataset datatype does not match the on-disk dataset (overwrite, not retype)",
));
}
let disk_ds = Dataspace::parse(®ion[ds_b..ds_e], LENGTH_SIZE)
.map_err(|_| Error::EditUnsupported("dataset header dataspace could not be parsed"))?;
if disk_ds.space_type != fd.ds.space_type
|| disk_ds.rank != fd.ds.rank
|| disk_ds.dimensions != fd.ds.dimensions
{
return Err(Error::EditUnsupported(
"write_dataset shape does not match the on-disk dataset (overwrite, not reshape)",
));
}
if le - lb < 2 {
return Err(Error::EditUnsupported("malformed data-layout message"));
}
let version = region[lb];
if version != 3 && version != 4 {
return Err(Error::EditUnsupported(
"an unsupported data-layout version cannot be overwritten in place yet",
));
}
match region[lb + 1] {
0 => Ok(WritePlan::Moving(MovingWrite::Compact {
region,
raw: fd.raw.clone(),
})),
1 => {
if le - lb < 18 {
return Err(Error::EditUnsupported("malformed contiguous data layout"));
}
let addr_off = lb + 2;
let data_addr =
u64::from_le_bytes(region[addr_off..addr_off + 8].try_into().unwrap());
let data_size = u64::from_le_bytes(region[lb + 10..lb + 18].try_into().unwrap());
if data_addr != UNDEF && data_size == fd.raw.len() as u64 {
if let Some(start) = data_addr
.checked_add(base)
.and_then(|a| usize::try_from(a).ok())
{
if start
.checked_add(fd.raw.len())
.is_some_and(|e| e as u64 <= src.len())
{
return Ok(WritePlan::InPlace {
data_addr: start,
raw: fd.raw.clone(),
});
}
}
}
let old_extent = if data_addr != UNDEF && data_size > 0 {
Some((data_addr + base, data_size))
} else {
None
};
Ok(WritePlan::Moving(MovingWrite::Contiguous {
region,
addr_off,
raw: fd.raw.clone(),
old_extent,
}))
}
2 => {
let dl =
DataLayout::parse(®ion[lb..le], OFFSET_SIZE, LENGTH_SIZE).map_err(|_| {
Error::EditUnsupported("dataset header data layout could not be parsed")
})?;
let DataLayout::Chunked {
version: lversion,
chunk_index_type,
..
} = dl
else {
return Err(Error::EditUnsupported("dataset is not chunked"));
};
if !chunk_index_enumerable(lversion, chunk_index_type) {
return Err(Error::EditUnsupported(
"a chunked dataset with a version-2 B-tree or unknown chunk index \
cannot be overwritten in place yet",
));
}
let ChunkedGeometry {
spatial,
element_size,
raw_size,
maxshape,
} = chunked_geometry(&fd.dt, &disk_ds, &dl)?;
let split = split_into_chunks(&fd.raw, &disk_ds.dimensions, &spatial, element_size);
let pipeline_message: Option<Vec<u8>> =
filter.map(|(fb, fe)| region[fb..fe].to_vec());
let new_chunk_bytes: Vec<Vec<u8>> = if let Some(pm) = &pipeline_message {
let pipeline = FilterPipeline::parse(pm).map_err(|_| {
Error::EditUnsupported("dataset filter pipeline could not be parsed")
})?;
if !pipeline_reencodable(&pipeline) {
return Err(Error::EditUnsupported(
"a chunked dataset using a filter this engine cannot re-encode \
cannot be overwritten in place yet",
));
}
let ctx = ChunkContext::from_datatype(&spatial, &fd.dt);
let mut encoded = Vec::with_capacity(split.len());
for (_, buf) in &split {
encoded.push(compress_chunk(buf, &pipeline, ctx)?);
}
encoded
} else {
split.into_iter().map(|(_, buf)| buf).collect()
};
let base_off = usize::try_from(base).map_err(|_| {
Error::EditUnsupported("userblock base address exceeds this platform")
})?;
if let Some(writes) = try_inplace_chunk_writes(
&BaseOffsetSource { inner: src, base },
&dl,
&disk_ds,
&spatial,
raw_size,
&new_chunk_bytes,
) {
let writes = writes
.into_iter()
.map(|(off, b)| (off + base_off, b))
.collect();
return Ok(WritePlan::InPlaceChunks { writes });
}
let meta = new_chunk_bytes
.iter()
.map(|c| ChunkMeta {
compressed_size: c.len() as u64,
filter_mask: 0,
})
.collect();
Ok(WritePlan::Moving(MovingWrite::Chunked {
region,
chunk_dims: spatial,
element_size,
raw_size,
maxshape,
pipeline_message,
meta,
chunk_bytes: new_chunk_bytes,
old_addr: addr,
}))
}
_ => Err(Error::EditUnsupported(
"an unsupported data-layout class cannot be overwritten in place yet",
)),
}
}
fn prepare_append<S: Source + ?Sized>(
src: &S,
addr: u64,
ab: &AppendBuilder,
base: u64,
) -> Result<MovingWrite, Error> {
if ab.dt_conflict {
return Err(Error::AppendUnsupported(
"append mixes element types in one builder; use one element type per \
append_dataset call",
));
}
let region = Self::gather_oh_messages(src, addr, base)?;
let mut datatype: Option<(usize, usize)> = None;
let mut dataspace: Option<(usize, usize)> = None;
let mut layout: Option<(usize, usize)> = None;
let mut filter: Option<(usize, usize)> = None;
let mut has_link = false;
let mut p = 0;
while let Some((msg_type, body, body_end)) = next_message(®ion, p)? {
match msg_type {
MessageType::Datatype => datatype = Some((body, body_end)),
MessageType::Dataspace => dataspace = Some((body, body_end)),
MessageType::DataLayout => layout = Some((body, body_end)),
MessageType::FilterPipeline => filter = Some((body, body_end)),
MessageType::Link | MessageType::LinkInfo | MessageType::SymbolTable => {
has_link = true;
}
_ => {}
}
p = body_end;
}
if has_link {
return Err(Error::AppendUnsupported(
"append target is a group, not a dataset",
));
}
let (dt_b, dt_e) =
datatype.ok_or(Error::AppendUnsupported("dataset header has no datatype"))?;
let (ds_b, ds_e) =
dataspace.ok_or(Error::AppendUnsupported("dataset header has no dataspace"))?;
let (lb, le) = layout.ok_or(Error::AppendUnsupported(
"dataset header has no data layout",
))?;
let (disk_dt, _) = Datatype::parse(®ion[dt_b..dt_e])
.map_err(|_| Error::AppendUnsupported("dataset header datatype could not be parsed"))?;
let disk_ds = Dataspace::parse(®ion[ds_b..ds_e], LENGTH_SIZE).map_err(|_| {
Error::AppendUnsupported("dataset header dataspace could not be parsed")
})?;
let dl = DataLayout::parse(®ion[lb..le], OFFSET_SIZE, LENGTH_SIZE).map_err(|_| {
Error::AppendUnsupported("dataset header data layout could not be parsed")
})?;
let DataLayout::Chunked {
version: lversion,
chunk_index_type,
btree_address,
..
} = &dl
else {
return Err(Error::AppendUnsupported(
"append requires a chunked dataset",
));
};
if *lversion != 4 || *chunk_index_type != Some(4) {
return Err(Error::AppendUnsupported(
"append requires an Extensible-Array-indexed chunked dataset (a single \
unlimited dimension under the latest format)",
));
}
if disk_ds.space_type != DataspaceType::Simple || disk_ds.dimensions.len() != 1 {
return Err(Error::AppendUnsupported(
"append requires a rank-1 dataset in this release",
));
}
match &disk_ds.max_dimensions {
Some(md) if md.first() == Some(&u64::MAX) => {}
_ => {
return Err(Error::AppendUnsupported(
"append requires a dataset that is unlimited along its first dimension",
));
}
}
let ChunkedGeometry {
spatial,
element_size,
raw_size,
..
} = chunked_geometry(&disk_dt, &disk_ds, &dl)?;
let chunk_elems = spatial[0];
if chunk_elems == 0 {
return Err(Error::AppendUnsupported(
"append requires a nonzero chunk length",
));
}
if ab.raw.len() % element_size != 0 {
return Err(Error::AppendUnsupported(
"appended byte length is not a whole number of elements",
));
}
match &ab.elem_dt {
Some(expected) if *expected != disk_dt => {
return Err(Error::AppendUnsupported(
"append datatype does not match the on-disk dataset (wrong element \
type or byte order)",
));
}
Some(_) => {}
None => {
if !datatype_is_raw_appendable(&disk_dt) {
return Err(Error::AppendUnsupported(
"append_raw onto this dataset's datatype (non-little-endian, \
variable-length, or reference) could misencode the bytes; use a \
typed append",
));
}
}
}
let new_elems = (ab.raw.len() / element_size) as u64;
let current_dim0 = disk_ds.dimensions[0];
let new_dim0 = current_dim0
.checked_add(new_elems)
.ok_or(Error::AppendUnsupported(
"append would overflow the dataset dimension",
))?;
let pipeline_message: Option<Vec<u8>> = filter.map(|(fb, fe)| region[fb..fe].to_vec());
let has_filters = pipeline_message.is_some();
let pipeline = match &pipeline_message {
Some(pm) => {
let parsed = FilterPipeline::parse(pm).map_err(|_| {
Error::AppendUnsupported("dataset filter pipeline could not be parsed")
})?;
if !pipeline_reencodable(&parsed) {
return Err(Error::AppendUnsupported(
"dataset uses a filter this engine cannot re-encode",
));
}
Some(parsed)
}
None => None,
};
if base > src.len() {
return Err(Error::AppendUnsupported(
"userblock base address past end-of-file",
));
}
let view = BaseOffsetSource { inner: src, base };
if let Some(idx_addr) = *btree_address {
let hdr =
ExtensibleArrayHeader::parse_from_source(&view, idx_addr, OFFSET_SIZE, LENGTH_SIZE)
.map_err(|_| {
Error::AppendUnsupported(
"dataset extensible-array header could not be parsed",
)
})?;
if (hdr.client_id == 1) != has_filters {
return Err(Error::AppendUnsupported(
"dataset filter metadata is inconsistent (chunk-index client id \
disagrees with the filter pipeline)",
));
}
}
let infos = enumerate_chunks_from_source(&view, &dl, &disk_ds, OFFSET_SIZE, LENGTH_SIZE)
.map_err(|_| Error::AppendUnsupported("dataset chunk index could not be enumerated"))?;
let grid = plan_dense_grid(infos, &disk_ds.dimensions, &spatial).ok_or(
Error::AppendUnsupported(
"dataset has a sparse or inconsistent chunk grid; cannot append",
),
)?;
let grid_order = grid.grid_order;
let n_full = usize::try_from(current_dim0 / chunk_elems)
.map_err(|_| Error::AppendUnsupported("chunk count exceeds this platform"))?;
let has_partial = current_dim0 % chunk_elems != 0;
let mut kept_chunks: Vec<WrittenChunk> = Vec::with_capacity(n_full);
for ci in grid_order.iter().take(n_full) {
kept_chunks.push(WrittenChunk {
address: ci.address,
compressed_size: u64::from(ci.chunk_size),
raw_size,
filter_mask: ci.filter_mask,
});
}
let mut tail_raw: Vec<u8> = Vec::new();
let mut old_tail_extent: Option<(u64, u64)> = None;
if has_partial {
let partial = &grid_order[n_full];
let len = partial.chunk_size as usize;
partial
.address
.checked_add(len as u64)
.filter(|&e| e <= view.len())
.ok_or(Error::AppendUnsupported(
"trailing chunk extends past end-of-file",
))?;
let stored = view
.read_exact_at(partial.address, len)
.map_err(|_| Error::AppendUnsupported("trailing chunk could not be read"))?;
let full = if let Some(pl) = &pipeline {
let ctx = ChunkContext::from_datatype(&spatial, &disk_dt);
decompress_chunk(&stored, pl, ctx, partial.filter_mask).map_err(Error::Format)?
} else {
stored
};
let live_elems = usize::try_from(current_dim0 % 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]);
old_tail_extent = Some((partial.address + base, u64::from(partial.chunk_size)));
}
tail_raw.extend_from_slice(&ab.raw);
let tail_len_elems = new_dim0 - (n_full as u64) * 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, &disk_dt);
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 mut grown = disk_ds.clone();
grown.dimensions[0] = new_dim0;
let new_dataspace_body = grown.serialize(LENGTH_SIZE);
#[expect(
clippy::cast_possible_truncation,
reason = "spatial chunk dims come from the on-disk u32 chunk_dimensions, so they fit u32"
)]
let chunk_dims_u32: Vec<u32> = spatial.iter().map(|&dm| dm as u32).collect();
Ok(MovingWrite::AppendedChunks {
region,
new_dataspace_body,
chunk_dims_u32,
element_size,
raw_size,
has_filters,
kept_chunks,
new_chunk_bytes,
old_addr: addr,
old_tail_extent,
})
}
fn read_object<S: Source + ?Sized>(src: &S, addr: u64, base: u64) -> Result<ObjModel, Error> {
let region = Self::gather_oh_messages(src, addr, base)?;
let mut dense = false;
let mut p = 0;
while let Some((msg_type, body, body_end)) = next_message(®ion, p)? {
if msg_type == MessageType::AttributeInfo {
let ai = crate::attribute_info::AttributeInfoMessage::parse(
®ion[body..body_end],
OFFSET_SIZE,
)
.map_err(|_| {
Error::EditUnsupported(
"a source attribute-info message could not be parsed for copying",
)
})?;
if ai.fractal_heap_address.is_some() {
dense = true;
}
}
p = body_end;
}
let dense_attrs = if dense {
let header =
ObjectHeader::parse_from_source(src, addr, OFFSET_SIZE, LENGTH_SIZE, base).map_err(|_| {
Error::EditUnsupported(
"a source object header with dense attributes could not be parsed for copying",
)
})?;
if base > src.len() {
return Err(Error::EditUnsupported(
"a source file's userblock is larger than the file itself",
));
}
let framed = BaseOffsetSource { inner: src, base };
let attrs = crate::attribute::extract_attributes_full_from_source(
&framed,
&header,
OFFSET_SIZE,
LENGTH_SIZE,
)
.map_err(|_| {
Error::EditUnsupported(
"a source object's dense (fractal-heap) attributes could not be read for copying",
)
})?;
crate::file_writer::dense_attrs_check(&attrs).map_err(Error::Format)?;
attrs
} else {
Vec::new()
};
let mut layout: Option<(usize, usize)> = None; let mut has_link_info = false;
let mut children: Vec<(String, u64)> = Vec::new();
let mut kept: Vec<u8> = Vec::new();
let mut p = 0;
while let Some((msg_type, body, body_end)) = next_message(®ion, p)? {
let mut keep = true;
match msg_type {
MessageType::AttributeInfo => {
if dense {
keep = false;
}
}
MessageType::Attribute => {
if dense {
keep = false;
}
}
MessageType::LinkInfo => {
has_link_info = true;
let mut q = body + 2;
if body_end - body >= 2 && region[body + 1] & 0x01 != 0 {
q += 8;
}
if q + 8 <= body_end {
let heap_addr = u64::from_le_bytes(region[q..q + 8].try_into().unwrap());
if heap_addr != u64::MAX {
return Err(Error::EditUnsupported(
"a group uses dense (fractal-heap) link storage (not supported in place yet)",
));
}
}
}
MessageType::Link => {
keep = false;
match LinkMessage::parse(®ion[body..body_end], OFFSET_SIZE) {
Ok(LinkMessage {
name,
link_target:
LinkTarget::Hard {
object_header_address,
},
..
}) => children.push((name, object_header_address)),
_ => {
return Err(Error::EditUnsupported(
"a group contains a soft/external link (not copyable in place yet)",
));
}
}
}
MessageType::DataLayout => {
layout = Some((kept.len() + (body - p), body_end - body));
}
_ => {}
}
if keep {
kept.extend_from_slice(®ion[p..body_end]);
}
p = body_end;
}
if let Some((lbody, lsize)) = layout {
let version = kept[lbody];
if !(version == 3 || version == 4) || lsize < 2 {
return Err(Error::EditUnsupported(
"an unsupported data-layout version cannot be copied in place yet",
));
}
let class = kept[lbody + 1];
match class {
0 => Ok(ObjModel::DatasetVerbatim {
region: kept,
dense_attrs,
}),
1 => {
if lbody + 18 > kept.len() {
return Err(Error::EditUnsupported("malformed contiguous data layout"));
}
let data_addr =
u64::from_le_bytes(kept[lbody + 2..lbody + 10].try_into().unwrap());
let data_size =
u64::from_le_bytes(kept[lbody + 10..lbody + 18].try_into().unwrap());
Ok(ObjModel::DatasetContiguous {
region: kept,
addr_off: lbody + 2,
data_addr,
data_size,
dense_attrs,
})
}
2 => Ok(ObjModel::DatasetChunked {
region: kept,
dense_attrs,
}),
_ => Err(Error::EditUnsupported(
"an unsupported data-layout class cannot be copied in place yet",
)),
}
} else if has_link_info {
ensure_group_info(&mut kept)?;
Ok(ObjModel::Group {
non_link_region: kept,
children,
dense_attrs,
})
} else {
Err(Error::EditUnsupported(
"an object is neither a contiguous/compact dataset nor a group",
))
}
}
fn read_copy_subtree<S: Source + ?Sized>(
src: &S,
addr: u64,
depth: u32,
cross_file: bool,
base: u64,
) -> Result<CopyTree, Error> {
if depth >= MAX_COPY_DEPTH {
return Err(Error::EditUnsupported(
"copy source nests too deeply (possible hard-link cycle)",
));
}
match Self::read_object(src, addr, base)? {
ObjModel::DatasetVerbatim {
region,
dense_attrs,
} => {
if cross_file {
reject_foreign_addresses(®ion)?;
reject_foreign_dense_attrs(&dense_attrs)?;
}
Ok(CopyTree::DatasetVerbatim {
region,
dense_attrs,
})
}
ObjModel::DatasetContiguous {
region,
addr_off,
data_addr,
data_size,
dense_attrs,
} => {
if cross_file {
reject_foreign_addresses(®ion)?;
reject_foreign_dense_attrs(&dense_attrs)?;
}
let start = data_addr
.checked_add(base)
.ok_or(Error::EditUnsupported("data address exceeds this platform"))?;
let len = usize::try_from(data_size)
.map_err(|_| Error::EditUnsupported("data size exceeds this platform"))?;
start
.checked_add(len as u64)
.filter(|&e| e <= src.len())
.ok_or(Error::EditUnsupported("dataset data is out of bounds"))?;
Ok(CopyTree::DatasetContiguous {
region,
addr_off,
data: src
.read_exact_at(start, len)
.map_err(|_| Error::EditUnsupported("dataset data is out of bounds"))?,
dense_attrs,
})
}
ObjModel::DatasetChunked {
region,
dense_attrs,
} => {
if cross_file {
reject_foreign_addresses(®ion)?;
reject_foreign_dense_attrs(&dense_attrs)?;
}
let ChunkedHeaderParts {
dt,
ds,
layout,
pipeline_message,
} = parse_chunked_header(®ion)?;
let DataLayout::Chunked {
version: lversion,
chunk_index_type,
..
} = layout
else {
return Err(Error::EditUnsupported("dataset is not chunked"));
};
if !chunk_index_enumerable(lversion, chunk_index_type) {
return Err(Error::EditUnsupported(
"a chunked dataset with a version-2 B-tree or unknown chunk index \
cannot be copied in place yet",
));
}
let ChunkedGeometry {
spatial: chunk_dims,
element_size,
raw_size,
maxshape,
} = chunked_geometry(&dt, &ds, &layout)?;
let dview = BaseOffsetSource { inner: src, base };
let infos =
enumerate_chunks_from_source(&dview, &layout, &ds, OFFSET_SIZE, LENGTH_SIZE)?;
let grid = plan_dense_grid(infos, &ds.dimensions, &chunk_dims).ok_or(
Error::EditUnsupported(
"a chunked dataset with unallocated (sparse) chunks cannot be copied in place yet",
),
)?;
if grid.grid_order.is_empty() {
return Err(Error::EditUnsupported(
"an empty chunked dataset cannot be copied in place yet",
));
}
let mut meta = Vec::with_capacity(grid.grid_order.len());
let mut chunk_bytes = Vec::with_capacity(grid.grid_order.len());
for ci in &grid.grid_order {
let len = ci.chunk_size as usize;
ci.address
.checked_add(len as u64)
.filter(|&e| e <= dview.len())
.ok_or(Error::EditUnsupported("chunk data is out of bounds"))?;
chunk_bytes.push(
dview
.read_exact_at(ci.address, len)
.map_err(|_| Error::EditUnsupported("chunk data is out of bounds"))?,
);
meta.push(ChunkMeta {
compressed_size: ci.chunk_size as u64,
filter_mask: ci.filter_mask,
});
}
Ok(CopyTree::DatasetChunked {
region,
chunk_dims,
element_size,
raw_size,
maxshape,
pipeline_message,
meta,
chunk_bytes,
dense_attrs,
})
}
ObjModel::Group {
non_link_region,
children,
dense_attrs,
} => {
if cross_file {
reject_foreign_addresses(&non_link_region)?;
reject_foreign_dense_attrs(&dense_attrs)?;
}
let mut kids = Vec::with_capacity(children.len());
for (name, child) in children {
let child = child.checked_add(base).ok_or(Error::EditUnsupported(
"child address exceeds this platform",
))?;
kids.push((
name,
Self::read_copy_subtree(src, child, depth + 1, cross_file, base)?,
));
}
Ok(CopyTree::Group {
non_link_region,
children: kids,
dense_attrs,
})
}
}
}
fn write_copy_subtree(&mut self, node: &CopyTree) -> Result<u64, Error> {
let base = self.superblock.base_address;
match node {
CopyTree::DatasetVerbatim {
region,
dense_attrs,
} => {
let mut region = region.clone();
self.append_dense_attrs(&mut region, dense_attrs)?;
let oh = build_v2_object_header(®ion);
self.alloc_or_append_typed(&oh, PageType::Meta)
}
CopyTree::DatasetContiguous {
region,
addr_off,
data,
dense_attrs,
} => {
let new_data_addr = self.alloc_or_append_typed(data, PageType::Raw)?;
let mut region = region.clone();
region[*addr_off..*addr_off + 8]
.copy_from_slice(&(new_data_addr - base).to_le_bytes());
self.append_dense_attrs(&mut region, dense_attrs)?;
let oh = build_v2_object_header(®ion);
self.alloc_or_append_typed(&oh, PageType::Meta)
}
CopyTree::DatasetChunked {
region,
chunk_dims,
element_size,
raw_size,
maxshape,
pipeline_message,
meta,
chunk_bytes,
dense_attrs,
} => self.write_chunked_relocatable(
region,
chunk_dims,
*element_size,
*raw_size,
maxshape.as_deref(),
pipeline_message.as_deref(),
meta,
chunk_bytes,
dense_attrs,
),
CopyTree::Group {
non_link_region,
children,
dense_attrs,
} => {
let mut region = non_link_region.clone();
for (name, child) in children {
let new_child = self.write_copy_subtree(child)?;
region.extend_from_slice(&encode_link_message(name, new_child - base));
}
self.append_dense_attrs(&mut region, dense_attrs)?;
let oh = build_v2_object_header(®ion);
self.alloc_or_append_typed(&oh, PageType::Meta)
}
}
}
#[expect(
clippy::too_many_arguments,
reason = "the chunked rebuild needs the full geometry, \
pipeline, and chunk payloads; bundling them into a struct would only move the list"
)]
fn write_chunked_relocatable(
&mut self,
region: &[u8],
chunk_dims: &[u64],
element_size: usize,
raw_size: u64,
maxshape: Option<&[u64]>,
pipeline_message: Option<&[u8]>,
meta: &[ChunkMeta],
chunk_bytes: &[Vec<u8>],
dense_attrs: &[crate::attribute::AttributeMessage],
) -> Result<u64, Error> {
self.begin_page(PageType::Raw)?;
let eof = self.image.len();
let stored_base = eof - self.superblock.base_address;
let layout = plan_chunked_data_verbatim(
meta,
chunk_dims,
element_size,
raw_size,
pipeline_message,
stored_base,
maxshape,
)?;
let mut buf = Vec::with_capacity(usize::try_from(layout.plan.total_len).unwrap_or(0));
emit_chunked_data_verbatim(
&mut buf,
&layout.plan,
&SliceChunkProvider {
chunks: chunk_bytes,
},
)?;
let written = self.append(&buf)?;
debug_assert_eq!(written, eof, "chunk blob must land at end-of-file",);
let mut new_region = replace_layout_message(region, &layout.layout_message)?;
self.append_dense_attrs(&mut new_region, dense_attrs)?;
let oh = build_v2_object_header(&new_region);
self.alloc_or_append_typed(&oh, PageType::Meta)
}
fn append_dense_attrs(
&mut self,
region: &mut Vec<u8>,
attrs: &[crate::attribute::AttributeMessage],
) -> Result<(), Error> {
if attrs.is_empty() {
return Ok(());
}
self.begin_page(PageType::Meta)?;
let eof = self.image.len();
let stored_base = eof - self.superblock.base_address;
let blob = crate::file_writer::build_dense_attrs(attrs, stored_base);
let written = self.append(&blob.blob)?;
debug_assert_eq!(
written, eof,
"dense attribute blob must land at end-of-file",
);
region.extend_from_slice(®ion_message(
MessageType::AttributeInfo,
&blob.attr_info_message,
));
Ok(())
}
fn write_moving(&mut self, mw: &MovingWrite) -> Result<u64, Error> {
let base = self.superblock.base_address;
match mw {
MovingWrite::Contiguous {
region,
addr_off,
raw,
..
} => {
let new_data_addr = self.alloc_or_append_typed(raw, PageType::Raw)?;
let mut region = region.clone();
region[*addr_off..*addr_off + 8]
.copy_from_slice(&(new_data_addr - base).to_le_bytes());
let size_off = *addr_off + 8;
region[size_off..size_off + 8].copy_from_slice(&(raw.len() as u64).to_le_bytes());
let oh = build_v2_object_header(®ion);
self.alloc_or_append_typed(&oh, PageType::Meta)
}
MovingWrite::Compact { region, raw } => {
let region = rebuild_compact_layout_region(region, raw)?;
let oh = build_v2_object_header(®ion);
self.alloc_or_append_typed(&oh, PageType::Meta)
}
MovingWrite::Chunked {
region,
chunk_dims,
element_size,
raw_size,
maxshape,
pipeline_message,
meta,
chunk_bytes,
..
} => self.write_chunked_relocatable(
region,
chunk_dims,
*element_size,
*raw_size,
maxshape.as_deref(),
pipeline_message.as_deref(),
meta,
chunk_bytes,
&[],
),
MovingWrite::AppendedChunks {
region,
new_dataspace_body,
chunk_dims_u32,
element_size,
raw_size,
has_filters,
kept_chunks,
new_chunk_bytes,
..
} => self.write_appended_chunks(
region,
new_dataspace_body,
chunk_dims_u32,
*element_size,
*raw_size,
*has_filters,
kept_chunks,
new_chunk_bytes,
),
MovingWrite::AttrEdit {
region,
pending_vl_attrs,
} => {
let mut region = region.clone();
for (msg, collections) in pending_vl_attrs {
let mut msg = msg.clone();
let addrs = self.place_vl_collections(collections)?;
patch_vl_refs(&mut msg.raw_data, &addrs);
region.extend_from_slice(®ion_message(
MessageType::Attribute,
&msg.serialize(LENGTH_SIZE),
));
}
let oh = build_v2_object_header(®ion);
self.alloc_or_append_typed(&oh, PageType::Meta)
}
}
}
#[expect(
clippy::too_many_arguments,
reason = "the append rebuild needs the header region, grown dataspace, chunk \
geometry, and both chunk sets; bundling them into a struct would only move the list"
)]
fn write_appended_chunks(
&mut self,
region: &[u8],
new_dataspace_body: &[u8],
chunk_dims_u32: &[u32],
element_size: usize,
raw_size: u64,
has_filters: bool,
kept_chunks: &[WrittenChunk],
new_chunk_bytes: &[Vec<u8>],
) -> Result<u64, Error> {
let base = self.superblock.base_address;
let mut combined: Vec<WrittenChunk> = kept_chunks.to_vec();
if !new_chunk_bytes.is_empty() {
self.begin_page(PageType::Raw)?;
}
for cb in new_chunk_bytes {
let abs = self.append(cb)?;
combined.push(WrittenChunk {
address: abs - base,
compressed_size: cb.len() as u64,
raw_size,
filter_mask: 0,
});
}
self.begin_page(PageType::Raw)?;
let ea_base = self.image.len() - base;
let ea_bytes =
build_extensible_array_at(&combined, OFFSET_SIZE, LENGTH_SIZE, has_filters, ea_base)
.map_err(Error::Format)?;
let written = self.append(&ea_bytes)?;
debug_assert_eq!(
written,
ea_base + base,
"extensible-array index must land at end-of-file",
);
#[expect(
clippy::cast_possible_truncation,
reason = "element size is a datatype byte width that fits u32"
)]
let layout_body = serialize_v4_extensible_array(
chunk_dims_u32,
ea_base,
OFFSET_SIZE,
element_size as u32,
);
let region = replace_dataspace_message(region, new_dataspace_body)?;
let region = replace_layout_message(®ion, &layout_body)?;
let oh = build_v2_object_header(®ion);
self.alloc_or_append_typed(&oh, PageType::Meta)
}
fn append(&mut self, bytes: &[u8]) -> Result<u64, Error> {
self.image.append(bytes)
}
fn write_at(&mut self, offset: usize, bytes: &[u8]) -> Result<(), Error> {
self.image.write_at(offset as u64, bytes)
}
fn begin_page(&mut self, ty: PageType) -> Result<(), Error> {
let Self { image, paged, .. } = self;
match paged.as_mut() {
Some(pg) => pg.begin(image.as_mut(), ty),
None => Ok(()),
}
}
fn append_typed(&mut self, bytes: &[u8], ty: PageType) -> Result<u64, Error> {
self.begin_page(ty)?;
self.append(bytes)
}
fn alloc_or_append_typed(&mut self, bytes: &[u8], ty: PageType) -> Result<u64, Error> {
if self.paged.is_some() {
return self.append_typed(bytes, ty);
}
self.alloc_or_append(bytes)
}
fn alloc_or_append(&mut self, bytes: &[u8]) -> Result<u64, Error> {
debug_assert!(
self.paged.is_none(),
"a paged file must allocate through alloc_or_append_typed"
);
if let Some(addr) = self.free.alloc(bytes.len() as u64) {
self.write_at(
usize::try_from(addr).map_err(|_| {
Error::EditUnsupported("free-region address exceeds this platform")
})?,
bytes,
)?;
Ok(addr)
} else {
self.append(bytes)
}
}
fn place_vl_collections(&mut self, collections: &[Vec<u8>]) -> Result<Vec<u64>, Error> {
collections
.iter()
.map(|collection| {
let addr = self.alloc_or_append_typed(collection, PageType::Meta)?;
Ok(addr - self.superblock.base_address)
})
.collect()
}
fn resolve_reference_target(
target: &ObjectRefTarget,
path_addr: &BTreeMap<PathKey, u64>,
nodes: &BTreeMap<PathKey, Node>,
add_targets: &[PathKey],
write_targets: &[PathKey],
pending_deletes: &[PathKey],
src: &(impl Source + ?Sized),
superblock: &Superblock,
) -> Result<u64, Error> {
let path = match target {
ObjectRefTarget::Raw(addr) => return Ok(*addr),
ObjectRefTarget::Path(path) => path,
};
let base = superblock.base_address;
let key = split_path(path);
if let Some(&addr) = path_addr.get(&key) {
return Ok(addr - base);
}
if nodes.contains_key(&key)
|| add_targets.iter().any(|t| is_prefix(t, &key))
|| write_targets.contains(&key)
|| pending_deletes.contains(&key)
{
return Err(Error::EditUnsupported(
"an object-reference dataset targets a path this commit is still writing; \
use separate commits",
));
}
match crate::group_v2::resolve_path_any_from_source(src, superblock, path) {
Ok(addr) => Ok(addr - base),
Err(_) => Ok(UNDEF),
}
}
fn preflight_reference_targets(
keys: &[PathKey],
flat: &BTreeMap<PathKey, Vec<FlatDataset>>,
nodes: &BTreeMap<PathKey, Node>,
add_targets: &[PathKey],
write_targets: &[PathKey],
pending_deletes: &[PathKey],
src: &(impl Source + ?Sized),
superblock: &Superblock,
) -> Result<(), Error> {
let mut by_depth = keys.to_vec();
by_depth.sort_by_key(|k| std::cmp::Reverse(k.len()));
let mut sim_addr: BTreeMap<PathKey, u64> = BTreeMap::new();
for key in &by_depth {
if let Some(datasets) = flat.get(key) {
let mut ordered: Vec<&FlatDataset> = datasets.iter().collect();
ordered.sort_by_key(|fd| fd.reference_targets.is_some());
for fd in ordered {
if let Some(patches) = &fd.reference_targets {
for patch in patches {
Self::resolve_reference_target(
&patch.target,
&sim_addr,
nodes,
add_targets,
write_targets,
pending_deletes,
src,
superblock,
)?;
}
}
let mut full = key.clone();
full.push(fd.name.clone());
sim_addr.insert(full, 0);
}
}
sim_addr.insert(key.clone(), 0);
}
Ok(())
}
fn build_chunked_dataset(&mut self, fd: &FlatDataset) -> Result<Vec<u8>, Error> {
self.begin_page(PageType::Raw)?;
let eof = self.image.len();
let stored_base = eof - self.superblock.base_address;
let chunk_dims = fd.chunk_options.resolve_chunk_dims(&fd.ds.dimensions);
let ctx = ChunkContext::from_datatype(&chunk_dims, &fd.dt);
let result = build_chunked_data_at_ext(
&fd.raw,
&fd.ds.dimensions,
ctx,
&fd.chunk_options,
stored_base,
fd.maxshape.as_deref(),
)?;
let written = self.append(&result.data_bytes)?;
debug_assert_eq!(written, eof, "chunk blob must land at end-of-file",);
Ok(build_chunked_dataset_oh(
&fd.dt,
&fd.ds,
&result.layout_message,
result.pipeline_message.as_deref(),
&fd.attrs,
None,
fd.fill.as_deref(),
)?)
}
fn oh_chunk_spans(&self, addr: usize) -> Result<Vec<(u64, u64)>, Error> {
Ok(
read_oh_chunks(&self.image(), addr as u64, self.superblock.base_address)?
.into_iter()
.map(|chunk| chunk.span)
.collect(),
)
}
fn count_incoming_hard_links(&self) -> Option<HashMap<u64, u32>> {
let os = self.superblock.offset_size;
let ls = self.superblock.length_size;
let base = self.superblock.base_address;
let mut counts: HashMap<u64, u32> = HashMap::new();
let mut visited: HashSet<u64> = HashSet::new();
let mut stack: Vec<u64> = vec![self.superblock.root_group_address];
let mut budget = MAX_LINK_GRAPH_NODES;
while let Some(addr) = stack.pop() {
if !visited.insert(addr) {
continue; }
if budget == 0 {
return None; }
budget -= 1;
let off = usize::try_from(addr).ok()?;
let header =
ObjectHeader::parse_from_source(&self.image(), off as u64, os, ls, base).ok()?;
let is_group = header.messages.iter().any(|m| {
matches!(
m.msg_type,
MessageType::SymbolTable | MessageType::Link | MessageType::LinkInfo
)
});
if !is_group {
continue;
}
let entries =
resolve_group_entries_from_source(&self.image(), &header, os, ls, base).ok()?;
for e in entries {
let child = e.object_header_address.checked_add(base)?;
*counts.entry(child).or_insert(0) += 1;
stack.push(child);
}
}
Some(counts)
}
fn collect_free_spans(
&self,
addr: usize,
depth: u32,
incoming: &HashMap<u64, u32>,
out: &mut Vec<(u64, u64, PageType)>,
) {
let base = self.superblock.base_address;
let file_len = self.image().len();
if depth >= MAX_COPY_DEPTH {
return;
}
if incoming.get(&(addr as u64)) != Some(&1) {
return;
}
let spans = match self.oh_chunk_spans(addr) {
Ok(s) => s,
Err(_) => return,
};
match Self::read_object(&self.image(), addr as u64, self.superblock.base_address) {
Ok(ObjModel::DatasetVerbatim { .. }) => out.extend(meta_spans(spans)),
Ok(ObjModel::DatasetContiguous {
data_addr,
data_size,
..
}) => {
out.extend(meta_spans(spans));
if data_addr != u64::MAX && data_size > 0 {
if let (Some(abs), Ok(len)) =
(data_addr.checked_add(base), usize::try_from(data_size))
{
if let Ok(start) = usize::try_from(abs) {
if start.checked_add(len).is_some_and(|e| e as u64 <= file_len) {
out.push((abs, data_size, PageType::Raw));
}
}
}
}
}
Ok(ObjModel::Group { children, .. }) => {
out.extend(meta_spans(spans));
for (_, child) in children {
if let Some(c) = child
.checked_add(base)
.and_then(|a| usize::try_from(a).ok())
{
self.collect_free_spans(c, depth + 1, incoming, out);
}
}
}
Ok(ObjModel::DatasetChunked { .. }) => {
if let Some(storage) = self.chunked_storage_spans(addr) {
out.extend(meta_spans(spans));
out.extend(storage);
}
}
Err(_) => {}
}
}
fn chunked_storage_spans(&self, addr: usize) -> Option<Vec<(u64, u64, PageType)>> {
let region =
Self::gather_oh_messages(&self.image(), addr as u64, self.superblock.base_address)
.ok()?;
let mut layout_msg: Option<(usize, usize)> = None;
let mut dataspace_msg: Option<(usize, usize)> = None;
let mut p = 0;
loop {
match next_message(®ion, p) {
Ok(Some((msg_type, body, body_end))) => {
match msg_type {
MessageType::DataLayout => layout_msg = Some((body, body_end)),
MessageType::Dataspace => dataspace_msg = Some((body, body_end)),
_ => {}
}
p = body_end;
}
Ok(None) => break,
Err(_) => return None,
}
}
let (lb, le) = layout_msg?;
let (db, de) = dataspace_msg?;
let layout = DataLayout::parse(®ion[lb..le], OFFSET_SIZE, LENGTH_SIZE).ok()?;
if !matches!(layout, DataLayout::Chunked { .. }) {
return None;
}
let dataspace = Dataspace::parse(®ion[db..de], LENGTH_SIZE).ok()?;
let base = self.superblock.base_address;
let split = crate::chunked_read::collect_chunked_storage_spans(
&BaseOffsetSource {
inner: &self.image(),
base,
},
&layout,
&dataspace,
OFFSET_SIZE,
LENGTH_SIZE,
)
.ok()?;
let mut spans: Vec<(u64, u64, PageType)> = Vec::new();
for (addr, len) in split.data.into_iter().chain(split.index) {
spans.push((addr.checked_add(base)?, len, PageType::Raw));
}
let mut plain: Vec<(u64, u64)> = spans.iter().map(|&(a, l, _)| (a, l)).collect();
if !spans_disjoint_in_bounds(&mut plain, self.image.len()) {
return None;
}
Some(spans)
}
fn chunked_index_spans(&self, addr: usize) -> Option<Vec<(u64, u64)>> {
let region =
Self::gather_oh_messages(&self.image(), addr as u64, self.superblock.base_address)
.ok()?;
let mut layout_msg: Option<(usize, usize)> = None;
let mut p = 0;
loop {
match next_message(®ion, p) {
Ok(Some((msg_type, body, body_end))) => {
if msg_type == MessageType::DataLayout {
layout_msg = Some((body, body_end));
}
p = body_end;
}
Ok(None) => break,
Err(_) => return None,
}
}
let (lb, le) = layout_msg?;
let layout = DataLayout::parse(®ion[lb..le], OFFSET_SIZE, LENGTH_SIZE).ok()?;
if !matches!(layout, DataLayout::Chunked { .. }) {
return None;
}
let base = self.superblock.base_address;
let mut spans = chunk_index_spans_from_source(
&BaseOffsetSource {
inner: &self.image(),
base,
},
&layout,
OFFSET_SIZE,
LENGTH_SIZE,
)
.ok()?;
for (a, _) in &mut spans {
*a = a.checked_add(base)?;
}
if !spans_disjoint_in_bounds(&mut spans, self.image.len()) {
return None;
}
Some(spans)
}
}
#[derive(Default)]
struct Node {
is_new: bool,
datasets: Vec<DatasetBuilder>,
attr_ops: Vec<AttrOp>,
deletes: Vec<String>,
copies: Vec<(String, CopyTree)>,
writes: Vec<(String, MovingWrite)>,
base_region: Vec<u8>,
existing_links: Vec<String>,
pending_vl_attrs: PendingVlAttrs,
}
enum AttrOp {
Set { name: String, value: AttrValue },
Remove { name: String },
}
enum ObjModel {
DatasetVerbatim {
region: Vec<u8>,
dense_attrs: Vec<crate::attribute::AttributeMessage>,
},
DatasetContiguous {
region: Vec<u8>,
addr_off: usize,
data_addr: u64,
data_size: u64,
dense_attrs: Vec<crate::attribute::AttributeMessage>,
},
DatasetChunked {
region: Vec<u8>,
dense_attrs: Vec<crate::attribute::AttributeMessage>,
},
Group {
non_link_region: Vec<u8>,
children: Vec<(String, u64)>,
dense_attrs: Vec<crate::attribute::AttributeMessage>,
},
}
enum CopyTree {
DatasetVerbatim {
region: Vec<u8>,
dense_attrs: Vec<crate::attribute::AttributeMessage>,
},
DatasetContiguous {
region: Vec<u8>,
addr_off: usize,
data: Vec<u8>,
dense_attrs: Vec<crate::attribute::AttributeMessage>,
},
DatasetChunked {
region: Vec<u8>,
chunk_dims: Vec<u64>,
element_size: usize,
raw_size: u64,
maxshape: Option<Vec<u64>>,
pipeline_message: Option<Vec<u8>>,
meta: Vec<ChunkMeta>,
chunk_bytes: Vec<Vec<u8>>,
dense_attrs: Vec<crate::attribute::AttributeMessage>,
},
Group {
non_link_region: Vec<u8>,
children: Vec<(String, CopyTree)>,
dense_attrs: Vec<crate::attribute::AttributeMessage>,
},
}
struct GroupInfo {
region: Vec<u8>,
link_names: Vec<String>,
}
enum WritePlan {
InPlace { data_addr: usize, raw: Vec<u8> },
InPlaceChunks { writes: Vec<(usize, Vec<u8>)> },
Moving(MovingWrite),
}
enum MovingWrite {
Contiguous {
region: Vec<u8>,
addr_off: usize,
raw: Vec<u8>,
old_extent: Option<(u64, u64)>,
},
Compact { region: Vec<u8>, raw: Vec<u8> },
Chunked {
region: Vec<u8>,
chunk_dims: Vec<u64>,
element_size: usize,
raw_size: u64,
maxshape: Option<Vec<u64>>,
pipeline_message: Option<Vec<u8>>,
meta: Vec<ChunkMeta>,
chunk_bytes: Vec<Vec<u8>>,
old_addr: u64,
},
AppendedChunks {
region: Vec<u8>,
new_dataspace_body: Vec<u8>,
chunk_dims_u32: Vec<u32>,
element_size: usize,
raw_size: u64,
has_filters: bool,
kept_chunks: Vec<WrittenChunk>,
new_chunk_bytes: Vec<Vec<u8>>,
old_addr: u64,
old_tail_extent: Option<(u64, u64)>,
},
AttrEdit {
region: Vec<u8>,
pending_vl_attrs: PendingVlAttrs,
},
}
struct FlatDataset {
name: String,
dt: crate::datatype::Datatype,
ds: Dataspace,
raw: Vec<u8>,
attrs: Vec<crate::attribute::AttributeMessage>,
chunk_options: ChunkOptions,
maxshape: Option<Vec<u64>>,
vl_attrs: Vec<(usize, Vec<Vec<u8>>)>,
vl_string_staging: Option<VlStringStaging>,
reference_targets: Option<Vec<ObjectRefPatch>>,
fill: Option<Vec<u8>>,
}
struct EditStore<'a> {
image: &'a mut dyn FileImage,
superblock: &'a mut Superblock,
sb_sig_off: usize,
paged: Option<&'a mut PagedEdit>,
}
impl EditStore<'_> {
fn append_into_raw_page(&mut self, bytes: &[u8]) -> Result<u64, Error> {
if let Some(pg) = self.paged.as_deref_mut() {
pg.begin(self.image, PageType::Raw)?;
}
self.image.append(bytes)
}
}
impl crate::source::Source for EditStore<'_> {
fn len(&self) -> u64 {
self.image.len()
}
fn read_at(&self, offset: u64, buf: &mut [u8]) -> Result<(), crate::error::FormatError> {
self.image.read_at(offset, buf)
}
fn read_metadata_at(
&self,
offset: u64,
len: usize,
) -> Result<Vec<u8>, crate::error::FormatError> {
self.image.read_metadata_at(offset, len)
}
}
impl Store for EditStore<'_> {
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> {
self.image.append(bytes)
}
fn append_raw(&mut self, bytes: &[u8]) -> Result<u64, Error> {
self.append_into_raw_page(bytes)
}
fn write_at(&mut self, offset: u64, bytes: &[u8]) -> Result<(), Error> {
self.image.write_at(offset, bytes)
}
fn patch_superblock_eof(&mut self) -> Result<(), Error> {
let eof = self.image.len();
self.superblock.eof_address = eof;
let bytes = self.superblock.serialize();
self.write_at(self.sb_sig_off as u64, &bytes)
}
fn sync(&mut self) -> Result<(), Error> {
self.image.sync_data()
}
}
fn paths_overlap(a: &[String], b: &[String]) -> bool {
a.starts_with(b) || b.starts_with(a)
}
pub(crate) fn as_inplace_error(e: Error) -> Error {
match e {
Error::AppendUnsupported(m) => Error::AppendInPlaceUnsupported(m),
other => other,
}
}
pub(crate) fn validate_gathered_append(st: &LocatedState, b: &AppendBuilder) -> Result<u64, Error> {
let raw = b.raw();
if raw.len() % st.element_size != 0 {
return Err(Error::AppendInPlaceUnsupported(
"appended byte length is not a whole number of elements",
));
}
match b.elem_dt() {
Some(expected) if *expected != st.datatype => {
return Err(Error::AppendInPlaceUnsupported(
"append datatype does not match the on-disk dataset (wrong element \
type or byte order)",
));
}
Some(_) => {}
None => {
if !datatype_is_raw_appendable(&st.datatype) {
return Err(Error::AppendInPlaceUnsupported(
"append_raw onto this dataset's datatype (non-little-endian, \
variable-length, or reference) could misencode the bytes; use a \
typed append",
));
}
}
}
Ok((raw.len() / st.element_size) as u64)
}
pub(crate) fn locate_dataset_state<F: Store>(
file: &F,
oh_addr: u64,
) -> Result<LocatedState, Error> {
let result = Located::locate_at(file, oh_addr, Error::AppendInPlaceUnsupported)?;
if result.located.chunk_elems == 0 {
return Err(Error::AppendInPlaceUnsupported(
"in-place append requires a nonzero chunk length",
));
}
let (dt_off, dt_size) = result.spans.datatype;
let dt_bytes = file
.read_metadata_at(dt_off, dt_size)
.map_err(|_| Error::AppendInPlaceUnsupported("dataset datatype could not be parsed"))?;
let (datatype, _) = Datatype::parse(&dt_bytes)
.map_err(|_| Error::AppendInPlaceUnsupported("dataset datatype could not be parsed"))?;
let pipeline = match result.spans.filter {
Some((fb, fsize)) => {
let fp_bytes = file.read_metadata_at(fb, fsize).map_err(|_| {
Error::AppendInPlaceUnsupported("dataset filter pipeline could not be parsed")
})?;
let parsed = FilterPipeline::parse(&fp_bytes).map_err(|_| {
Error::AppendInPlaceUnsupported("dataset filter pipeline could not be parsed")
})?;
if !pipeline_reencodable(&parsed) {
return Err(Error::AppendInPlaceUnsupported(
"dataset uses a filter this engine cannot re-encode",
));
}
Some(parsed)
}
None => None,
};
let element_size = result.located.elem_bytes;
let spatial = vec![result.located.chunk_elems];
Ok(LocatedState {
loc: result.located,
datatype,
spatial,
element_size,
pipeline,
})
}
fn split_path(path: &str) -> PathKey {
path.split('/')
.filter(|s| !s.is_empty())
.map(String::from)
.collect()
}
fn ensure_ancestors(nodes: &mut BTreeMap<PathKey, Node>, path: &[String]) {
for len in 0..=path.len() {
nodes.entry(path[..len].to_vec()).or_default();
}
}
fn spans_disjoint_in_bounds(spans: &mut [(u64, u64)], eof: u64) -> bool {
for &(addr, len) in spans.iter() {
match addr.checked_add(len) {
Some(end) if len > 0 && end <= eof => {}
_ => return false,
}
}
spans.sort_unstable_by_key(|&(addr, _)| addr);
spans.windows(2).all(|w| w[0].0 + w[0].1 <= w[1].0)
}
fn retain_disjoint_in_bounds(spans: &mut Vec<(u64, u64, PageType)>, eof: u64) {
spans.retain(|&(addr, len, _)| len > 0 && addr.checked_add(len).is_some_and(|e| e <= eof));
spans.sort_unstable_by_key(|&(addr, _, _)| addr);
let mut kept_end = 0u64;
spans.retain(|&(addr, len, _)| {
if addr >= kept_end {
kept_end = addr + len;
true
} else {
false }
});
}
fn meta_spans(spans: Vec<(u64, u64)>) -> impl Iterator<Item = (u64, u64, PageType)> {
spans.into_iter().map(|(a, l)| (a, l, PageType::Meta))
}
fn flatten_dataset(db: DatasetBuilder) -> Result<FlatDataset, Error> {
if db.name.is_empty() {
return Err(Error::EditUnsupported("dataset path has an empty name"));
}
let dt = db
.datatype
.ok_or(Error::EditUnsupported("dataset has no datatype/data"))?;
let shape = db
.shape
.ok_or(Error::EditUnsupported("dataset has no shape"))?;
let is_empty = shape.contains(&0);
let chunked = db.chunk_options.is_chunked() || db.maxshape.is_some();
if is_empty && chunked {
return Err(Error::EditUnsupported(
"chunked or extensible empty (zero-element) datasets cannot be added in place yet",
));
}
if db.vl_string_staging.is_some() && chunked {
return Err(Error::EditUnsupported(
"chunked or extensible variable-length-string datasets cannot be added in place yet",
));
}
if db.reference_targets.is_some() && chunked {
return Err(Error::EditUnsupported(
"chunked or extensible object-reference datasets cannot be added in place yet",
));
}
let raw = if is_empty {
db.data.unwrap_or_default()
} else {
db.data
.ok_or(Error::EditUnsupported("dataset has no data"))?
};
let elem = dt.type_size() as u64;
if elem > 0 {
let expected = shape
.iter()
.try_fold(1u64, |acc, &d| acc.checked_mul(d))
.and_then(|n| n.checked_mul(elem));
match expected {
Some(expected) if raw.len() as u64 == expected => {}
Some(_) => {
return Err(Error::EditUnsupported(
"dataset data length does not match its shape",
));
}
None => {
return Err(Error::EditUnsupported(
"dataset shape is too large to address on this platform",
));
}
}
}
if chunked {
db.chunk_options
.validate_geometry(&shape, db.maxshape.as_deref())
.map_err(Error::EditUnsupported)?;
#[cfg(not(feature = "deflate"))]
if db.chunk_options.deflate_level.is_some() {
return Err(Error::EditUnsupported(
"deflate compression requires the `deflate` crate feature",
));
}
let chunk_dims = db.chunk_options.resolve_chunk_dims(&shape);
let ctx = ChunkContext::from_datatype(&chunk_dims, &dt);
db.chunk_options
.build_pipeline(
ctx.element_size,
&chunk_dims,
ctx.element_type,
ctx.scale_offset_type,
)
.map_err(|_| {
Error::EditUnsupported(
"this dataset's filter pipeline cannot be added in place \
(an unsupported filter, an incompatible datatype, or a \
compression feature that is not enabled)",
)
})?;
}
if make_link(&db.name, 0).serialize(OFFSET_SIZE).len() > OBJECT_HEADER_MESSAGE_MAX {
return Err(Error::EditUnsupported(
"dataset name is too long to encode as a link message",
));
}
let ds = Dataspace {
space_type: if shape.is_empty() {
DataspaceType::Scalar
} else {
DataspaceType::Simple
},
#[expect(
clippy::cast_possible_truncation,
reason = "dataspace rank fits the 1-byte dimensionality field (HDF5 caps rank at 32)"
)]
rank: shape.len() as u8,
dimensions: shape,
max_dimensions: db.maxshape.clone(),
};
let mut attrs: Vec<crate::attribute::AttributeMessage> = Vec::with_capacity(db.attrs.len());
for (n, v) in &db.attrs {
attrs.push(build_attr_message(n, v));
}
let mut vl_attrs: Vec<(usize, Vec<Vec<u8>>)> = Vec::new();
for (i, (_, v)) in db.attrs.iter().enumerate() {
if let AttrValue::VarLenAsciiArray(strings) = v {
let str_refs: Vec<&str> = strings.iter().map(String::as_str).collect();
vl_attrs.push((i, build_global_heap_collections(&str_refs)));
}
}
#[cfg(feature = "provenance")]
if let Some(ref prov) = db.provenance {
let p = crate::provenance::Provenance {
creator: prov.creator.clone(),
timestamp: prov.timestamp.clone(),
source: prov.source.clone(),
};
attrs.extend(p.build_attrs(&raw));
}
for a in &attrs {
if a.serialize(LENGTH_SIZE).len() > OBJECT_HEADER_MESSAGE_MAX {
return Err(Error::EditUnsupported(
"dataset attribute is too large to encode in place",
));
}
}
if attrs.len() > MAX_COMPACT_ATTRS {
return Err(Error::EditUnsupported(
"datasets with dense (many) attributes cannot be added in place yet",
));
}
if let Some(fill) = &db.fill {
let expected = elem.to_usize()?;
if fill.len() != expected {
return Err(Error::Format(FormatError::FillValueSizeMismatch {
expected,
actual: fill.len(),
}));
}
}
Ok(FlatDataset {
name: db.name,
dt,
ds,
raw,
attrs,
chunk_options: db.chunk_options,
maxshape: db.maxshape,
vl_attrs,
vl_string_staging: db.vl_string_staging,
reference_targets: db.reference_targets,
fill: db.fill,
})
}
const GROUP_INFO_BODY: [u8; 2] = [0, 0];
fn chunk_index_enumerable(version: u8, chunk_index_type: Option<u8>) -> bool {
matches!((version, chunk_index_type), (3, _) | (4, Some(1..=4)))
}
pub(crate) fn pipeline_reencodable(pipeline: &FilterPipeline) -> bool {
pipeline.filters.iter().all(|f| match f.filter_id {
FILTER_DEFLATE | FILTER_SHUFFLE | FILTER_FLETCHER32 | FILTER_SCALEOFFSET | FILTER_LZF => {
true
}
#[cfg(feature = "zfp")]
crate::filter_pipeline::FILTER_ZFP => true,
_ => false,
})
}
fn replace_layout_message(region: &[u8], new_layout_body: &[u8]) -> Result<Vec<u8>, Error> {
let mut out = Vec::with_capacity(region.len());
let mut p = 0;
let mut replaced = false;
while let Some((msg_type, _body, body_end)) = next_message(region, p)? {
if msg_type == MessageType::DataLayout && !replaced {
out.extend_from_slice(®ion_message(MessageType::DataLayout, new_layout_body));
replaced = true;
} else {
out.extend_from_slice(®ion[p..body_end]);
}
p = body_end;
}
if !replaced {
return Err(Error::EditUnsupported(
"chunked dataset header has no data-layout message to relocate",
));
}
Ok(out)
}
fn replace_dataspace_message(region: &[u8], new_dataspace_body: &[u8]) -> Result<Vec<u8>, Error> {
let mut out = Vec::with_capacity(region.len());
let mut p = 0;
let mut replaced = false;
while let Some((msg_type, _body, body_end)) = next_message(region, p)? {
if msg_type == MessageType::Dataspace && !replaced {
out.extend_from_slice(®ion_message(MessageType::Dataspace, new_dataspace_body));
replaced = true;
} else {
out.extend_from_slice(®ion[p..body_end]);
}
p = body_end;
}
if !replaced {
return Err(Error::AppendUnsupported(
"dataset header has no dataspace message to grow",
));
}
Ok(out)
}
pub(crate) fn datatype_is_raw_appendable(dt: &Datatype) -> bool {
match dt {
Datatype::FixedPoint { byte_order, .. }
| Datatype::FloatingPoint { byte_order, .. }
| Datatype::Time { byte_order, .. }
| Datatype::BitField { byte_order, .. } => *byte_order == DatatypeByteOrder::LittleEndian,
Datatype::String { .. } | Datatype::Opaque { .. } => true,
Datatype::Enumeration { base_type, .. } | Datatype::Array { base_type, .. } => {
datatype_is_raw_appendable(base_type)
}
Datatype::Compound { members, .. } => members
.iter()
.all(|m| datatype_is_raw_appendable(&m.datatype)),
Datatype::VariableLength { .. } | Datatype::Reference { .. } => false,
}
}
struct ChunkedHeaderParts {
dt: crate::datatype::Datatype,
ds: Dataspace,
layout: DataLayout,
pipeline_message: Option<Vec<u8>>,
}
fn parse_chunked_header(region: &[u8]) -> Result<ChunkedHeaderParts, Error> {
let mut datatype: Option<(usize, usize)> = None;
let mut dataspace: Option<(usize, usize)> = None;
let mut layout: Option<(usize, usize)> = None;
let mut pipeline: Option<(usize, usize)> = None;
let mut p = 0;
while let Some((msg_type, body, body_end)) = next_message(region, p)? {
match msg_type {
MessageType::Datatype => datatype = Some((body, body_end)),
MessageType::Dataspace => dataspace = Some((body, body_end)),
MessageType::DataLayout => layout = Some((body, body_end)),
MessageType::FilterPipeline => pipeline = Some((body, body_end)),
_ => {}
}
p = body_end;
}
let (dt_b, dt_e) = datatype.ok_or(Error::EditUnsupported("dataset header has no datatype"))?;
let (ds_b, ds_e) =
dataspace.ok_or(Error::EditUnsupported("dataset header has no dataspace"))?;
let (lb, le) = layout.ok_or(Error::EditUnsupported("dataset header has no data layout"))?;
let (dt, _) = crate::datatype::Datatype::parse(®ion[dt_b..dt_e])
.map_err(|_| Error::EditUnsupported("dataset header datatype could not be parsed"))?;
let ds = Dataspace::parse(®ion[ds_b..ds_e], LENGTH_SIZE)
.map_err(|_| Error::EditUnsupported("dataset header dataspace could not be parsed"))?;
let dl = DataLayout::parse(®ion[lb..le], OFFSET_SIZE, LENGTH_SIZE)
.map_err(|_| Error::EditUnsupported("dataset header data layout could not be parsed"))?;
if !matches!(dl, DataLayout::Chunked { .. }) {
return Err(Error::EditUnsupported("dataset is not chunked"));
}
let pipeline_message = pipeline.map(|(b, e)| region[b..e].to_vec());
Ok(ChunkedHeaderParts {
dt,
ds,
layout: dl,
pipeline_message,
})
}
struct ChunkedGeometry {
spatial: Vec<u64>,
element_size: usize,
raw_size: u64,
maxshape: Option<Vec<u64>>,
}
fn chunked_geometry(
dt: &crate::datatype::Datatype,
ds: &Dataspace,
layout: &DataLayout,
) -> Result<ChunkedGeometry, Error> {
let DataLayout::Chunked {
chunk_dimensions, ..
} = layout
else {
return Err(Error::EditUnsupported("dataset is not chunked"));
};
let rank = ds.dimensions.len();
if chunk_dimensions.len() <= rank {
return Err(Error::EditUnsupported(
"chunked layout has malformed dimensions",
));
}
let spatial: Vec<u64> = chunk_dimensions[..rank]
.iter()
.map(|&c| u64::from(c))
.collect();
let element_size = dt.type_size() as usize;
if element_size == 0 {
return Err(Error::EditUnsupported(
"chunked dataset has a zero element size",
));
}
let raw_size = spatial
.iter()
.copied()
.product::<u64>()
.saturating_mul(element_size as u64);
let maxshape = ds
.max_dimensions
.as_ref()
.filter(|ms| *ms != &ds.dimensions)
.cloned();
Ok(ChunkedGeometry {
spatial,
element_size,
raw_size,
maxshape,
})
}
fn try_inplace_chunk_writes<S: Source + ?Sized>(
src: &S,
layout: &DataLayout,
ds: &Dataspace,
spatial: &[u64],
raw_size: u64,
new_bytes: &[Vec<u8>],
) -> Option<Vec<(usize, Vec<u8>)>> {
let infos = enumerate_chunks_from_source(src, layout, ds, OFFSET_SIZE, LENGTH_SIZE).ok()?;
let grid = plan_dense_grid(infos, &ds.dimensions, spatial)?;
if grid.grid_order.len() != new_bytes.len() {
return None;
}
let mut writes = Vec::with_capacity(new_bytes.len() + 1);
let mut spans: Vec<(u64, u64)> = Vec::with_capacity(new_bytes.len() + 1);
let mut any_shrunk = false;
for (ci, bytes) in grid.grid_order.iter().zip(new_bytes.iter()) {
if ci.filter_mask != 0 {
return None;
}
let new_len = bytes.len() as u64;
let slot = u64::from(ci.chunk_size);
if new_len > slot {
return None;
}
if new_len < slot {
any_shrunk = true;
}
let start = usize::try_from(ci.address).ok()?;
start
.checked_add(bytes.len())
.filter(|&e| e as u64 <= src.len())?;
writes.push((start, bytes.clone()));
spans.push((ci.address, new_len));
}
if any_shrunk {
let (index_addr, index_bytes) =
try_rebuild_index_in_place(src, layout, raw_size, &grid.grid_order, new_bytes)?;
spans.push((index_addr as u64, index_bytes.len() as u64));
writes.push((index_addr, index_bytes));
}
if !spans_disjoint_in_bounds(&mut spans, src.len()) {
return None;
}
Some(writes)
}
fn try_rebuild_index_in_place<S: Source + ?Sized>(
src: &S,
layout: &DataLayout,
raw_size: u64,
grid_order: &[crate::chunked_read::ChunkInfo],
new_bytes: &[Vec<u8>],
) -> Option<(usize, Vec<u8>)> {
let DataLayout::Chunked {
btree_address: Some(index_addr),
chunk_index_type,
version,
..
} = layout
else {
return None;
};
let written: Vec<crate::chunked_write::WrittenChunk> = grid_order
.iter()
.zip(new_bytes)
.map(|(ci, b)| crate::chunked_write::WrittenChunk {
address: ci.address,
compressed_size: b.len() as u64,
raw_size,
filter_mask: 0,
})
.collect();
let new_index = match (version, chunk_index_type) {
(4, Some(3)) => crate::chunked_write::build_fixed_array_at(
&written,
OFFSET_SIZE,
LENGTH_SIZE,
true,
*index_addr,
),
(4, Some(4)) => crate::chunked_write::build_extensible_array_at(
&written,
OFFSET_SIZE,
LENGTH_SIZE,
true,
*index_addr,
)
.ok()?,
_ => return None,
};
let mut spans =
crate::chunked_read::chunk_index_spans_from_source(src, layout, OFFSET_SIZE, LENGTH_SIZE)
.ok()?;
if spans.is_empty() {
return None;
}
spans.sort_unstable_by_key(|&(a, _)| a);
if spans[0].0 != *index_addr {
return None;
}
let mut end = *index_addr;
for &(a, l) in &spans {
if a != end {
return None; }
end = a.checked_add(l)?;
}
if new_index.len() as u64 != end - *index_addr {
return None;
}
let start = usize::try_from(*index_addr).ok()?;
start
.checked_add(new_index.len())
.filter(|&e| e as u64 <= src.len())?;
Some((start, new_index))
}
struct SliceChunkProvider<'a> {
chunks: &'a [Vec<u8>],
}
impl ChunkProvider for SliceChunkProvider<'_> {
fn chunk_bytes(&self, index: usize, out: &mut Vec<u8>) -> Result<(), FormatError> {
let chunk = self.chunks.get(index).ok_or_else(|| {
FormatError::ChunkedReadError("chunk index out of range for in-memory provider".into())
})?;
out.extend_from_slice(chunk);
Ok(())
}
}
fn region_message(msg_type: MessageType, body: &[u8]) -> Vec<u8> {
let mut m = Vec::with_capacity(4 + body.len());
#[expect(
clippy::cast_possible_truncation,
reason = "message type ids are a small enum that fits the 1-byte v2 type field"
)]
m.push(msg_type.to_u16() as u8);
#[expect(
clippy::cast_possible_truncation,
reason = "callers pass bodies that fit the 2-byte message-size field (see doc comment)"
)]
m.extend_from_slice(&(body.len() as u16).to_le_bytes());
m.push(0); m.extend_from_slice(body);
m
}
fn fresh_group_region() -> Vec<u8> {
let mut li = Vec::with_capacity(18);
li.push(0); li.push(0); li.extend_from_slice(&u64::MAX.to_le_bytes()); li.extend_from_slice(&u64::MAX.to_le_bytes()); let mut region = region_message(MessageType::LinkInfo, &li);
region.extend_from_slice(®ion_message(MessageType::GroupInfo, &GROUP_INFO_BODY));
region
}
fn ensure_group_info(region: &mut Vec<u8>) -> Result<(), Error> {
let mut p = 0;
while let Some((msg_type, _body, body_end)) = next_message(region, p)? {
if msg_type == MessageType::GroupInfo {
return Ok(());
}
p = body_end;
}
region.extend_from_slice(®ion_message(MessageType::GroupInfo, &GROUP_INFO_BODY));
Ok(())
}
fn encode_link_message(name: &str, addr: u64) -> Vec<u8> {
let body = make_link(name, addr).serialize(OFFSET_SIZE);
region_message(MessageType::Link, &body)
}
fn patch_link_target(region: &mut [u8], name: &str, new_addr: u64) -> Result<(), Error> {
let mut p = 0;
while let Some((msg_type, body, body_end)) = next_message(region, p)? {
if msg_type == MessageType::Link {
if let Ok(link) = LinkMessage::parse(®ion[body..body_end], OFFSET_SIZE) {
if link.name == name {
return match link.link_target {
LinkTarget::Hard { .. } => {
let ofs = body_end - OFFSET_SIZE as usize;
region[ofs..body_end].copy_from_slice(&new_addr.to_le_bytes());
Ok(())
}
_ => Err(Error::EditUnsupported(
"a group on the edited path is reached by a soft/external link",
)),
};
}
}
}
p = body_end;
}
Err(Error::EditUnsupported(
"expected child link not found in parent group",
))
}
const COMPACT_LAYOUT_PREAMBLE: usize = 4;
fn rebuild_compact_layout_region(region: &[u8], raw: &[u8]) -> Result<Vec<u8>, Error> {
if raw.len() > OBJECT_HEADER_MESSAGE_MAX - COMPACT_LAYOUT_PREAMBLE {
return Err(Error::EditUnsupported(
"compact dataset data is too large to overwrite in place",
));
}
let mut out = Vec::with_capacity(region.len() + raw.len());
let mut p = 0;
let mut replaced = false;
while let Some((msg_type, body, body_end)) = next_message(region, p)? {
if msg_type == MessageType::DataLayout {
if body_end - body < 2 || region[body + 1] != 0 {
return Err(Error::EditUnsupported(
"compact-layout overwrite found a non-compact data layout",
));
}
let mut layout = Vec::with_capacity(COMPACT_LAYOUT_PREAMBLE + raw.len());
layout.push(region[body]); layout.push(0); #[expect(
clippy::cast_possible_truncation,
reason = "raw.len() bounded below the u16 inline-size field above"
)]
layout.extend_from_slice(&(raw.len() as u16).to_le_bytes());
layout.extend_from_slice(raw);
out.push(region[p]);
#[expect(
clippy::cast_possible_truncation,
reason = "the guard above bounds COMPACT_LAYOUT_PREAMBLE + raw.len(), this \
body's exact length, to the 2-byte message-size field"
)]
out.extend_from_slice(&(layout.len() as u16).to_le_bytes());
out.push(region[p + 3]);
out.extend_from_slice(&layout);
replaced = true;
} else {
out.extend_from_slice(®ion[p..body_end]);
}
p = body_end;
}
if p < region.len() {
out.extend_from_slice(®ion[p..]);
}
if !replaced {
return Err(Error::EditUnsupported(
"compact dataset header has no data-layout message",
));
}
Ok(out)
}
fn remove_link_from_region(region: &[u8], name: &str) -> Result<Vec<u8>, Error> {
let mut out = Vec::with_capacity(region.len());
let mut p = 0;
let mut removed = false;
while let Some((msg_type, body, body_end)) = next_message(region, p)? {
let mut skip = false;
if msg_type == MessageType::Link {
if let Ok(link) = LinkMessage::parse(®ion[body..body_end], OFFSET_SIZE) {
if link.name == name {
skip = true;
removed = true;
}
}
}
if !skip {
out.extend_from_slice(®ion[p..body_end]);
}
p = body_end;
}
if p < region.len() {
out.extend_from_slice(®ion[p..]);
}
if !removed {
return Err(Error::EditUnsupported(
"link to delete not found in its parent group",
));
}
Ok(out)
}
fn apply_group_attr_ops(region: &[u8], ops: &[AttrOp]) -> Result<(Vec<u8>, PendingVlAttrs), Error> {
let mut out = region.to_vec();
let mut pending_vl: PendingVlAttrs = Vec::new();
let mut wrote_attr = false;
for op in ops {
match op {
AttrOp::Set { name, value } => {
wrote_attr = true;
pending_vl.retain(|(msg, _)| &msg.name != name);
if let AttrValue::VarLenAsciiArray(strings) = value {
out = remove_attr_from_region(&out, name, false)?;
let msg = build_attr_message(name, value);
if msg.serialize(LENGTH_SIZE).len() > OBJECT_HEADER_MESSAGE_MAX {
return Err(Error::EditUnsupported(
"attribute is too large to encode in place",
));
}
let str_refs: Vec<&str> = strings.iter().map(String::as_str).collect();
pending_vl.push((msg, build_global_heap_collections(&str_refs)));
} else {
out = set_attr_in_region(&out, name, value)?;
}
}
AttrOp::Remove { name } => {
let before = pending_vl.len();
pending_vl.retain(|(msg, _)| &msg.name != name);
if pending_vl.len() == before {
out = remove_attr_from_region(&out, name, true)?;
}
}
}
}
if wrote_attr && compact_attr_count(&out)? + pending_vl.len() > MAX_COMPACT_ATTRS {
return Err(Error::EditUnsupported(
"attributes would exceed compact storage; dense attribute edits are not supported in place yet",
));
}
Ok((out, pending_vl))
}
fn attribute_info_is_dense(body: &[u8]) -> bool {
match crate::attribute_info::AttributeInfoMessage::parse(body, OFFSET_SIZE) {
Ok(ai) => ai.fractal_heap_address.is_some(),
Err(_) => true,
}
}
fn set_attr_in_region(region: &[u8], name: &str, value: &AttrValue) -> Result<Vec<u8>, Error> {
let new_msg = encode_attr_message(name, value)?;
let mut out = Vec::with_capacity(region.len() + new_msg.len());
let mut p = 0;
while let Some((msg_type, body, body_end)) = next_message(region, p)? {
match msg_type {
MessageType::AttributeInfo => {
if attribute_info_is_dense(®ion[body..body_end]) {
return Err(Error::EditUnsupported(
"a target object uses dense (fractal-heap) attribute storage (not supported in place yet)",
));
}
}
MessageType::Attribute => {
let attr_name = parse_compact_attr_name(region, p, body, body_end)?;
if attr_name == name {
p = body_end;
continue;
}
}
_ => {}
}
out.extend_from_slice(®ion[p..body_end]);
p = body_end;
}
out.extend_from_slice(&new_msg);
if p < region.len() {
out.extend_from_slice(®ion[p..]);
}
Ok(out)
}
fn remove_attr_from_region(region: &[u8], name: &str, required: bool) -> Result<Vec<u8>, Error> {
let mut out = Vec::with_capacity(region.len());
let mut p = 0;
let mut removed = false;
while let Some((msg_type, body, body_end)) = next_message(region, p)? {
let mut skip = false;
match msg_type {
MessageType::AttributeInfo => {
if attribute_info_is_dense(®ion[body..body_end]) {
return Err(Error::EditUnsupported(
"a target object uses dense (fractal-heap) attribute storage (not supported in place yet)",
));
}
}
MessageType::Attribute => {
let attr_name = parse_compact_attr_name(region, p, body, body_end)?;
if attr_name == name {
skip = true;
removed = true;
}
}
_ => {}
}
if !skip {
out.extend_from_slice(®ion[p..body_end]);
}
p = body_end;
}
if p < region.len() {
out.extend_from_slice(®ion[p..]);
}
if !removed && required {
return Err(Error::EditUnsupported("attribute to remove was not found"));
}
Ok(out)
}
fn compact_attr_count(region: &[u8]) -> Result<usize, Error> {
let mut count = 0usize;
let mut p = 0;
while let Some((msg_type, body, body_end)) = next_message(region, p)? {
if msg_type == MessageType::AttributeInfo
&& attribute_info_is_dense(®ion[body..body_end])
{
return Err(Error::EditUnsupported(
"a target object uses dense (fractal-heap) attribute storage (not supported in place yet)",
));
}
if msg_type == MessageType::Attribute {
count += 1;
}
p = body_end;
}
Ok(count)
}
fn parse_compact_attr_name(
region: &[u8],
msg_start: usize,
body: usize,
body_end: usize,
) -> Result<String, Error> {
if region[msg_start + 3] != 0 {
return Err(Error::EditUnsupported(
"a target object has a shared attribute message (not editable in place yet)",
));
}
crate::attribute::AttributeMessage::parse(®ion[body..body_end], LENGTH_SIZE)
.map(|attr| attr.name)
.map_err(|_| Error::EditUnsupported("a target object has an unreadable attribute message"))
}
fn encode_attr_message(name: &str, value: &AttrValue) -> Result<Vec<u8>, Error> {
debug_assert!(
!matches!(value, AttrValue::VarLenAsciiArray(_)),
"VarLenAsciiArray must be intercepted by apply_group_attr_ops before reaching encode_attr_message"
);
let body = build_attr_message(name, value).serialize(LENGTH_SIZE);
if body.len() > OBJECT_HEADER_MESSAGE_MAX {
return Err(Error::EditUnsupported(
"group attribute is too large to encode in place",
));
}
Ok(region_message(MessageType::Attribute, &body))
}
fn is_prefix(a: &[String], b: &[String]) -> bool {
a.len() <= b.len() && b[..a.len()] == *a
}
pub(crate) fn rewrite_extension_region_bytes(
region: &[u8],
info: &FileSpaceInfo,
) -> Result<Vec<u8>, Error> {
let new_body = info.serialize();
let new_len = u16::try_from(new_body.len())
.map_err(|_| Error::EditUnsupported("File Space Info message too large"))?;
let mut out = Vec::with_capacity(region.len());
let mut p = 0;
let mut replaced = false;
while let Some((msg_type, _body, body_end)) = next_message(region, p)? {
if msg_type == MessageType::FileSpaceInfo {
out.push(region[p]); out.extend_from_slice(&new_len.to_le_bytes());
out.push(region[p + 3]); out.extend_from_slice(&new_body);
replaced = true;
} else {
out.extend_from_slice(®ion[p..body_end]);
}
p = body_end;
}
if !replaced {
return Err(Error::EditUnsupported(
"a persisting file's superblock extension has no File Space Info message",
));
}
Ok(out)
}
fn oh_region_at(prefix: &[u8], addr: u64, file_len: u64) -> Result<(u64, u64), Error> {
if prefix.len() < 6 || &prefix[..4] != b"OHDR" || prefix[4] != 2 {
return Err(Error::EditUnsupported(
"an object does not use a version 2 object header",
));
}
let flags = prefix[5];
if flags & 0x04 != 0 {
return Err(Error::EditUnsupported(
"an object tracks message creation order (not supported in place yet)",
));
}
let mut pos = 6usize;
if flags & 0x20 != 0 {
pos += 16; }
if flags & 0x10 != 0 {
pos += 4; }
let size_width = match flags & 0x03 {
0 => 1usize,
1 => 2,
2 => 4,
_ => 8,
};
if prefix.len() < pos + size_width {
return Err(Error::EditUnsupported("truncated object header"));
}
let chunk0_size = read_le(&prefix[pos..pos + size_width]) as u64;
pos += size_width;
let region_start = addr
.checked_add(pos as u64)
.ok_or(Error::EditUnsupported("truncated object header"))?;
let region_end = region_start
.checked_add(chunk0_size)
.filter(|e| e.checked_add(4).is_some_and(|end| end <= file_len))
.ok_or(Error::EditUnsupported("truncated object header"))?;
Ok((region_start, region_end))
}
struct OhChunk {
span: (u64, u64),
buf: Vec<u8>,
messages_start: usize,
}
impl OhChunk {
fn message_region(&self) -> (&[u8], usize) {
(&self.buf, self.messages_start)
}
}
fn read_oh_chunk0<S: Source + ?Sized>(src: &S, addr: u64) -> Result<OhChunk, Error> {
let file_len = src.len();
let window = file_len
.saturating_sub(addr)
.min(OH_PREFIX_MAX as u64)
.to_usize()?;
let prefix = src.read_metadata_at(addr, window)?;
let (rs, re) = oh_region_at(&prefix, addr, file_len)?;
let len = (re - addr).to_usize()?;
Ok(OhChunk {
span: (addr, len as u64 + 4),
buf: src.read_metadata_at(addr, len)?,
messages_start: (rs - addr).to_usize()?,
})
}
fn read_oh_chunks<S: Source + ?Sized>(
src: &S,
addr: u64,
base: u64,
) -> Result<Vec<OhChunk>, Error> {
let mut chunks = vec![read_oh_chunk0(src, addr)?];
let mut i = 0;
while i < chunks.len() {
if chunks.len() > MAX_OH_CHUNKS {
return Err(Error::EditUnsupported(
"object header has too many continuation chunks",
));
}
let mut found = Vec::new();
let (region, mut p) = chunks[i].message_region();
while let Some((msg_type, body, body_end)) = next_message(region, p)? {
if msg_type == MessageType::ObjectHeaderContinuation {
found.push(read_oh_continuation(src, region, body, body_end, base)?);
}
p = body_end;
}
i += 1;
chunks.extend(found);
}
Ok(chunks)
}
fn read_oh_continuation<S: Source + ?Sized>(
src: &S,
region: &[u8],
body: usize,
body_end: usize,
base: u64,
) -> Result<OhChunk, Error> {
if body_end - body < (OFFSET_SIZE + LENGTH_SIZE) as usize {
return Err(Error::EditUnsupported("malformed continuation message"));
}
let off = u64::from_le_bytes(region[body..body + 8].try_into().unwrap());
let len = u64::from_le_bytes(region[body + 8..body + 16].try_into().unwrap());
let off = off
.checked_add(base)
.ok_or(Error::EditUnsupported("continuation address overflow"))?;
let end = off
.checked_add(len)
.filter(|&e| e <= src.len() && len >= 8)
.ok_or(Error::EditUnsupported("continuation block out of bounds"))?;
let want = (end - off)
.to_usize()
.map_err(|_| Error::EditUnsupported("continuation length exceeds this platform"))?;
let mut buf = src.read_metadata_at(off, want)?;
if buf[..4] != *b"OCHK" {
return Err(Error::EditUnsupported(
"invalid continuation block signature",
));
}
buf.truncate(want - 4);
Ok(OhChunk {
span: (off, len),
buf,
messages_start: 4,
})
}
pub(crate) fn next_message(
region: &[u8],
p: usize,
) -> Result<Option<(MessageType, usize, usize)>, Error> {
if p + 4 > region.len() {
return Ok(None);
}
let msg_type = MessageType::from_u16(region[p] as u16);
let msg_size = u16::from_le_bytes([region[p + 1], region[p + 2]]) as usize;
let body = p + 4;
let body_end = body + msg_size;
if body_end > region.len() {
return Err(Error::EditUnsupported("malformed object header message"));
}
Ok(Some((msg_type, body, body_end)))
}
const MSG_FLAG_SHARED: u8 = 0x02;
fn reject_foreign_addresses(region: &[u8]) -> Result<(), Error> {
let mut p = 0;
while let Some((msg_type, body, body_end)) = next_message(region, p)? {
if region[p + 3] & MSG_FLAG_SHARED != 0 {
return Err(Error::EditUnsupported(
"a shared (committed/SOHM) object-header message cannot be copied to another file yet",
));
}
match msg_type {
MessageType::Datatype => {
let (dt, _) =
crate::datatype::Datatype::parse(®ion[body..body_end]).map_err(|_| {
Error::EditUnsupported("a source datatype could not be parsed for copying")
})?;
if datatype_copies_foreign_address(&dt) {
return Err(Error::EditUnsupported(
"variable-length or reference datasets cannot be copied to another file yet",
));
}
}
MessageType::Attribute => {
let attr =
crate::attribute::AttributeMessage::parse(®ion[body..body_end], LENGTH_SIZE)
.map_err(|_| {
Error::EditUnsupported(
"a source attribute could not be parsed for copying",
)
})?;
if datatype_copies_foreign_address(&attr.datatype) {
return Err(Error::EditUnsupported(
"variable-length or reference attributes cannot be copied to another file yet",
));
}
}
_ => {}
}
p = body_end;
}
Ok(())
}
fn reject_foreign_dense_attrs(attrs: &[crate::attribute::AttributeMessage]) -> Result<(), Error> {
for attr in attrs {
if datatype_copies_foreign_address(&attr.datatype) {
return Err(Error::EditUnsupported(
"variable-length or reference dense (fractal-heap) attributes cannot be copied to another file yet",
));
}
}
Ok(())
}
fn datatype_copies_foreign_address(dt: &crate::datatype::Datatype) -> bool {
use crate::datatype::Datatype;
match dt {
Datatype::VariableLength { .. } | Datatype::Reference { .. } => true,
Datatype::Compound { members, .. } => members
.iter()
.any(|m| datatype_copies_foreign_address(&m.datatype)),
Datatype::Array { base_type, .. } | Datatype::Enumeration { base_type, .. } => {
datatype_copies_foreign_address(base_type)
}
_ => false,
}
}
pub(crate) fn build_v2_object_header(region: &[u8]) -> Vec<u8> {
let total = region.len();
let (flags, width) = if total <= 255 {
(0u8, 1usize)
} else if total <= 65535 {
(1u8, 2)
} else {
(2u8, 4)
};
let mut buf = Vec::with_capacity(8 + total + 4);
buf.extend_from_slice(b"OHDR");
buf.push(2); buf.push(flags);
#[expect(
clippy::cast_possible_truncation,
reason = "width was selected just above to be the smallest field that holds total"
)]
match width {
1 => buf.push(total as u8),
2 => buf.extend_from_slice(&(total as u16).to_le_bytes()),
_ => buf.extend_from_slice(&(total as u32).to_le_bytes()),
}
buf.extend_from_slice(region);
let checksum = jenkins_lookup3(&buf);
buf.extend_from_slice(&checksum.to_le_bytes());
buf
}
#[expect(
clippy::cast_possible_truncation,
reason = "callers parse in-file sizes/offsets bounded by the in-memory image; downstream \
slicing is length-checked, so a malformed oversized field errors rather than reads OOB"
)]
fn read_le(bytes: &[u8]) -> usize {
let mut v = 0u64;
for (i, &b) in bytes.iter().enumerate() {
v |= (b as u64) << (8 * i);
}
v as usize
}
#[cfg(test)]
mod tests {
use super::*;
fn region_types(region: &[u8]) -> Vec<MessageType> {
let mut out = Vec::new();
let mut p = 0;
while let Some((mt, _, end)) = next_message(region, p).unwrap() {
out.push(mt);
p = end;
}
out
}
#[test]
fn append_inplace_crash_consistency_partial_tail_prefix() {
use crate::reader::File as PureFile;
use crate::writer::FileBuilder;
use tempfile::tempdir;
let build = |path: &std::path::Path, n: i32, chunk: u64| {
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.write(path).unwrap();
};
for (n, chunk, add) in [(6i32, 4u64, 5i32), (9, 2, 6)] {
let dir = tempdir().unwrap();
let base = dir.path().join("base.h5");
build(&base, n, chunk);
for max_phase in 1u8..=4 {
let p = dir.path().join(format!("crash_{n}_{chunk}_{max_phase}.h5"));
std::fs::copy(&base, &p).unwrap();
{
let mut s = WriteEngine::open_with_locking(&p, FileLocking::Enabled).unwrap();
s.append_inplace_i32_phased("d", &(n..n + add).collect::<Vec<_>>(), max_phase)
.unwrap();
}
let expected_len = if max_phase == 4 { n + add } else { n };
let f = PureFile::from_bytes(std::fs::read(&p).unwrap()).unwrap();
assert_eq!(
f.dataset("d").unwrap().read_i32().unwrap(),
(0..expected_len).collect::<Vec<_>>(),
"inconsistent view after crash at phase {max_phase} (n={n}, chunk={chunk})"
);
}
}
}
fn build_unit_chunked(path: &std::path::Path, n: i32) {
use crate::writer::FileBuilder;
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(&[1]);
b.write(path).unwrap();
}
fn append_stopped_at(
base: &std::path::Path,
out: &std::path::Path,
values: std::ops::Range<i32>,
max_phase: u8,
) {
std::fs::copy(base, out).unwrap();
let mut s = WriteEngine::open_with_locking(out, FileLocking::Enabled).unwrap();
s.append_inplace_i32_phased("d", &values.collect::<Vec<_>>(), max_phase)
.unwrap();
}
#[test]
fn append_inplace_crash_consistency_across_ea_boundaries() {
use crate::reader::File as PureFile;
use tempfile::tempdir;
let dir = tempdir().unwrap();
let base = dir.path().join("base.h5");
let (n, target) = (50i32, 250i32);
build_unit_chunked(&base, n);
for max_phase in 1u8..=4 {
let p = dir.path().join(format!("crash_ea_{max_phase}.h5"));
append_stopped_at(&base, &p, n..target, max_phase);
let expected_len = if max_phase == 4 { target } else { n };
let f = PureFile::from_bytes(std::fs::read(&p).unwrap()).unwrap();
assert_eq!(
f.dataset("d").unwrap().read_i32().unwrap(),
(0..expected_len).collect::<Vec<_>>(),
"inconsistent view after crash at phase {max_phase}"
);
}
}
#[test]
fn paged_open_seeds_each_manager_by_slot() {
use crate::writer::FileBuilder;
use tempfile::tempdir;
let dir = tempdir().unwrap();
let path = dir.path().join("paged_seed.h5");
let mut b = FileBuilder::new();
b.create_dataset("d")
.with_i32_data(&(0..1000).collect::<Vec<i32>>())
.with_shape(&[1000]);
b.with_file_space_strategy(FileSpaceStrategy::Page, true, 0)
.with_file_space_page_size(4096);
b.write(&path).unwrap();
let on_disk: u64 = crate::reader::File::open(&path)
.unwrap()
.persisted_free_space()
.iter()
.map(|&(_, l)| l)
.sum();
let s = WriteEngine::open_with_locking(&path, FileLocking::Enabled).unwrap();
let pg = s.paged.as_ref().expect("a paged file installs paged state");
assert_eq!(pg.page_size, 4096);
assert!(
!pg.meta.sections().is_empty(),
"SUPER (slot 0) sections seed the metadata list"
);
assert!(
!pg.raw_small.sections().is_empty(),
"DRAW (slot 2) sections seed the small-raw list, not the metadata list"
);
let mut all = pg.all_sections();
let flat: u64 = all.iter().map(|&(_, l)| l).sum();
assert_eq!(
flat, on_disk,
"the split lists hold exactly the file's free space"
);
all.sort_by_key(|&(a, _)| a);
let mut prev_end = 0u64;
for (addr, len) in all {
assert!(addr >= prev_end, "the per-type lists do not overlap");
prev_end = addr + len;
}
}
#[test]
fn deleted_chunk_index_is_freed_into_a_raw_manager() {
use crate::writer::FileBuilder;
use tempfile::tempdir;
let dir = tempdir().unwrap();
let path = dir.path().join("paged_chunk_index_free.h5");
let page = 4096u64;
let mut b = FileBuilder::new();
for name in ["drop", "keep"] {
b.create_dataset(name)
.with_i32_data(&(0..200).collect::<Vec<i32>>())
.with_shape(&[200])
.with_chunks(&[50]);
}
b.with_file_space_strategy(FileSpaceStrategy::Page, true, 0)
.with_file_space_page_size(page);
b.write(&path).unwrap();
{
let mut s = WriteEngine::open_with_locking(&path, FileLocking::Enabled).unwrap();
s.delete("/drop");
s.commit().unwrap();
}
let live_raw_pages: Vec<u64> = {
let f = crate::reader::File::open(&path).unwrap();
let ds = f.dataset("keep").unwrap();
let mut pages: Vec<u64> = ds
.chunks()
.unwrap()
.iter()
.filter(|c| c.storage_size > 0)
.flat_map(|c| (c.address / page)..=((c.address + c.storage_size - 1) / page))
.collect();
pages.sort_unstable();
pages.dedup();
pages
};
assert!(!live_raw_pages.is_empty(), "expected live raw pages");
let s = WriteEngine::open_with_locking(&path, FileLocking::Enabled).unwrap();
let pg = s.paged.as_ref().expect("a paged file installs paged state");
for (addr, len) in pg.meta.sections() {
for p in (addr / page)..=((addr + len - 1) / page) {
assert!(
!live_raw_pages.contains(&p),
"metadata free section ({addr}, {len}) sits in page {p}, which still \
holds live raw chunk data"
);
}
}
let reclaimed: u64 = pg.all_sections().iter().map(|&(_, l)| l).sum();
assert!(reclaimed > 0, "the delete reclaimed nothing");
}
#[test]
fn failed_paged_commit_leaves_the_free_lists_untouched() {
use crate::writer::FileBuilder;
use tempfile::tempdir;
let dir = tempdir().unwrap();
let path = dir.path().join("paged_failed_commit.h5");
let mut b = FileBuilder::new();
b.create_dataset("keep")
.with_i32_data(&(0..200).collect::<Vec<i32>>())
.with_shape(&[200]);
b.create_dataset("drop")
.with_i32_data(&(0..200).collect::<Vec<i32>>())
.with_shape(&[200]);
b.with_file_space_strategy(FileSpaceStrategy::Page, true, 0)
.with_file_space_page_size(4096);
b.write(&path).unwrap();
let mut s = WriteEngine::open_with_locking(&path, FileLocking::Enabled).unwrap();
let before = s.space_accounting().reusable_free_space;
let good_ext = s.superblock.superblock_extension_address;
s.superblock.superblock_extension_address = Some(0);
s.delete("/drop");
assert!(
s.commit().is_err(),
"a commit with an unreadable extension must fail"
);
assert_eq!(
s.space_accounting().reusable_free_space,
before,
"a failed commit must not record still-live regions as free"
);
s.superblock.superblock_extension_address = good_ext;
s.delete("/drop");
s.commit()
.expect("the session is usable after a failed commit");
drop(s);
let f = crate::reader::File::open(&path).unwrap();
let kept = f.dataset("keep").unwrap().read_i32().unwrap();
assert_eq!(kept, (0..200).collect::<Vec<i32>>(), "keep survives intact");
let freed: u64 = f.persisted_free_space().iter().map(|&(_, l)| l).sum();
let live_end = f.file_size();
assert!(
freed < live_end,
"the recorded free space cannot cover the whole file"
);
}
#[test]
fn append_inplace_crash_consistency_paged_prefix() {
use crate::reader::File as PureFile;
use tempfile::tempdir;
let dir = tempdir().unwrap();
let base = dir.path().join("base.h5");
let (start, target) = (131_000i32, 132_000i32);
build_unit_chunked(&base, start);
for max_phase in 1u8..=4 {
let p = dir.path().join(format!("crash_paged_{max_phase}.h5"));
append_stopped_at(&base, &p, start..target, max_phase);
let expected_len = if max_phase == 4 { target } else { start };
let f = PureFile::from_bytes(std::fs::read(&p).unwrap()).unwrap();
assert_eq!(
f.dataset("d").unwrap().read_i32().unwrap(),
(0..expected_len).collect::<Vec<_>>(),
"inconsistent paged view after crash at phase {max_phase}"
);
}
}
#[test]
#[cfg(not(target_pointer_width = "32"))]
fn append_inplace_crash_consistency_c_library_reads_prefix() {
use tempfile::tempdir;
for (n, chunk, add) in [(6i32, 4u64, 5i32), (8, 2, 6), (50, 1, 200)] {
let dir = tempdir().unwrap();
let base = dir.path().join("base.h5");
{
use crate::writer::FileBuilder;
let mut b = FileBuilder::new();
b.create_dataset("d")
.with_i32_data(&(0..n).collect::<Vec<i32>>())
.with_shape(&[n as u64])
.with_maxshape(&[u64::MAX])
.with_chunks(&[chunk]);
b.write(&base).unwrap();
}
for max_phase in 1u8..=4 {
let p = dir
.path()
.join(format!("crash_c_{n}_{chunk}_{max_phase}.h5"));
append_stopped_at(&base, &p, n..n + add, max_phase);
let expected_len = if max_phase == 4 { n + add } else { n };
let f = hdf5::File::open(&p).unwrap();
assert_eq!(
f.dataset("d").unwrap().read_raw::<i32>().unwrap(),
(0..expected_len).collect::<Vec<_>>(),
"C library saw an inconsistent view after crash at phase {max_phase} \
(n={n}, chunk={chunk})"
);
f.close().unwrap();
}
}
}
#[test]
#[cfg(not(target_pointer_width = "32"))]
fn append_inplace_recover_and_reappend_after_phase3_crash() {
use crate::reader::File as PureFile;
use tempfile::tempdir;
let dir = tempdir().unwrap();
let path = dir.path().join("phase3_recover.h5");
let n = 50i32;
build_unit_chunked(&path, n);
{
let mut s = WriteEngine::open_with_locking(&path, FileLocking::Enabled).unwrap();
s.append_inplace_i32_phased("d", &(1000..1200).collect::<Vec<_>>(), 3)
.unwrap();
}
let committed: Vec<i32> = (0..n).collect();
let pf = PureFile::from_bytes(std::fs::read(&path).unwrap()).unwrap();
assert_eq!(
pf.dataset("d").unwrap().read_i32().unwrap(),
committed,
"phase-3 crash exposed uncommitted data to the pure reader"
);
{
let f = hdf5::File::open(&path).unwrap();
assert_eq!(
f.dataset("d").unwrap().read_raw::<i32>().unwrap(),
committed,
"phase-3 crash exposed uncommitted data to the C library"
);
f.close().unwrap();
}
{
let mut s = WriteEngine::open_with_locking(&path, FileLocking::Enabled).unwrap();
s.append_inplace_i32_phased("d", &(n..150).collect::<Vec<_>>(), 4)
.unwrap();
}
let expected: Vec<i32> = (0..150).collect();
let pf = PureFile::from_bytes(std::fs::read(&path).unwrap()).unwrap();
assert_eq!(
pf.dataset("d").unwrap().read_i32().unwrap(),
expected,
"recovery did not roll forward correctly (pure reader)"
);
let f = hdf5::File::open(&path).unwrap();
assert_eq!(
f.dataset("d").unwrap().read_raw::<i32>().unwrap(),
expected,
"recovery did not roll forward correctly (C library)"
);
f.close().unwrap();
}
#[test]
fn raw_appendable_recurses_into_aggregates() {
use crate::datatype::{CompoundMember, DatatypeByteOrder};
let f64_with = |byte_order| Datatype::FloatingPoint {
size: 8,
byte_order,
bit_offset: 0,
bit_precision: 64,
exponent_location: 52,
exponent_size: 11,
mantissa_location: 0,
mantissa_size: 52,
exponent_bias: 1023,
};
let le_f64 = f64_with(DatatypeByteOrder::LittleEndian);
let be_f64 = f64_with(DatatypeByteOrder::BigEndian);
assert!(datatype_is_raw_appendable(&le_f64));
assert!(!datatype_is_raw_appendable(&be_f64));
let be_member = Datatype::Compound {
size: 8,
members: vec![CompoundMember {
name: "x".into(),
byte_offset: 0,
datatype: be_f64.clone(),
}],
};
assert!(!datatype_is_raw_appendable(&be_member));
let le_member = Datatype::Compound {
size: 8,
members: vec![CompoundMember {
name: "x".into(),
byte_offset: 0,
datatype: le_f64.clone(),
}],
};
assert!(datatype_is_raw_appendable(&le_member));
assert!(!datatype_is_raw_appendable(&Datatype::Array {
base_type: Box::new(be_f64.clone()),
dimensions: vec![4],
}));
assert!(!datatype_is_raw_appendable(&Datatype::VariableLength {
is_string: false,
padding: None,
charset: None,
base_type: Box::new(le_f64.clone()),
}));
assert!(!datatype_is_raw_appendable(&Datatype::Reference {
size: 8,
ref_type: crate::datatype::ReferenceType::Object,
}));
}
#[test]
fn fresh_group_region_pairs_link_info_with_group_info() {
let types = region_types(&fresh_group_region());
assert_eq!(types, vec![MessageType::LinkInfo, MessageType::GroupInfo]);
}
#[test]
fn ensure_group_info_appends_when_missing() {
let li_body = {
let mut b = vec![0u8, 0];
b.extend_from_slice(&u64::MAX.to_le_bytes());
b.extend_from_slice(&u64::MAX.to_le_bytes());
b
};
let mut region = region_message(MessageType::LinkInfo, &li_body);
ensure_group_info(&mut region).unwrap();
assert_eq!(
region_types(®ion),
vec![MessageType::LinkInfo, MessageType::GroupInfo]
);
let mut p = 0;
while let Some((mt, body, end)) = next_message(®ion, p).unwrap() {
if mt == MessageType::GroupInfo {
assert_eq!(®ion[body..end], &GROUP_INFO_BODY);
}
p = end;
}
}
#[test]
fn ensure_group_info_is_idempotent() {
let mut region = fresh_group_region();
let before = region.clone();
ensure_group_info(&mut region).unwrap();
assert_eq!(region, before);
}
#[test]
fn reject_foreign_addresses_refuses_any_shared_message() {
let mut shared = region_message(MessageType::Dataspace, &[0u8; 8]);
shared[3] = MSG_FLAG_SHARED; let err = reject_foreign_addresses(&shared).unwrap_err();
assert!(err.to_string().contains("shared"), "got: {err}");
let plain = region_message(MessageType::Dataspace, &[0u8; 8]);
reject_foreign_addresses(&plain).unwrap();
}
fn compact_layout_body(version: u8, data: &[u8]) -> Vec<u8> {
let mut b = vec![version, 0];
b.extend_from_slice(&(data.len() as u16).to_le_bytes());
b.extend_from_slice(data);
b
}
#[test]
fn rebuild_compact_layout_replaces_inline_data_only() {
let mut region = region_message(MessageType::Dataspace, &[0xAB; 8]);
region.extend_from_slice(®ion_message(
MessageType::DataLayout,
&compact_layout_body(3, &[1, 2, 3, 4]),
));
region.extend_from_slice(®ion_message(MessageType::Attribute, &[0xCD; 5]));
let out = rebuild_compact_layout_region(®ion, &[9, 8, 7, 6]).unwrap();
assert_eq!(
region_types(&out),
vec![
MessageType::Dataspace,
MessageType::DataLayout,
MessageType::Attribute,
]
);
let mut p = 0;
while let Some((mt, body, end)) = next_message(&out, p).unwrap() {
match mt {
MessageType::Dataspace => assert_eq!(&out[body..end], &[0xAB; 8]),
MessageType::DataLayout => {
assert_eq!(out[body], 3, "version preserved");
assert_eq!(out[body + 1], 0, "still compact");
let size = u16::from_le_bytes([out[body + 2], out[body + 3]]) as usize;
assert_eq!(size, 4);
assert_eq!(&out[body + 4..body + 4 + size], &[9, 8, 7, 6]);
}
MessageType::Attribute => assert_eq!(&out[body..end], &[0xCD; 5]),
other => panic!("unexpected message {other:?}"),
}
p = end;
}
}
#[test]
fn rebuild_compact_layout_refuses_non_compact() {
let mut region = region_message(MessageType::DataLayout, &{
let mut b = vec![3u8, 1]; b.extend_from_slice(&0u64.to_le_bytes());
b.extend_from_slice(&0u64.to_le_bytes());
b
});
region.extend_from_slice(®ion_message(MessageType::Dataspace, &[0; 8]));
let err = rebuild_compact_layout_region(®ion, &[1, 2]).unwrap_err();
assert!(err.to_string().contains("non-compact"), "got: {err}");
}
#[test]
fn a_refused_open_never_builds_the_image() {
use crate::writer::FileBuilder;
let dir = tempfile::tempdir().unwrap();
let flagged = dir.path().join("flagged.h5");
let mut b = FileBuilder::new();
b.create_dataset("d").with_i32_data(&[1, 2, 3]);
b.write(&flagged).unwrap();
let ancient = dir.path().join("ancient.h5");
std::fs::copy(&flagged, &ancient).unwrap();
for (path, version, flags) in [(&flagged, 3, SWMR_WRITE_FLAGS), (&ancient, 9, 0)] {
let mut data = std::fs::read(path).unwrap();
let off = signature::find_signature(&data).unwrap();
let mut sb = Superblock::parse(&data, off).unwrap();
sb.version = version;
sb.consistency_flags = flags;
let bytes = sb.serialize();
data[off..off + bytes.len()].copy_from_slice(&bytes);
std::fs::write(path, &data).unwrap();
}
for path in [&flagged, &ancient] {
let built = std::cell::Cell::new(false);
let err = match WriteEngine::open_imaged(path, Some(FileLocking::Enabled), |h, len| {
built.set(true);
Ok(Box::new(HandleImage::new(
h,
len,
MetadataCacheConfig::disabled(),
)))
}) {
Err(e) => e,
Ok(_) => panic!("{} must be refused", path.display()),
};
assert!(
!built.get(),
"{} was refused with {err:?}, but the image was built first",
path.display()
);
}
}
#[test]
fn a_stale_consistency_flag_is_refused_then_cleared_by_a_commit() {
use crate::writer::FileBuilder;
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("stale_flag.h5");
let mut b = FileBuilder::new();
b.create_dataset("d").with_i32_data(&[1, 2, 3]);
b.write(&path).unwrap();
{
let mut data = std::fs::read(&path).unwrap();
let off = signature::find_signature(&data).unwrap();
let mut sb = Superblock::parse(&data, off).unwrap();
assert!(
sb.version >= 2,
"FileBuilder should emit a v2/v3 superblock"
);
sb.consistency_flags = 0x05;
let bytes = sb.serialize();
data[off..off + bytes.len()].copy_from_slice(&bytes);
std::fs::write(&path, &data).unwrap();
assert_eq!(
Superblock::parse(&data, off).unwrap().consistency_flags,
0x05
);
}
match WriteEngine::open_with_locking(&path, FileLocking::Enabled) {
Err(Error::FileMarkedInUse(_)) => {}
Err(e) => panic!("expected the flag refusal, got {e:?}"),
Ok(_) => panic!("a flagged file must not be edited in place"),
}
{
let mut data = std::fs::read(&path).unwrap();
let off = signature::find_signature(&data).unwrap();
let mut sb = Superblock::parse(&data, off).unwrap();
sb.version = 2;
sb.consistency_flags = crate::file_lock::WRITE_ACCESS;
let bytes = sb.serialize();
data[off..off + bytes.len()].copy_from_slice(&bytes);
std::fs::write(&path, &data).unwrap();
}
{
let mut s = WriteEngine::open_with_locking(&path, FileLocking::Enabled)
.expect("the gate skips a v2 superblock, so this opens");
let mut b = DatasetBuilder::new("e");
b.with_i32_data(&[4, 5]);
s.stage_created_dataset("e", b);
s.commit().unwrap();
}
let data = std::fs::read(&path).unwrap();
let off = signature::find_signature(&data).unwrap();
assert_eq!(
Superblock::parse(&data, off).unwrap().consistency_flags,
0,
"commit must clear the stale consistency flag"
);
}
#[test]
fn add_vlen_string_dataset_with_null_elements_via_edit_session() {
use crate::type_builders::VlStringElement;
use crate::writer::FileBuilder;
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("vlen_null.h5");
let mut b = FileBuilder::new();
b.create_dataset("seed").with_i32_data(&[0]);
b.write(&path).unwrap();
let datatype =
crate::type_builders::make_vlen_string_type(crate::datatype::CharacterSet::Utf8);
let elements = vec![
VlStringElement::Bytes(b"alpha".to_vec()),
VlStringElement::Null,
VlStringElement::Bytes(b"gamma".to_vec()),
];
{
let mut s = WriteEngine::open_with_locking(&path, FileLocking::Enabled).unwrap();
let mut b = DatasetBuilder::new("labels");
b.with_vlen_string_elements(datatype, &elements).unwrap();
s.stage_created_dataset("labels", b);
s.commit().unwrap();
}
let file = crate::reader::File::open(&path).unwrap();
let ds = file.dataset("labels").unwrap();
assert_eq!(
ds.read_string().unwrap(),
vec!["alpha".to_string(), String::new(), "gamma".to_string()]
);
}
#[test]
fn edit_session_root_group_base_address_overflow_is_rejected() {
use crate::writer::FileBuilder;
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("edit_root_overflow.h5");
const UB: u64 = 512;
let mut b = FileBuilder::new();
b.with_userblock(UB);
b.create_dataset("d").with_i32_data(&[1, 2, 3]);
b.write(&path).unwrap();
let mut data = std::fs::read(&path).unwrap();
let off = signature::find_signature(&data).unwrap();
let mut sb = Superblock::parse(&data, off).unwrap();
assert_eq!(sb.base_address, UB, "userblock file must have base == UB");
sb.root_group_address = u64::MAX;
let bytes = sb.serialize();
data[off..off + bytes.len()].copy_from_slice(&bytes);
std::fs::write(&path, &data).unwrap();
let err = WriteEngine::open_with_locking(&path, FileLocking::Enabled)
.err()
.expect("open must fail");
match err {
Error::Format(FormatError::OffsetOverflow { offset, length }) => {
assert_eq!(offset, u64::MAX);
assert_eq!(length, UB);
}
other => panic!("expected root-group address overflow, got {other:?}"),
}
}
use tempfile::tempdir;
fn build_appendable(path: &Path, n: i32, chunk: u64) {
let data: Vec<i32> = (0..n).collect();
let mut b = crate::writer::FileBuilder::new();
b.create_dataset("d")
.with_i32_data(&data)
.with_shape(&[n as u64])
.with_maxshape(&[u64::MAX])
.with_chunks(&[chunk]);
b.write(path).unwrap();
}
fn open_bounded_session(path: &Path) -> WriteEngine {
WriteEngine::open_rw_with_strategy(
path,
crate::source::MetadataCacheConfig::disabled(),
FileLocking::Enabled,
MemoryStrategy::Bounded,
)
.unwrap()
}
fn dataset_addr(engine: &WriteEngine) -> u64 {
crate::group_v2::resolve_path_any_from_source(&engine.image(), engine.superblock(), "d")
.unwrap()
}
#[test]
fn bounded_append_crash_consistency_partial_tail_prefix() {
let dir = tempdir().unwrap();
for (case, (n, chunk, add)) in [(0usize, (6i32, 4u64, 5i32)), (1, (8, 2, 6))] {
let base = dir.path().join(std::format!("base_{case}.h5"));
build_appendable(&base, n, chunk);
for max_phase in 1u8..=4 {
let p = dir.path().join(std::format!("crash_{case}_{max_phase}.h5"));
std::fs::copy(&base, &p).unwrap();
{
let mut engine = open_bounded_session(&p);
let addr = dataset_addr(&engine);
let mut b = AppendBuilder::new();
b.append_i32(&(n..n + add).collect::<Vec<_>>());
engine
.append_inplace_gathered(AppendTarget::Header(addr), &b, max_phase)
.unwrap();
}
let expected_len = if max_phase == 4 { n + add } else { n };
let got = crate::File::open(&p)
.unwrap()
.dataset("d")
.unwrap()
.read_i32()
.unwrap();
assert_eq!(
got,
(0..expected_len).collect::<Vec<_>>(),
"case {case} phase {max_phase}"
);
}
}
}
#[test]
fn bounded_multi_batch_append_commits_every_batch() {
let dir = tempdir().unwrap();
let p = dir.path().join("multibatch.h5");
build_appendable(&p, 5, 512);
let total = 700_000i32;
{
let mut engine = open_bounded_session(&p);
let addr = dataset_addr(&engine);
let mut b = AppendBuilder::new();
b.append_i32(&(5..total).collect::<Vec<_>>());
engine
.append_inplace_gathered(AppendTarget::Header(addr), &b, 4)
.unwrap();
}
let got = crate::File::open(&p)
.unwrap()
.dataset("d")
.unwrap()
.read_i32()
.unwrap();
assert_eq!(got.len(), total as usize);
assert!(got.iter().enumerate().all(|(i, &v)| v == i as i32));
}
#[test]
fn only_a_bounded_session_batches_a_large_append() {
let dir = tempdir().unwrap();
let p = dir.path().join("batching.h5");
build_appendable(&p, 8, 4);
let bounded_batch = {
let mut engine = open_bounded_session(&p);
engine
.append_geometry(AppendTarget::Path("d"))
.unwrap()
.full_batch_elems
};
let mirror_batch = {
let mut engine = WriteEngine::open_with_locking(&p, FileLocking::Enabled).unwrap();
engine
.append_geometry(AppendTarget::Path("d"))
.unwrap()
.full_batch_elems
};
assert_eq!(
mirror_batch,
u64::MAX,
"a mirror session must take the whole append as one crash-atomic batch"
);
assert!(
bounded_batch < u64::MAX,
"a bounded session must cap a batch, got {bounded_batch}"
);
assert_eq!(
bounded_batch % 4,
0,
"a batch must be a whole number of chunks"
);
}
#[test]
fn bounded_persist_append_without_finalize_is_readable() {
let dir = tempdir().unwrap();
let p = dir.path().join("persist_crash.h5");
let mut b = crate::writer::FileBuilder::new();
b.with_file_space_strategy(crate::FileSpaceStrategy::FsmAggr, true, 1);
b.create_dataset("d")
.with_i32_data(&(0..6).collect::<Vec<i32>>())
.with_shape(&[6])
.with_maxshape(&[u64::MAX])
.with_chunks(&[4]);
b.write(&p).unwrap();
{
let mut engine = open_bounded_session(&p);
assert!(engine.persist.is_some(), "persist state is armed at open");
let addr = dataset_addr(&engine);
let mut ab = AppendBuilder::new();
ab.append_i32(&(6..20).collect::<Vec<_>>());
engine
.append_inplace_gathered(AppendTarget::Header(addr), &ab, 4)
.unwrap();
}
let got = crate::File::open(&p)
.unwrap()
.dataset("d")
.unwrap()
.read_i32()
.unwrap();
assert_eq!(got, (0..20).collect::<Vec<_>>());
}
#[test]
fn bounded_paged_reopen_after_crash_realigns_and_stays_readable() {
let dir = tempdir().unwrap();
let p = dir.path().join("paged_crash.h5");
let mut b = crate::writer::FileBuilder::new();
b.with_file_space_strategy(crate::FileSpaceStrategy::Page, true, 0)
.with_file_space_page_size(4096);
b.create_dataset("d")
.with_i32_data(&(0..64).collect::<Vec<i32>>())
.with_shape(&[64])
.with_maxshape(&[u64::MAX])
.with_chunks(&[64]);
b.write(&p).unwrap();
{
let mut engine = open_bounded_session(&p);
let addr = dataset_addr(&engine);
let mut ab = AppendBuilder::new();
ab.append_i32(&(64..2000).collect::<Vec<_>>());
engine
.append_inplace_gathered(AppendTarget::Header(addr), &ab, 4)
.unwrap();
}
assert_ne!(
std::fs::metadata(&p).unwrap().len() % 4096,
0,
"a crashed (un-finalized) paged session leaves the file non-page-aligned"
);
{
let mut engine = open_bounded_session(&p);
let addr = dataset_addr(&engine);
let mut ab = AppendBuilder::new();
ab.append_i32(&(2000..2500).collect::<Vec<_>>());
engine
.append_inplace_gathered(AppendTarget::Header(addr), &ab, 4)
.unwrap();
engine.finalize_persist().unwrap();
engine.sync().unwrap();
}
assert_eq!(
std::fs::metadata(&p).unwrap().len() % 4096,
0,
"reopen + append + finalize re-aligns the paged file"
);
let got = crate::File::open(&p)
.unwrap()
.dataset("d")
.unwrap()
.read_i32()
.unwrap();
assert_eq!(got, (0..2500).collect::<Vec<_>>());
}
#[test]
fn a_bounded_commit_reads_far_less_than_the_file() {
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
let dir = tempdir().unwrap();
let p = dir.path().join("bulk.h5");
let rows = 2_000_000i32;
build_appendable(&p, rows, 8192);
let file_len = std::fs::metadata(&p).unwrap().len();
assert!(file_len > 4 << 20, "file is only {file_len} bytes");
let read_bytes = Arc::new(AtomicU64::new(0));
{
let mut engine =
WriteEngine::open_bounded_counting(&p, Arc::clone(&read_bytes)).unwrap();
engine.create_group("g");
engine.commit().unwrap();
}
let read = read_bytes.load(Ordering::Relaxed);
assert!(
read > 0,
"the commit read nothing, so the test proves nothing"
);
assert!(
read < 64 << 10,
"a bounded commit read {read} bytes of a {file_len}-byte file"
);
}
#[test]
fn a_commit_after_an_append_pads_the_raw_page_the_append_left() {
const PAGE: u64 = 4096;
let dir = tempdir().unwrap();
let p = dir.path().join("paged_interleave.h5");
let mut b = crate::writer::FileBuilder::new();
b.with_file_space_strategy(crate::FileSpaceStrategy::Page, true, 0)
.with_file_space_page_size(PAGE);
b.create_dataset("d")
.with_i32_data(&(0..64).collect::<Vec<i32>>())
.with_shape(&[64])
.with_maxshape(&[u64::MAX])
.with_chunks(&[64]);
b.write(&p).unwrap();
let mut engine = WriteEngine::open_with_locking(&p, FileLocking::Enabled).unwrap();
let mut ab = AppendBuilder::new();
ab.append_i32(&(64..2000).collect::<Vec<_>>());
engine
.append_inplace_gathered(AppendTarget::Path("d"), &ab, 4)
.unwrap();
assert_eq!(
engine.paged.as_ref().unwrap().last,
Some(PageType::Raw),
"the append must record that the tail page now holds raw data"
);
assert_ne!(
engine.image.len() % PAGE,
0,
"the append must leave a partially-filled page for the commit to pad"
);
engine.create_group("g");
engine.commit().unwrap();
let pg = engine.paged.as_ref().expect("the file is paged");
let raw_free = pg.raw_small.sections();
assert!(
!raw_free.is_empty(),
"the commit packed metadata into the raw page the append left open"
);
for (addr, len) in raw_free {
assert_eq!(
(addr + len) % PAGE,
0,
"padding {addr}+{len} does not reach a page boundary"
);
}
drop(engine);
assert_eq!(
crate::File::open(&p)
.unwrap()
.dataset("d")
.unwrap()
.read_i32()
.unwrap(),
(0..2000).collect::<Vec<_>>()
);
}
#[test]
fn an_inplace_append_to_a_paged_file_allocates_only_raw_pages() {
const PAGE: u64 = 4096;
let dir = tempdir().unwrap();
let p = dir.path().join("paged_raw.h5");
let mut b = crate::writer::FileBuilder::new();
b.with_file_space_strategy(crate::FileSpaceStrategy::Page, true, 0)
.with_file_space_page_size(PAGE);
b.create_dataset("d")
.with_i32_data(&(0..64).collect::<Vec<i32>>())
.with_shape(&[64])
.with_maxshape(&[u64::MAX])
.with_chunks(&[64]);
b.write(&p).unwrap();
let mut engine = WriteEngine::open_with_locking(&p, FileLocking::Enabled).unwrap();
let before = engine.image().len();
for range in [64..2000, 2000..4000] {
let mut ab = AppendBuilder::new();
ab.append_i32(&range.collect::<Vec<_>>());
engine
.append_inplace_gathered(AppendTarget::Path("d"), &ab, 4)
.unwrap();
}
let pg = engine.paged.as_ref().expect("the file is paged");
assert_eq!(
pg.last,
Some(PageType::Raw),
"the append left the tail page holding something other than raw data"
);
assert!(
pg.meta_pad.is_empty() && pg.raw_pad.is_empty(),
"an in-place append switched page type: meta_pad={:?} raw_pad={:?}",
pg.meta_pad,
pg.raw_pad
);
let addr = crate::group_v2::resolve_path_any_from_source(
&engine.image(),
engine.superblock(),
"d",
)
.unwrap();
let spans = engine
.chunked_storage_spans(addr.to_usize().unwrap())
.expect("a chunked dataset has reclaimable spans");
let fresh = spans.iter().filter(|&&(a, _, _)| a >= before).count();
assert!(
fresh > 0,
"the append allocated nothing above {before}, so the assertion above proves nothing"
);
assert!(
spans.iter().all(|&(_, _, ty)| ty == PageType::Raw),
"the reclaim tags every chunked span raw; a metadata tag here would need \
the placement rule above to change with it"
);
}
}