use std::collections::HashMap;
use std::io::{Read, Seek, SeekFrom, Write};
use std::path::Path;
use crate::chunk_index_inplace::{Store, apply_ea_append, plan_ea_append};
use crate::edit::{
AppendBuilder, LocatedState, as_inplace_error, build_v2_object_header, locate_dataset_state,
read_single_chunk_ext_region, rewrite_extension_region_bytes, validate_gathered_append,
};
use crate::error::{Error, FormatError};
use crate::file_lock::{self, FileLocking};
use crate::file_space_info::{FileSpaceInfo, FileSpaceStrategy, NUM_FILE_FSM_MANAGERS};
use crate::free_space::FreeList;
use crate::free_space_manager::{
FreeSection, SECT_CLASS_LARGE, SECT_CLASS_SIMPLE, SECT_CLASS_SMALL, fshd_len, fsse_len,
read_persisted_sections_source, serialize_file_fsm,
};
use crate::message_type::MessageType;
use crate::object_header::ObjectHeader;
use crate::signature;
use crate::source::{MetadataCacheConfig, MetadataReadCache, Source};
use crate::superblock::Superblock;
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,
}
fn read_at_handle(
handle: &std::fs::File,
len: u64,
offset: u64,
buf: &mut [u8],
) -> Result<(), FormatError> {
let end = offset
.checked_add(buf.len() as u64)
.ok_or(FormatError::OffsetOverflow {
offset,
length: buf.len() as u64,
})?;
if end > len {
return Err(FormatError::UnexpectedEof {
expected: end.to_usize().unwrap_or(usize::MAX),
available: len.to_usize().unwrap_or(usize::MAX),
});
}
let mut h = handle;
h.seek(SeekFrom::Start(offset))
.map_err(|e| FormatError::Source(std::format!("{e}")))?;
h.read_exact(buf)
.map_err(|e| FormatError::Source(std::format!("{e}")))?;
Ok(())
}
use crate::convert::TryToUsize;
struct RawSource<'a> {
handle: &'a std::fs::File,
len: u64,
}
impl Source for RawSource<'_> {
fn len(&self) -> u64 {
self.len
}
fn read_at(&self, offset: u64, buf: &mut [u8]) -> Result<(), FormatError> {
read_at_handle(self.handle, self.len, offset, buf)
}
}
pub(crate) struct BoundedStore {
handle: std::fs::File,
len: u64,
sb_sig_off: u64,
superblock: Superblock,
metadata_cache: Option<(MetadataCacheConfig, std::sync::Mutex<MetadataReadCache>)>,
paged: Option<PagedAppend>,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
enum PageType {
Meta,
Raw,
}
struct PagedAppend {
page_size: u64,
last: Option<PageType>,
meta_pad: Vec<(u64, u64)>,
raw_pad: Vec<(u64, u64)>,
}
impl BoundedStore {
pub(crate) fn superblock(&self) -> &Superblock {
&self.superblock
}
fn append_typed(&mut self, bytes: &[u8], ty: PageType) -> Result<u64, Error> {
let pad = match &self.paged {
Some(pg) if self.len % pg.page_size != 0 => {
let pad_len = pg.page_size - self.len % pg.page_size;
match pg.last {
Some(prev) if prev != ty => Some((Some(prev), pad_len)),
Some(_) => None, None => Some((None, pad_len)),
}
}
_ => None,
};
if let Some((prev, pad_len)) = pad {
let pad_at = self.len;
self.append_at_eof(&vec![0u8; pad_len.to_usize()?])?;
if let Some(pg) = self.paged.as_mut() {
match prev {
Some(PageType::Meta) => pg.meta_pad.push((pad_at, pad_len)),
Some(PageType::Raw) => pg.raw_pad.push((pad_at, pad_len)),
None => {} }
}
}
let addr = self.append_at_eof(bytes)?;
if let Some(pg) = self.paged.as_mut() {
pg.last = Some(ty);
}
Ok(addr)
}
fn pad_to_page(&mut self) -> Result<u64, Error> {
let pad = match &self.paged {
Some(pg) if self.len % pg.page_size != 0 => {
Some((pg.last, pg.page_size - self.len % pg.page_size))
}
_ => None,
};
if let Some((last, pad_len)) = pad {
let pad_at = self.len;
self.append_at_eof(&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(self.len)
}
fn append_at_eof(&mut self, bytes: &[u8]) -> Result<u64, Error> {
let addr = self.len;
self.handle.seek(SeekFrom::Start(addr)).map_err(Error::Io)?;
self.handle.write_all(bytes).map_err(Error::Io)?;
self.len += bytes.len() as u64;
Ok(addr)
}
fn write_at_raw(&mut self, offset: u64, bytes: &[u8]) -> Result<(), Error> {
let end = offset
.checked_add(bytes.len() as u64)
.filter(|&e| e <= self.len)
.ok_or(Error::Format(FormatError::UnexpectedEof {
expected: offset.to_usize().unwrap_or(usize::MAX),
available: self.len.to_usize().unwrap_or(usize::MAX),
}))?;
debug_assert!(end <= self.len);
if let Some((_, cache)) = &self.metadata_cache {
cache
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.invalidate_overlapping(offset, bytes.len());
}
self.handle
.seek(SeekFrom::Start(offset))
.map_err(Error::Io)?;
self.handle.write_all(bytes).map_err(Error::Io)?;
Ok(())
}
fn write_persist_tail(
&mut self,
old_ext_region: &[u8],
strategy: FileSpaceStrategy,
threshold: u64,
page_size: u64,
sections: &[FreeSection],
) -> Result<Vec<(u64, u64)>, Error> {
let os = self.superblock.offset_size;
let new_root = self.superblock.root_group_address;
let ext_addr = self.len;
let placeholder =
FileSpaceInfo::persistent_single_manager(strategy, threshold, page_size, 0, 0);
let ext_len = build_v2_object_header(&rewrite_extension_region_bytes(
old_ext_region,
&placeholder,
)?)
.len() as u64;
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(&rewrite_extension_region_bytes(old_ext_region, &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(&rewrite_extension_region_bytes(old_ext_region, &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(sections, 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_bytes(&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_bytes(&fshd)?;
debug_assert_eq!(wf, fshd_addr);
new_old_blocks.push((fshd_addr, fshd.len() as u64));
let ws = self.append_bytes(&fsse)?;
new_old_blocks.push((ws, fsse.len() as u64));
}
Store::sync(self)?;
self.repoint_persist_superblock(new_root, final_eof, ext_addr)?;
Store::sync(self)?;
Ok(new_old_blocks)
}
#[allow(clippy::too_many_arguments)]
fn write_persist_tail_paged(
&mut self,
old_ext_region: &[u8],
strategy: FileSpaceStrategy,
threshold: u64,
page_size: u64,
meta: &[FreeSection],
raw_small: &[FreeSection],
raw_large: &[FreeSection],
) -> Result<Vec<(u64, u64)>, Error> {
let os = self.superblock.offset_size;
let new_root = self.superblock.root_group_address;
let ext_addr = self.len;
debug_assert_eq!(
ext_addr % page_size,
0,
"extension begins on a page boundary"
);
let placeholder = FileSpaceInfo::persistent_managers(
strategy,
threshold,
page_size,
[u64::MAX; NUM_FILE_FSM_MANAGERS],
0,
);
let ext_len = build_v2_object_header(&rewrite_extension_region_bytes(
old_ext_region,
&placeholder,
)?)
.len() as u64;
let mut slot0 = Vec::new();
let mut slot2 = Vec::new();
let mut slot6 = Vec::new();
for s in split_at_pages(meta, page_size) {
if s.size < page_size {
slot0.push(s)
} else {
slot6.push(s)
}
}
for s in split_at_pages(raw_small, page_size) {
if s.size < page_size {
slot2.push(s)
} else {
slot6.push(s)
}
}
for s in split_at_pages(raw_large, page_size) {
slot6.push(s);
}
slot6.sort_by_key(|s| s.addr);
let managers: [(usize, u8, &[FreeSection]); 3] = [
(0, SECT_CLASS_SMALL, &slot0),
(2, SECT_CLASS_SMALL, &slot2),
(6, SECT_CLASS_LARGE, &slot6),
];
if managers.iter().all(|(_, _, s)| s.is_empty()) {
let info = FileSpaceInfo::persistent_empty(strategy, threshold, page_size);
let ext_oh =
build_v2_object_header(&rewrite_extension_region_bytes(old_ext_region, &info)?);
let written_ext = self.append_bytes(&ext_oh)?;
debug_assert_eq!(written_ext, ext_addr);
let final_eof = align_up(ext_addr + ext_oh.len() as u64, page_size);
self.pad_zeros_to(final_eof)?;
Store::sync(self)?;
self.repoint_persist_superblock(new_root, final_eof, ext_addr)?;
Store::sync(self)?;
return Ok(vec![(ext_addr, ext_oh.len() as u64)]);
}
let mut slots = [u64::MAX; NUM_FILE_FSM_MANAGERS];
let mut blocks: Vec<(usize, u64, u64, u8, &[FreeSection])> = Vec::new(); let mut cursor = ext_addr + ext_len;
for &(slot, class, sections) in &managers {
if sections.is_empty() {
continue;
}
let fshd_addr = cursor;
let fsse_addr = fshd_addr + fshd_len(os);
let section_sizes: Vec<u64> = sections.iter().map(|s| s.size).collect();
cursor = fsse_addr + fsse_len(§ion_sizes, os);
slots[slot] = fshd_addr;
blocks.push((slot, fshd_addr, fsse_addr, class, sections));
}
let end_of_managers = cursor;
let final_eof = align_up(end_of_managers, page_size);
let eoa_pre_fsm = final_eof;
let info =
FileSpaceInfo::persistent_managers(strategy, threshold, page_size, slots, eoa_pre_fsm);
let ext_oh =
build_v2_object_header(&rewrite_extension_region_bytes(old_ext_region, &info)?);
debug_assert_eq!(
ext_oh.len() as u64,
ext_len,
"extension length must be stable across the placeholder and real messages"
);
let written_ext = self.append_bytes(&ext_oh)?;
debug_assert_eq!(written_ext, ext_addr);
let mut new_old_blocks = vec![(ext_addr, ext_oh.len() as u64)];
for &(_slot, fshd_addr, fsse_addr, class, sections) in &blocks {
let (fshd, fsse) = serialize_file_fsm(sections, fshd_addr, fsse_addr, os, class);
let wf = self.append_bytes(&fshd)?;
debug_assert_eq!(wf, fshd_addr);
new_old_blocks.push((fshd_addr, fshd.len() as u64));
let ws = self.append_bytes(&fsse)?;
debug_assert_eq!(ws, fsse_addr);
new_old_blocks.push((ws, fsse.len() as u64));
}
debug_assert_eq!(self.len, end_of_managers);
self.pad_zeros_to(final_eof)?;
Store::sync(self)?;
self.repoint_persist_superblock(new_root, final_eof, ext_addr)?;
Store::sync(self)?;
Ok(new_old_blocks)
}
fn pad_zeros_to(&mut self, target: u64) -> Result<(), Error> {
if target > self.len {
let pad = (target - self.len).to_usize()?;
self.append_at_eof(&vec![0u8; pad])?;
}
debug_assert_eq!(self.len, target);
Ok(())
}
fn repoint_persist_superblock(
&mut self,
new_root: u64,
new_eof: u64,
new_ext_addr: u64,
) -> Result<(), Error> {
self.superblock.root_group_address = new_root;
self.superblock.eof_address = new_eof;
self.superblock.superblock_extension_address = Some(new_ext_addr);
self.superblock.consistency_flags = 0;
self.len = new_eof;
let bytes = self.superblock.serialize();
self.write_at_raw(self.sb_sig_off, &bytes)
}
}
impl Source for BoundedStore {
fn len(&self) -> u64 {
self.len
}
fn read_at(&self, offset: u64, buf: &mut [u8]) -> Result<(), FormatError> {
read_at_handle(&self.handle, self.len, offset, buf)
}
fn read_metadata_at(&self, offset: u64, len: usize) -> Result<Vec<u8>, FormatError> {
let Some((config, cache)) = &self.metadata_cache else {
return self.read_exact_at(offset, len);
};
if len == 0 || len > config.max_entry_bytes() || len > config.max_bytes() {
return self.read_exact_at(offset, len);
}
if let Some(bytes) = cache
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(offset, len)
{
return Ok(bytes);
}
let bytes = self.read_exact_at(offset, len)?;
cache
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(offset, len, bytes.clone(), config.max_bytes());
Ok(bytes)
}
}
impl Store for BoundedStore {
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.append_at_eof(bytes)
}
fn append_raw(&mut self, bytes: &[u8]) -> Result<u64, Error> {
self.append_typed(bytes, PageType::Raw)
}
fn append_meta(&mut self, bytes: &[u8]) -> Result<u64, Error> {
self.append_typed(bytes, PageType::Meta)
}
fn write_at(&mut self, offset: u64, bytes: &[u8]) -> Result<(), Error> {
self.write_at_raw(offset, bytes)
}
fn patch_superblock_eof(&mut self) -> Result<(), Error> {
self.superblock.eof_address = self.len;
let bytes = self.superblock.serialize();
self.write_at_raw(self.sb_sig_off, &bytes)
}
fn sync(&mut self) -> Result<(), Error> {
self.handle.flush().map_err(Error::Io)?;
self.handle.sync_data().map_err(Error::Io)?;
Ok(())
}
}
struct BoundedPersistState {
strategy: FileSpaceStrategy,
threshold: u64,
page_size: u64,
free: PersistFree,
old_blocks: Vec<(u64, u64)>,
old_ext_addr: u64,
}
enum PersistFree {
Flat(FreeList),
Paged {
meta: FreeList,
raw_small: FreeList,
raw_large: FreeList,
},
}
impl PagedAppend {
fn new(page_size: u64) -> Self {
Self {
page_size,
last: None,
meta_pad: Vec::new(),
raw_pad: Vec::new(),
}
}
}
pub(crate) struct BoundedEngine {
store: BoundedStore,
located: HashMap<u64, LocatedState>,
persist: Option<BoundedPersistState>,
original_len: u64,
}
impl BoundedEngine {
pub(crate) fn open(path: &Path, metadata_cache: MetadataCacheConfig) -> Result<Self, Error> {
let handle = std::fs::OpenOptions::new()
.read(true)
.write(true)
.open(path)
.map_err(Error::Io)?;
file_lock::acquire_exclusive(&handle, FileLocking::Enabled, path)?;
let len = handle.metadata().map_err(Error::Io)?.len();
let raw = RawSource {
handle: &handle,
len,
};
let sb_sig_off = signature::find_signature_in(&raw)?;
let superblock = Superblock::parse_from_source(&raw, sb_sig_off)?;
if superblock.version < 2 {
return Err(Error::EditUnsupported(
"bounded read-write access requires a latest-format file (v2/v3 superblock); \
use File::open_rw",
));
}
if superblock.offset_size != 8 || superblock.length_size != 8 {
return Err(Error::EditUnsupported(
"bounded read-write access requires 8-byte offsets and lengths",
));
}
if superblock.base_address != 0 || sb_sig_off != 0 {
return Err(Error::EditUnsupported(
"bounded read-write access does not support a file with a userblock \
(non-zero base address); use File::open_rw",
));
}
let persist = load_bounded_persist(&raw, &superblock)?;
let paged = match &persist {
Some(ps) if ps.strategy == FileSpaceStrategy::Page => {
Some(PagedAppend::new(ps.page_size))
}
_ => None,
};
Ok(Self {
store: BoundedStore {
handle,
len,
sb_sig_off,
superblock,
metadata_cache: metadata_cache.is_enabled().then(|| {
(
metadata_cache,
std::sync::Mutex::new(MetadataReadCache::new()),
)
}),
paged,
},
located: HashMap::new(),
persist,
original_len: len,
})
}
pub(crate) fn store(&self) -> &BoundedStore {
&self.store
}
pub(crate) fn sync(&mut self) -> Result<(), Error> {
Store::sync(&mut self.store)
}
pub(crate) fn finalize_persist(&mut self) -> Result<(), Error> {
if self.persist.is_none() || self.store.len() == self.original_len {
return Ok(());
}
let Some(ps) = self.persist.take() else {
return Ok(());
};
let BoundedPersistState {
strategy,
threshold,
page_size,
free,
old_blocks,
old_ext_addr,
} = ps;
let (region, _ext_len) = read_single_chunk_ext_region(&self.store, old_ext_addr)?;
let (new_old_blocks, new_free) = match free {
PersistFree::Flat(mut free) => {
for &(a, l) in &old_blocks {
free.free(a, l);
}
let sections = free_sections(&free);
let nob = self
.store
.write_persist_tail(®ion, strategy, threshold, page_size, §ions)?;
(nob, PersistFree::Flat(free))
}
PersistFree::Paged {
mut meta,
mut raw_small,
raw_large,
} => {
self.store.pad_to_page()?;
if let Some(pg) = self.store.paged.as_mut() {
for &(a, l) in &pg.meta_pad {
meta.free(a, l);
}
for &(a, l) in &pg.raw_pad {
raw_small.free(a, l);
}
pg.meta_pad.clear();
pg.raw_pad.clear();
}
for &(a, l) in &old_blocks {
meta.free(a, l);
}
let nob = self.store.write_persist_tail_paged(
®ion,
strategy,
threshold,
page_size,
&free_sections(&meta),
&free_sections(&raw_small),
&free_sections(&raw_large),
)?;
(
nob,
PersistFree::Paged {
meta,
raw_small,
raw_large,
},
)
}
};
self.persist = Some(BoundedPersistState {
strategy,
threshold,
page_size,
free: new_free,
old_blocks: new_old_blocks,
old_ext_addr: self
.store
.superblock
.superblock_extension_address
.unwrap_or(old_ext_addr),
});
Ok(())
}
pub(crate) fn append_geometry(&mut self, oh_addr: u64) -> Result<AppendGeometry, Error> {
let Self { store, located, .. } = self;
let st = match located.entry(oh_addr) {
std::collections::hash_map::Entry::Occupied(e) => e.into_mut(),
std::collections::hash_map::Entry::Vacant(e) => {
e.insert(locate_dataset_state(&*store, oh_addr)?)
}
};
let chunk_elems = st.loc.chunk_elems.max(1);
let batch_chunks = (APPEND_BATCH_BYTES / (st.loc.chunk_bytes.max(1) as u64)).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: batch_chunks * chunk_elems,
})
}
pub(crate) fn append_gathered(
&mut self,
oh_addr: u64,
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 Self { store, located, .. } = self;
let st = match located.entry(oh_addr) {
std::collections::hash_map::Entry::Occupied(e) => e.into_mut(),
std::collections::hash_map::Entry::Vacant(e) => {
e.insert(locate_dataset_state(&*store, oh_addr)?)
}
};
let new_elems = validate_gathered_append(st, b)?;
if new_elems == 0 {
return Ok(());
}
let raw = b.raw();
let chunk_elems = st.loc.chunk_elems.max(1);
if st.pipeline.is_some()
&& (st.loc.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 elem_bytes = st.element_size as u64;
let batch_chunks = (APPEND_BATCH_BYTES / (st.loc.chunk_bytes.max(1) as u64)).max(1);
let full_batch_elems = batch_chunks * chunk_elems;
let mut done = 0u64;
while done < new_elems {
let to_boundary = (chunk_elems - st.loc.current_dim % chunk_elems) % chunk_elems;
let take = (new_elems - done).min(to_boundary + full_batch_elems);
let start = (done * elem_bytes).to_usize()?;
let end = ((done + take) * elem_bytes).to_usize()?;
let batch = &raw[start..end];
let plan = plan_ea_append(
&*store,
&st.loc,
&st.datatype,
&st.spatial,
st.element_size,
st.pipeline.as_ref(),
batch,
take,
)
.map_err(as_inplace_error)?;
apply_ea_append(store, &mut st.loc, &plan, max_phase).map_err(as_inplace_error)?;
if max_phase < 4 {
return Ok(());
}
done += take;
}
Ok(())
}
}
fn ext_file_space_info(raw: &RawSource<'_>, superblock: &Superblock) -> Option<FileSpaceInfo> {
let rel = superblock.superblock_extension_address?;
if rel == u64::MAX {
return None;
}
let header = ObjectHeader::parse_from_source(
raw,
rel,
superblock.offset_size,
superblock.length_size,
0,
)
.ok()?;
let msg = header
.messages
.iter()
.find(|m| m.msg_type == MessageType::FileSpaceInfo)?;
FileSpaceInfo::parse(&msg.data, superblock.offset_size, superblock.length_size).ok()
}
fn load_bounded_persist(
raw: &RawSource<'_>,
superblock: &Superblock,
) -> Result<Option<BoundedPersistState>, Error> {
let Some(info) = ext_file_space_info(raw, superblock) else {
return Ok(None);
};
let paged = info.strategy == FileSpaceStrategy::Page;
if paged && !info.persist {
return Err(Error::EditUnsupported(
"bounded read-write of a paged file (H5F_FSPACE_STRATEGY_PAGE) requires \
persisted free space (persist=true); use File::open_rw",
));
}
if !info.persist {
return Ok(None);
}
let os = superblock.offset_size;
let ext_addr = superblock
.superblock_extension_address
.expect("ext_file_space_info returned Some, so the extension address is set");
let file_len = raw.len();
let mut old_blocks = Vec::new();
let free = if paged {
let mut tagged: Vec<(FreeSection, u8)> = Vec::new(); for (k, &m) in info.manager_addrs.iter().enumerate() {
if m == u64::MAX {
continue;
}
let (secs, blocks) = read_persisted_sections_source(raw, &[m], 0, os)
.unwrap_or_else(|_| (Vec::new(), Vec::new()));
old_blocks.extend(blocks);
let which = match k {
2 => 1u8,
6 => 2u8,
_ => 0u8, };
for s in secs {
tagged.push((s, which));
}
}
tagged.sort_by_key(|(s, _)| s.addr);
let mut meta = FreeList::new();
let mut raw_small = FreeList::new();
let mut raw_large = FreeList::new();
let mut prev_end = 0u64;
for (s, which) 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;
match which {
1 => raw_small.free(s.addr, s.size),
2 => raw_large.free(s.addr, s.size),
_ => meta.free(s.addr, s.size),
}
}
PersistFree::Paged {
meta,
raw_small,
raw_large,
}
} else {
let (mut sections, manager_blocks) =
read_persisted_sections_source(raw, &info.manager_addrs, 0, os)
.unwrap_or_else(|_| (Vec::new(), Vec::new()));
old_blocks.extend(manager_blocks);
let mut free = FreeList::new();
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;
free.free(s.addr, s.size);
}
PersistFree::Flat(free)
};
let (_region, ext_len) = read_single_chunk_ext_region(raw, ext_addr)?;
old_blocks.push((ext_addr, ext_len));
Ok(Some(BoundedPersistState {
strategy: info.strategy,
threshold: info.threshold,
page_size: info.page_size,
free,
old_blocks,
old_ext_addr: ext_addr,
}))
}
fn free_sections(free: &FreeList) -> Vec<FreeSection> {
free.sections()
.into_iter()
.map(|(addr, size)| FreeSection { addr, size })
.collect()
}
fn align_up(value: u64, page: u64) -> u64 {
value.div_ceil(page) * page
}
fn split_at_pages(sections: &[FreeSection], page: u64) -> Vec<FreeSection> {
let mut out = Vec::new();
for s in sections {
let end = s.addr.saturating_add(s.size);
let mut start = s.addr;
while start < end {
let boundary = (start / page + 1) * page;
let piece_end = end.min(boundary);
out.push(FreeSection {
addr: start,
size: piece_end - start,
});
start = piece_end;
}
}
out
}
#[cfg(test)]
mod tests {
use super::*;
use crate::group_v2;
use crate::writer::FileBuilder;
use tempfile::tempdir;
fn build(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();
}
fn dataset_addr(engine: &BoundedEngine) -> u64 {
group_v2::resolve_path_any_from_source(engine.store(), &engine.store().superblock, "d")
.unwrap()
}
#[test]
fn 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(&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 =
BoundedEngine::open(&p, MetadataCacheConfig::disabled()).unwrap();
let addr = dataset_addr(&engine);
let mut b = AppendBuilder::new();
b.append_i32(&(n..n + add).collect::<Vec<_>>());
engine.append_gathered(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 multi_batch_append_commits_every_batch() {
let dir = tempdir().unwrap();
let p = dir.path().join("multibatch.h5");
build(&p, 5, 512);
let total = 700_000i32;
{
let mut engine = BoundedEngine::open(&p, MetadataCacheConfig::disabled()).unwrap();
let addr = dataset_addr(&engine);
let mut b = AppendBuilder::new();
b.append_i32(&(5..total).collect::<Vec<_>>());
engine.append_gathered(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 persist_append_without_finalize_is_readable() {
let dir = tempdir().unwrap();
let p = dir.path().join("persist_crash.h5");
let mut b = 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 = BoundedEngine::open(&p, MetadataCacheConfig::disabled()).unwrap();
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_gathered(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 paged_reopen_after_crash_realigns_and_stays_readable() {
let dir = tempdir().unwrap();
let p = dir.path().join("paged_crash.h5");
let mut b = 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 = BoundedEngine::open(&p, MetadataCacheConfig::disabled()).unwrap();
let addr = dataset_addr(&engine);
let mut ab = AppendBuilder::new();
ab.append_i32(&(64..2000).collect::<Vec<_>>());
engine.append_gathered(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 = BoundedEngine::open(&p, MetadataCacheConfig::disabled()).unwrap();
let addr = dataset_addr(&engine);
let mut ab = AppendBuilder::new();
ab.append_i32(&(2000..2500).collect::<Vec<_>>());
engine.append_gathered(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 split_at_pages_splits_on_boundaries_preserving_total() {
let page = 4096;
assert_eq!(
split_at_pages(
&[FreeSection {
addr: 3740,
size: 4452
}],
page
),
vec![
FreeSection {
addr: 3740,
size: 356
}, FreeSection {
addr: 4096,
size: 4096
}, ]
);
assert_eq!(
split_at_pages(
&[FreeSection {
addr: 100,
size: 200
}],
page
),
vec![FreeSection {
addr: 100,
size: 200
}]
);
let whole = split_at_pages(
&[FreeSection {
addr: 0,
size: 3 * 4096,
}],
page,
);
assert_eq!(whole.len(), 3);
assert!(whole.iter().all(|s| s.size == 4096));
let pieces = split_at_pages(
&[FreeSection {
addr: 5000,
size: 10000,
}],
page,
);
assert_eq!(pieces.iter().map(|s| s.size).sum::<u64>(), 10000);
let mut prev = pieces[0].addr;
for s in &pieces {
assert_eq!(s.addr, prev);
prev = s.addr + s.size;
}
assert_eq!(prev, 15000);
}
}