facade: File::open_swmr, Dataset::refresh and bounded retries for live SWMR files
Range-read milestone M5 (docs/design/swmr.md): read a file a libhdf5 SWMR writer (h5py f.swmr_mode = True) is still appending to, as h5py's File(path, 'r', swmr=True) does. - File::open_swmr / open_storage_swmr: the file is read through the new FileStorage (pread / seek_read, never mapped, len() the current length), reads are bounded by the length at each read, and the chunk cache is not used (a cached index hides new chunks; a cached edge chunk reads as fill where the writer has since written). A file open for writing without SWMR is refused with Error::Locked, as libhdf5 refuses it. - Dataset::refresh re-reads the object header (H5Drefresh). - Operations that fail with an error a racing write can cause (every format error but a wrong path/selection, an unsupported feature or a bad argument) are run again, up to 100 attempts (libhdf5's default metadata read attempts for SWMR readers), 1 us doubling to 10 ms apart; nested operations retry as a whole. swmr_retries() counts them, swmr_writer_active() re-reads the superblock flags. Tests: FileStorage growth and retry policy (unit); the mid-write copy through open_swmr; a storage that garbles reads (retried, given up after the attempts, never returned); and a live h5py writer appending to Extensible-Array (plain and gzip) and v2-B-tree datasets for 2500 steps while two Rust reader threads and h5py's SWMR reader check every value, then the closed file read equal to h5py. Leaving the chunk cache on in live mode fails the live test. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
This commit is contained in:
@@ -47,6 +47,7 @@ pub mod lazy;
|
||||
#[cfg(feature = "mmap")]
|
||||
pub mod mmap_file;
|
||||
pub mod reader;
|
||||
mod swmr;
|
||||
pub mod types;
|
||||
pub mod vlen;
|
||||
pub mod writer;
|
||||
@@ -57,6 +58,7 @@ pub use lazy::{LazyDataset, LazyFile, LazyGroup};
|
||||
#[cfg(feature = "mmap")]
|
||||
pub use mmap_file::{MmapDataset, MmapFile, MmapGroup};
|
||||
pub use reader::{Dataset, File, Group, SharedStorage, VdsResolver};
|
||||
pub use swmr::{FileStorage, SWMR_READ_ATTEMPTS};
|
||||
pub use types::{AttrValue, DType};
|
||||
pub use vlen::VlenValue;
|
||||
pub use writer::FileBuilder;
|
||||
|
||||
+229
-12
@@ -115,6 +115,10 @@ struct FileData {
|
||||
/// parser reads asks for it, and the patched/overlay checks and range
|
||||
/// conversions behind it cost a local metadata walk a few percent.
|
||||
contiguous: Option<WholeView>,
|
||||
/// A file a SWMR writer may still be appending to
|
||||
/// ([`File::open_swmr`]): reads are bounded by the storage's current
|
||||
/// length, not by `end`, and nothing is ever read as one slice.
|
||||
live: bool,
|
||||
}
|
||||
|
||||
/// A borrow of the HDF5 data held by a [`FileData`]'s own `backing` or
|
||||
@@ -138,7 +142,7 @@ impl FileData {
|
||||
/// bytes past the recorded end of file are not read, as in libhdf5.
|
||||
fn new(mut backing: Backing) -> Result<(Self, Superblock), Error> {
|
||||
if let Backing::Storage(storage) = backing {
|
||||
return Self::new_storage(storage);
|
||||
return Self::new_storage(storage, false);
|
||||
}
|
||||
let whole = backing.whole_file().unwrap_or_default();
|
||||
let (user_block, hdf5) = signature::split_user_block(whole)?;
|
||||
@@ -173,14 +177,15 @@ impl FileData {
|
||||
overlay: Vec::new(),
|
||||
image_error,
|
||||
contiguous: None,
|
||||
live: false,
|
||||
};
|
||||
data.contiguous = data.find_contiguous();
|
||||
Ok((data, superblock))
|
||||
}
|
||||
|
||||
/// [`Self::new`] for a [`Storage`] backend: the same checks, through
|
||||
/// reads of the storage.
|
||||
fn new_storage(storage: SharedStorage) -> Result<(Self, Superblock), Error> {
|
||||
/// reads of the storage. A `live` file is never read as one slice.
|
||||
fn new_storage(storage: SharedStorage, live: bool) -> Result<(Self, Superblock), Error> {
|
||||
let file_len = storage.len();
|
||||
let base = signature::find_signature_in(&*storage)?;
|
||||
let mut data = Self {
|
||||
@@ -191,6 +196,7 @@ impl FileData {
|
||||
overlay: Vec::new(),
|
||||
image_error: None,
|
||||
contiguous: None,
|
||||
live,
|
||||
};
|
||||
// Worked out again below, once the end of file and any cache image
|
||||
// are known.
|
||||
@@ -230,6 +236,9 @@ impl FileData {
|
||||
|
||||
/// [`Self::contiguous`], worked out from `backing` and `patched`.
|
||||
fn find_contiguous(&self) -> Option<WholeView> {
|
||||
if self.live {
|
||||
return None;
|
||||
}
|
||||
let bytes = self.compute_contiguous()?;
|
||||
Some(WholeView {
|
||||
ptr: bytes.as_ptr(),
|
||||
@@ -295,6 +304,12 @@ impl Storage for FileData {
|
||||
if let Some(all) = self.contiguous() {
|
||||
return all.read_at(offset, len);
|
||||
}
|
||||
if self.live {
|
||||
// No end fixed at open: the storage's reads end where the file
|
||||
// ends now.
|
||||
let bytes = cut_to(self.remote()?.read_at(self.base + offset, len)?, len);
|
||||
return Ok(self.with_overlay(offset, bytes));
|
||||
}
|
||||
let size = self.end - self.base;
|
||||
let len = usize::try_from(size.saturating_sub(offset)).map_or(len, |avail| avail.min(len));
|
||||
if len == 0 {
|
||||
@@ -305,6 +320,11 @@ impl Storage for FileData {
|
||||
}
|
||||
|
||||
fn len(&self) -> u64 {
|
||||
if self.live
|
||||
&& let Backing::Storage(s) = &self.backing
|
||||
{
|
||||
return s.len().saturating_sub(self.base);
|
||||
}
|
||||
self.end - self.base
|
||||
}
|
||||
|
||||
@@ -312,7 +332,7 @@ impl Storage for FileData {
|
||||
if let Some(all) = self.contiguous() {
|
||||
return all.read_ranges(ranges);
|
||||
}
|
||||
let size = self.end - self.base;
|
||||
let size = Storage::len(self);
|
||||
let shifted: Vec<Range<u64>> = ranges
|
||||
.iter()
|
||||
.map(|r| {
|
||||
@@ -375,6 +395,9 @@ pub struct File {
|
||||
/// Resolves external Virtual Dataset source files instead of
|
||||
/// `base_dir` (see [`File::set_vds_resolver`]).
|
||||
vds_resolver: Option<VdsResolver>,
|
||||
/// Attempts per operation on a live file, and the retries made (see
|
||||
/// [`File::open_swmr`]).
|
||||
swmr: crate::swmr::Retries,
|
||||
}
|
||||
|
||||
impl File {
|
||||
@@ -394,6 +417,7 @@ impl File {
|
||||
chunk_cache: ChunkCache::new(),
|
||||
base_dir,
|
||||
vds_resolver: None,
|
||||
swmr: crate::swmr::Retries::default(),
|
||||
})
|
||||
}
|
||||
#[cfg(not(feature = "mmap"))]
|
||||
@@ -428,6 +452,7 @@ impl File {
|
||||
chunk_cache: ChunkCache::new(),
|
||||
base_dir: None,
|
||||
vds_resolver: None,
|
||||
swmr: crate::swmr::Retries::default(),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -463,9 +488,130 @@ impl File {
|
||||
chunk_cache: ChunkCache::new(),
|
||||
base_dir: None,
|
||||
vds_resolver: None,
|
||||
swmr: crate::swmr::Retries::default(),
|
||||
})
|
||||
}
|
||||
|
||||
/// Open a file that a libhdf5 SWMR writer (h5py `f.swmr_mode = True`)
|
||||
/// may still be appending to, as libhdf5's SWMR reader does
|
||||
/// (`H5F_ACC_SWMR_READ`, h5py `File(path, "r", swmr=True)`).
|
||||
///
|
||||
/// The file is read with positioned reads ([`FileStorage`](crate::FileStorage)),
|
||||
/// never mapped, and reads are bounded by the file's length at the
|
||||
/// time of each read rather than a length fixed at open. The chunk cache
|
||||
/// is not used, so every read reads the chunk index and chunks as they
|
||||
/// are now. A [`Dataset`] handle keeps the extent it was opened (or
|
||||
/// last [refreshed](Dataset::refresh)) with, as in libhdf5: call
|
||||
/// [`Dataset::refresh`] to see the writer's appends.
|
||||
///
|
||||
/// An operation that fails with an error a concurrent write can cause
|
||||
/// (a checksum mismatch, a short read, a bad signature or version byte,
|
||||
/// a chunk or chunk index that does not decode: every format error but
|
||||
/// a wrong path or selection, an unsupported feature or a bad
|
||||
/// argument) is run again from the start, up to
|
||||
/// [`swmr_read_attempts`](Self::swmr_read_attempts) times (100 by
|
||||
/// default, libhdf5's default for SWMR readers), with a pause of 1 µs
|
||||
/// doubling up to 10 ms between attempts. Data is only returned from an
|
||||
/// attempt in which every structure read verified, so a torn read is
|
||||
/// an error (after the last attempt), never data.
|
||||
///
|
||||
/// A file whose superblock says it is open for writing without SWMR is
|
||||
/// refused with [`Error::Locked`], as libhdf5 refuses it: such a writer
|
||||
/// does not order its writes for readers. Any other file opens (one
|
||||
/// whose writer has closed it reads like [`File::open`]).
|
||||
pub fn open_swmr<P: AsRef<std::path::Path>>(path: P) -> Result<Self, Error> {
|
||||
let storage = crate::swmr::FileStorage::open(path.as_ref()).map_err(Error::Io)?;
|
||||
let mut f = Self::open_storage_swmr(Arc::new(storage))?;
|
||||
f.base_dir = path.as_ref().parent().map(|p| p.to_path_buf());
|
||||
Ok(f)
|
||||
}
|
||||
|
||||
/// [`File::open_swmr`] over any [`Storage`] whose [`Storage::len`]
|
||||
/// follows the file as it grows. A storage that caches blocks, or pins
|
||||
/// the file's length at open (`clawhdf5-remote`'s `BlockCache` and
|
||||
/// `HttpStorage`), does not show the writer's appends.
|
||||
pub fn open_storage_swmr(storage: SharedStorage) -> Result<Self, Error> {
|
||||
let swmr = crate::swmr::Retries::default();
|
||||
let (data, superblock) = swmr.retry(|| FileData::new_storage(storage.clone(), true))?;
|
||||
if superblock.version >= 3 && superblock.is_write_access() && !superblock.is_swmr_write() {
|
||||
return Err(Error::Locked(
|
||||
"the file is open for writing without SWMR (libhdf5: \"file is already open \
|
||||
for write\"); a SWMR reader needs a SWMR writer"
|
||||
.into(),
|
||||
));
|
||||
}
|
||||
Ok(Self {
|
||||
data,
|
||||
superblock,
|
||||
chunk_cache: ChunkCache::new(),
|
||||
base_dir: None,
|
||||
vds_resolver: None,
|
||||
swmr,
|
||||
})
|
||||
}
|
||||
|
||||
/// Whether the file was opened with [`File::open_swmr`] or
|
||||
/// [`File::open_storage_swmr`].
|
||||
pub fn is_swmr_read(&self) -> bool {
|
||||
self.data.live
|
||||
}
|
||||
|
||||
/// How many times an operation on a SWMR-read file is tried (see
|
||||
/// [`File::open_swmr`]); 100 unless set. Files not opened for SWMR
|
||||
/// reading try every operation once.
|
||||
pub fn swmr_read_attempts(&self) -> u32 {
|
||||
self.swmr.attempts
|
||||
}
|
||||
|
||||
/// Set how many times an operation on a SWMR-read file is tried before
|
||||
/// its error is returned (libhdf5's `H5Pset_metadata_read_attempts`); at
|
||||
/// least 1.
|
||||
pub fn set_swmr_read_attempts(&mut self, attempts: u32) {
|
||||
self.swmr.attempts = attempts.max(1);
|
||||
}
|
||||
|
||||
/// How many times an operation on this SWMR-read file has been run
|
||||
/// again because a concurrent write made it fail (the counterpart of
|
||||
/// libhdf5's `H5Fget_metadata_read_retry_info`), counting the attempts
|
||||
/// made while opening it. Always 0 for other files.
|
||||
pub fn swmr_retries(&self) -> u64 {
|
||||
self.swmr.retries()
|
||||
}
|
||||
|
||||
/// Whether a SWMR writer has the file open now: the superblock is read
|
||||
/// again and its SWMR-write flag returned. libhdf5 clears the flag when
|
||||
/// the writer closes the file, so a reader can stop following it then
|
||||
/// (and one last [`Dataset::refresh`] sees the final extents).
|
||||
pub fn swmr_writer_active(&self) -> Result<bool, Error> {
|
||||
self.retry(|| {
|
||||
let sb = Superblock::parse_in(&self.data, 0)?;
|
||||
Ok(sb.version >= 3 && sb.is_swmr_write())
|
||||
})
|
||||
}
|
||||
|
||||
/// Run `op`, again on errors a SWMR writer can cause when the file is
|
||||
/// live (see [`File::open_swmr`]); once otherwise.
|
||||
fn retry<T>(&self, op: impl FnMut() -> Result<T, Error>) -> Result<T, Error> {
|
||||
if self.data.live {
|
||||
self.swmr.retry(op)
|
||||
} else {
|
||||
let mut op = op;
|
||||
op()
|
||||
}
|
||||
}
|
||||
|
||||
/// The chunk cache reads of this file go through: the file's own, or
|
||||
/// for a live file a fresh one per read, since a cached chunk index
|
||||
/// would hide the writer's new chunks and a cached edge chunk would
|
||||
/// read as fill where the writer has written since.
|
||||
fn with_chunk_cache<T>(&self, op: impl FnOnce(&ChunkCache) -> T) -> T {
|
||||
if self.data.live {
|
||||
op(&ChunkCache::new())
|
||||
} else {
|
||||
op(&self.chunk_cache)
|
||||
}
|
||||
}
|
||||
|
||||
/// Resolve external Virtual Dataset source files (their names as the
|
||||
/// mappings store them) with `resolver`, instead of reading them from
|
||||
/// the directory of the file (for [`File::open`]) or refusing them (for
|
||||
@@ -488,6 +634,7 @@ impl File {
|
||||
///
|
||||
/// The path uses `/` separators (e.g., `"group1/values"`).
|
||||
pub fn dataset(&self, path: &str) -> Result<Dataset<'_>, Error> {
|
||||
self.retry(|| {
|
||||
let addr = with_bytes!(self.data.meta()?, |d| group_v2::resolve_path_any_in(
|
||||
d,
|
||||
&self.superblock,
|
||||
@@ -499,9 +646,11 @@ impl File {
|
||||
}
|
||||
Dataset {
|
||||
file: self,
|
||||
address: addr,
|
||||
header: hdr,
|
||||
}
|
||||
.check_open()
|
||||
})
|
||||
}
|
||||
|
||||
/// A `Dataset` handle for the object header at `address` (an address
|
||||
@@ -511,15 +660,18 @@ impl File {
|
||||
/// are scanned); keep the address instead to open the same dataset
|
||||
/// repeatedly.
|
||||
pub fn dataset_at(&self, address: u64) -> Result<Dataset<'_>, Error> {
|
||||
self.retry(|| {
|
||||
let hdr = self.parse_header(address)?;
|
||||
if !has_message(&hdr, MessageType::DataLayout) {
|
||||
return Err(Error::NotADataset(format!("object at address {address}")));
|
||||
}
|
||||
Dataset {
|
||||
file: self,
|
||||
address,
|
||||
header: hdr,
|
||||
}
|
||||
.check_open()
|
||||
})
|
||||
}
|
||||
|
||||
/// A `Group` handle for the object header at `address` (from
|
||||
@@ -538,11 +690,11 @@ impl File {
|
||||
/// The path uses `/` separators (e.g., `"sensors"`).
|
||||
/// Use `"/"` or `""` for the root group.
|
||||
pub fn group(&self, path: &str) -> Result<Group<'_>, Error> {
|
||||
let addr = with_bytes!(self.data.meta()?, |d| group_v2::resolve_path_any_in(
|
||||
d,
|
||||
&self.superblock,
|
||||
path
|
||||
))?;
|
||||
let addr = self.retry(|| {
|
||||
Ok(with_bytes!(self.data.meta()?, |d| {
|
||||
group_v2::resolve_path_any_in(d, &self.superblock, path)
|
||||
})?)
|
||||
})?;
|
||||
Ok(Group {
|
||||
file: self,
|
||||
address: addr,
|
||||
@@ -803,15 +955,19 @@ impl<'f> Group<'f> {
|
||||
|
||||
/// Get a dataset within this group by name.
|
||||
pub fn dataset(&self, name: &str) -> Result<Dataset<'f>, Error> {
|
||||
let hdr = self.file.parse_header(self.child_address(name)?)?;
|
||||
self.file.retry(|| {
|
||||
let address = self.child_address(name)?;
|
||||
let hdr = self.file.parse_header(address)?;
|
||||
if !has_message(&hdr, MessageType::DataLayout) {
|
||||
return Err(Error::NotADataset(name.to_string()));
|
||||
}
|
||||
Dataset {
|
||||
file: self.file,
|
||||
address,
|
||||
header: hdr,
|
||||
}
|
||||
.check_open()
|
||||
})
|
||||
}
|
||||
|
||||
/// Get a subgroup within this group by name.
|
||||
@@ -842,6 +998,7 @@ impl<'f> Group<'f> {
|
||||
/// name index rather than by listing the group (see
|
||||
/// [`group_v2::resolve_child`]).
|
||||
fn child_address(&self, name: &str) -> Result<u64, Error> {
|
||||
self.file.retry(|| {
|
||||
with_bytes!(self.file.data.meta()?, |d| group_v2::resolve_child_in(
|
||||
d,
|
||||
&self.file.superblock,
|
||||
@@ -849,6 +1006,7 @@ impl<'f> Group<'f> {
|
||||
name
|
||||
))
|
||||
.map_err(Error::Format)
|
||||
})
|
||||
}
|
||||
|
||||
/// This group's children that can be opened, as `(name, object header
|
||||
@@ -869,10 +1027,12 @@ impl<'f> Group<'f> {
|
||||
/// [`group_v2::resolve_group_children`]); dangling, external and
|
||||
/// user-defined links are left out.
|
||||
fn children(&self) -> Result<Vec<GroupEntry>, Error> {
|
||||
self.file.retry(|| {
|
||||
with_bytes!(self.file.data.meta()?, |d| {
|
||||
group_v2::resolve_group_children_in(d, &self.file.superblock, self.address)
|
||||
})
|
||||
.map_err(Error::Format)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -884,6 +1044,8 @@ impl<'f> Group<'f> {
|
||||
#[derive(Debug)]
|
||||
pub struct Dataset<'f> {
|
||||
file: &'f File,
|
||||
/// Address of the object header, for [`Dataset::refresh`].
|
||||
address: u64,
|
||||
header: ObjectHeader,
|
||||
}
|
||||
|
||||
@@ -906,6 +1068,31 @@ impl<'f> Dataset<'f> {
|
||||
Ok(self)
|
||||
}
|
||||
|
||||
/// Read the dataset's object header again, as libhdf5's `H5Drefresh`
|
||||
/// (h5py `Dataset.refresh()`) does: afterwards [`shape`](Self::shape)
|
||||
/// and every read use the dataset's extent as it is in the file now.
|
||||
/// For a file a SWMR writer is appending to ([`File::open_swmr`]) this
|
||||
/// is how a reader sees the appends; a transient failure is retried
|
||||
/// (see there). On error the handle keeps its old header.
|
||||
pub fn refresh(&mut self) -> Result<(), Error> {
|
||||
let file = self.file;
|
||||
let address = self.address;
|
||||
let fresh = file.retry(|| {
|
||||
let hdr = file.parse_header(address)?;
|
||||
if !has_message(&hdr, MessageType::DataLayout) {
|
||||
return Err(Error::NotADataset(format!("object at address {address}")));
|
||||
}
|
||||
Dataset {
|
||||
file,
|
||||
address,
|
||||
header: hdr,
|
||||
}
|
||||
.check_open()
|
||||
})?;
|
||||
self.header = fresh.header;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Returns the shape (dimensions) of the dataset.
|
||||
pub fn shape(&self) -> Result<Vec<u64>, Error> {
|
||||
let ds = self.dataspace()?;
|
||||
@@ -1096,6 +1283,15 @@ impl<'f> Dataset<'f> {
|
||||
if matches!(selection, clawhdf5_format::selection::Selection::All) {
|
||||
return self.read_raw();
|
||||
}
|
||||
self.file.retry(|| self.read_selection_once(selection))
|
||||
}
|
||||
|
||||
/// [`read_selection`](Self::read_selection) of a selection other than
|
||||
/// `All`, once.
|
||||
fn read_selection_once(
|
||||
&self,
|
||||
selection: &clawhdf5_format::selection::Selection,
|
||||
) -> Result<Vec<u8>, Error> {
|
||||
let dt = self.datatype()?;
|
||||
let ds = self.dataspace()?;
|
||||
let dl = self.data_layout()?;
|
||||
@@ -1182,6 +1378,17 @@ impl<'f> Dataset<'f> {
|
||||
if matches!(selection, clawhdf5_format::selection::Selection::All) {
|
||||
return full();
|
||||
}
|
||||
self.file
|
||||
.retry(|| self.read_typed_selection_once(selection, convert))
|
||||
}
|
||||
|
||||
/// [`read_typed_selection`](Self::read_typed_selection) of a selection
|
||||
/// other than `All`, once.
|
||||
fn read_typed_selection_once<T: data_read::NativeElement>(
|
||||
&self,
|
||||
selection: &clawhdf5_format::selection::Selection,
|
||||
convert: fn(&[u8], &Datatype) -> Result<Vec<T>, FormatError>,
|
||||
) -> Result<Vec<T>, Error> {
|
||||
let dt = self.datatype()?;
|
||||
if T::is_native(&dt) && self.file.data.contiguous().is_some() {
|
||||
if let Ok(Some(raw)) = self.read_raw_ref() {
|
||||
@@ -1587,6 +1794,8 @@ impl<'f> Dataset<'f> {
|
||||
}
|
||||
let ds = self.dataspace()?;
|
||||
let pipeline = self.filter_pipeline()?;
|
||||
self.file.retry(|| {
|
||||
self.file.with_chunk_cache(|cache| {
|
||||
Ok(data_read::read_chunked_native_in::<T, _>(
|
||||
&self.header.messages,
|
||||
&self.file.data,
|
||||
@@ -1596,11 +1805,17 @@ impl<'f> Dataset<'f> {
|
||||
pipeline.as_ref(),
|
||||
self.file.offset_size(),
|
||||
self.file.length_size(),
|
||||
Some(&self.file.chunk_cache),
|
||||
Some(cache),
|
||||
)?)
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
fn read_raw(&self) -> Result<Vec<u8>, Error> {
|
||||
self.file.retry(|| self.read_raw_once())
|
||||
}
|
||||
|
||||
fn read_raw_once(&self) -> Result<Vec<u8>, Error> {
|
||||
let dt = self.datatype()?;
|
||||
let ds = self.dataspace()?;
|
||||
let dl = self.data_layout()?;
|
||||
@@ -1622,6 +1837,7 @@ impl<'f> Dataset<'f> {
|
||||
self.file.offset_size(),
|
||||
self.file.length_size(),
|
||||
|| {
|
||||
self.file.with_chunk_cache(|cache| {
|
||||
Ok(data_read::read_raw_data_cached_in(
|
||||
&self.file.data,
|
||||
&dl,
|
||||
@@ -1630,8 +1846,9 @@ impl<'f> Dataset<'f> {
|
||||
pipeline.as_ref(),
|
||||
self.file.offset_size(),
|
||||
self.file.length_size(),
|
||||
&self.file.chunk_cache,
|
||||
cache,
|
||||
)?)
|
||||
})
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,298 @@
|
||||
//! Reading files a SWMR writer is still appending to (see
|
||||
//! `docs/design/swmr.md` and [`File::open_swmr`](crate::File::open_swmr)):
|
||||
//! a [`FileStorage`] whose length grows with the file, and the bounded
|
||||
//! retries of operations a concurrent write can make fail.
|
||||
|
||||
use std::borrow::Cow;
|
||||
use std::cell::Cell;
|
||||
use std::path::Path;
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::time::Duration;
|
||||
|
||||
use clawhdf5_format::error::FormatError;
|
||||
use clawhdf5_format::storage::Storage;
|
||||
|
||||
use crate::error::Error;
|
||||
|
||||
/// How many times a live file ([`File::open_swmr`](crate::File::open_swmr))
|
||||
/// tries an operation that fails with an error a concurrent write can cause
|
||||
/// before returning the error: 100, libhdf5's default number of metadata
|
||||
/// read attempts for SWMR access (`H5Pset_metadata_read_attempts`).
|
||||
pub const SWMR_READ_ATTEMPTS: u32 = 100;
|
||||
|
||||
/// Longest pause between two attempts.
|
||||
const MAX_PAUSE: Duration = Duration::from_millis(10);
|
||||
|
||||
/// A local file read with positioned reads (`pread` on Unix, `seek_read` on
|
||||
/// Windows), never mapped, whose [`Storage::len`] is the file's length at
|
||||
/// the time of the call: a [`Storage`] for a file that another process is
|
||||
/// appending to. A read past the end is short, as the trait allows.
|
||||
#[derive(Debug)]
|
||||
pub struct FileStorage {
|
||||
file: std::fs::File,
|
||||
#[cfg(not(any(unix, windows)))]
|
||||
lock: std::sync::Mutex<()>,
|
||||
}
|
||||
|
||||
impl FileStorage {
|
||||
/// Open the file at `path` for reading.
|
||||
pub fn open<P: AsRef<Path>>(path: P) -> std::io::Result<Self> {
|
||||
Ok(Self::new(std::fs::File::open(path)?))
|
||||
}
|
||||
|
||||
/// A storage over an open file.
|
||||
pub fn new(file: std::fs::File) -> Self {
|
||||
Self {
|
||||
file,
|
||||
#[cfg(not(any(unix, windows)))]
|
||||
lock: std::sync::Mutex::new(()),
|
||||
}
|
||||
}
|
||||
|
||||
/// Up to `buf.len()` bytes at `offset`; fewer only at the end of file.
|
||||
fn read_into(&self, offset: u64, buf: &mut [u8]) -> std::io::Result<usize> {
|
||||
let mut got = 0;
|
||||
while got < buf.len() {
|
||||
match self.read_once(offset + got as u64, &mut buf[got..]) {
|
||||
Ok(0) => break,
|
||||
Ok(n) => got += n,
|
||||
Err(e) if e.kind() == std::io::ErrorKind::Interrupted => {}
|
||||
Err(e) => return Err(e),
|
||||
}
|
||||
}
|
||||
Ok(got)
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn read_once(&self, offset: u64, buf: &mut [u8]) -> std::io::Result<usize> {
|
||||
std::os::unix::fs::FileExt::read_at(&self.file, buf, offset)
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
fn read_once(&self, offset: u64, buf: &mut [u8]) -> std::io::Result<usize> {
|
||||
std::os::windows::fs::FileExt::seek_read(&self.file, buf, offset)
|
||||
}
|
||||
|
||||
#[cfg(not(any(unix, windows)))]
|
||||
fn read_once(&self, offset: u64, buf: &mut [u8]) -> std::io::Result<usize> {
|
||||
use std::io::{Read, Seek, SeekFrom};
|
||||
let _guard = self.lock.lock().unwrap_or_else(|p| p.into_inner());
|
||||
let mut f = &self.file;
|
||||
f.seek(SeekFrom::Start(offset))?;
|
||||
f.read(buf)
|
||||
}
|
||||
}
|
||||
|
||||
impl Storage for FileStorage {
|
||||
fn read_at(&self, offset: u64, len: usize) -> Result<Cow<'_, [u8]>, FormatError> {
|
||||
// Never allocate more than the file holds, whatever a (possibly
|
||||
// hostile) size field asked for.
|
||||
let avail = self.len().saturating_sub(offset);
|
||||
let len = usize::try_from(avail).map_or(len, |a| a.min(len));
|
||||
let mut buf = vec![0u8; len];
|
||||
let got = self
|
||||
.read_into(offset, &mut buf)
|
||||
.map_err(|e| FormatError::Storage(format!("read of {len} bytes at {offset}: {e}")))?;
|
||||
buf.truncate(got);
|
||||
Ok(Cow::Owned(buf))
|
||||
}
|
||||
|
||||
fn len(&self) -> u64 {
|
||||
self.file.metadata().map_or(0, |m| m.len())
|
||||
}
|
||||
}
|
||||
|
||||
/// Whether `e` can be caused by reading a structure while a SWMR writer
|
||||
/// rewrites or has not yet finished writing it — a checksum mismatch, a
|
||||
/// short read, a bad signature, a chunk index or chunk that does not
|
||||
/// decode, ... (see [`is_transient_format`]) — so that the operation is
|
||||
/// worth running again.
|
||||
pub(crate) fn is_transient(e: &Error) -> bool {
|
||||
match e {
|
||||
Error::Format(f) => is_transient_format(f),
|
||||
Error::Io(io) => io.kind() == std::io::ErrorKind::UnexpectedEof,
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
|
||||
/// A read that raced a write can garble any field of a structure, and the
|
||||
/// parsers report that as whatever check fails first (the checksum, a
|
||||
/// signature, a version byte, a size), so every format error counts except
|
||||
/// those that bytes read later cannot change: a name or selection the
|
||||
/// caller got wrong, a feature this reader does not support, an argument
|
||||
/// that does not fit.
|
||||
fn is_transient_format(e: &FormatError) -> bool {
|
||||
use FormatError as F;
|
||||
!matches!(
|
||||
e,
|
||||
F::PathNotFound(_)
|
||||
| F::SelectionOutOfBounds(_)
|
||||
| F::UnsupportedFilter(_)
|
||||
| F::ExternalDataFilesUnsupported
|
||||
| F::ExternalLinkUnsupported { .. }
|
||||
| F::ContiguousStorageRequired(_)
|
||||
| F::TypeMismatch { .. }
|
||||
| F::DataSizeMismatch { .. }
|
||||
| F::SerializationError(_)
|
||||
| F::CompressionError(_)
|
||||
| F::DuplicateDatasetName(_)
|
||||
| F::InvalidLinkName
|
||||
)
|
||||
}
|
||||
|
||||
thread_local! {
|
||||
/// Set while this thread runs an operation under [`retry`], so an
|
||||
/// operation made of retried operations retries as a whole, not each
|
||||
/// part up to the limit.
|
||||
static RETRYING: Cell<bool> = const { Cell::new(false) };
|
||||
}
|
||||
|
||||
/// A live file's retry policy: attempts per operation, and a count of the
|
||||
/// retries made.
|
||||
#[derive(Debug)]
|
||||
pub(crate) struct Retries {
|
||||
pub(crate) attempts: u32,
|
||||
retried: AtomicU64,
|
||||
}
|
||||
|
||||
impl Default for Retries {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
attempts: SWMR_READ_ATTEMPTS,
|
||||
retried: AtomicU64::new(0),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Retries {
|
||||
/// [`retry`] with this policy, counting the retries.
|
||||
pub(crate) fn retry<T>(&self, op: impl FnMut() -> Result<T, Error>) -> Result<T, Error> {
|
||||
retry(self.attempts, &self.retried, op)
|
||||
}
|
||||
|
||||
/// Retries made so far.
|
||||
pub(crate) fn retries(&self) -> u64 {
|
||||
self.retried.load(Ordering::Relaxed)
|
||||
}
|
||||
}
|
||||
|
||||
/// Run `op` up to `attempts` times while it fails with a transient error
|
||||
/// ([`is_transient`]), pausing 1 µs, 2 µs, 4 µs, … up to 10 ms between
|
||||
/// attempts, and adding each retry to `retried`; return its first success
|
||||
/// or last error. Inside another `retry` on the same thread, `op` runs
|
||||
/// once.
|
||||
fn retry<T>(
|
||||
attempts: u32,
|
||||
retried: &AtomicU64,
|
||||
mut op: impl FnMut() -> Result<T, Error>,
|
||||
) -> Result<T, Error> {
|
||||
if RETRYING.with(Cell::get) {
|
||||
return op();
|
||||
}
|
||||
struct Reset;
|
||||
impl Drop for Reset {
|
||||
fn drop(&mut self) {
|
||||
RETRYING.with(|r| r.set(false));
|
||||
}
|
||||
}
|
||||
RETRYING.with(|r| r.set(true));
|
||||
let _reset = Reset;
|
||||
let mut pause = Duration::from_micros(1);
|
||||
let mut attempt = 1;
|
||||
loop {
|
||||
match op() {
|
||||
Err(e) if attempt < attempts && is_transient(&e) => {
|
||||
retried.fetch_add(1, Ordering::Relaxed);
|
||||
std::thread::sleep(pause);
|
||||
pause = (pause * 2).min(MAX_PAUSE);
|
||||
attempt += 1;
|
||||
}
|
||||
result => return result,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn file_storage_reads_what_the_file_holds_now() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let path = dir.path().join("grow.bin");
|
||||
std::fs::write(&path, b"hello").unwrap();
|
||||
let s = FileStorage::open(&path).unwrap();
|
||||
assert_eq!(s.len(), 5);
|
||||
assert_eq!(&*s.read_at(1, 3).unwrap(), b"ell");
|
||||
assert_eq!(&*s.read_at(3, 10).unwrap(), b"lo");
|
||||
assert!(s.read_at(9, 4).unwrap().is_empty());
|
||||
// The file grows after the storage was opened.
|
||||
use std::io::Write;
|
||||
std::fs::OpenOptions::new()
|
||||
.append(true)
|
||||
.open(&path)
|
||||
.unwrap()
|
||||
.write_all(b", world")
|
||||
.unwrap();
|
||||
assert_eq!(s.len(), 12);
|
||||
assert_eq!(&*s.read_at(3, 100).unwrap(), b"lo, world");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn retry_runs_again_only_for_transient_errors() {
|
||||
let n = AtomicU64::new(0);
|
||||
let mut calls = 0;
|
||||
let r: Result<u32, Error> = retry(5, &n, || {
|
||||
calls += 1;
|
||||
if calls < 3 {
|
||||
Err(Error::Format(FormatError::ChecksumMismatch {
|
||||
expected: 1,
|
||||
computed: 2,
|
||||
}))
|
||||
} else {
|
||||
Ok(7)
|
||||
}
|
||||
});
|
||||
assert_eq!(r.unwrap(), 7);
|
||||
assert_eq!(calls, 3);
|
||||
assert_eq!(n.load(Ordering::Relaxed), 2);
|
||||
|
||||
// Gives up after `attempts`.
|
||||
let mut calls = 0;
|
||||
let r: Result<(), Error> = retry(4, &n, || {
|
||||
calls += 1;
|
||||
Err(Error::Format(FormatError::UnexpectedEof {
|
||||
expected: 8,
|
||||
available: 0,
|
||||
}))
|
||||
});
|
||||
assert!(r.is_err());
|
||||
assert_eq!(calls, 4);
|
||||
assert_eq!(n.load(Ordering::Relaxed), 5);
|
||||
|
||||
// A permanent error is returned at once.
|
||||
let mut calls = 0;
|
||||
let r: Result<(), Error> = retry(4, &n, || {
|
||||
calls += 1;
|
||||
Err(Error::Format(FormatError::UnsupportedFilter(999)))
|
||||
});
|
||||
assert!(r.is_err());
|
||||
assert_eq!(calls, 1);
|
||||
assert_eq!(n.load(Ordering::Relaxed), 5);
|
||||
|
||||
// Nested: the inner operation runs once per outer attempt.
|
||||
let mut inner = 0;
|
||||
let mut outer = 0;
|
||||
let r: Result<(), Error> = retry(3, &n, || {
|
||||
outer += 1;
|
||||
retry(3, &n, || {
|
||||
inner += 1;
|
||||
Err(Error::Format(FormatError::SignatureNotFound))
|
||||
})
|
||||
});
|
||||
assert!(r.is_err());
|
||||
assert_eq!((outer, inner), (3, 3));
|
||||
// The flag is reset afterwards.
|
||||
assert!(!RETRYING.with(Cell::get));
|
||||
}
|
||||
}
|
||||
@@ -6,10 +6,13 @@
|
||||
//! The live tests need python3 with h5py; they are skipped without it,
|
||||
//! unless `CLAWHDF5_REQUIRE_INTEROP=1`.
|
||||
|
||||
use std::io::{BufRead, BufReader, Read};
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::process::Command;
|
||||
use std::process::{Command, Stdio};
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
|
||||
|
||||
use clawhdf5::File;
|
||||
use clawhdf5::{File, Selection};
|
||||
|
||||
fn python() -> String {
|
||||
std::env::var("CLAWHDF5_PYTHON").unwrap_or_else(|_| "python3".to_string())
|
||||
@@ -77,6 +80,18 @@ fn a_copy_taken_mid_write_reads_past_its_recorded_end_of_file() {
|
||||
check_mid_write_copy(&File::open(&path).unwrap());
|
||||
check_mid_write_copy(&File::open_buffered(&path).unwrap());
|
||||
check_mid_write_copy(&File::from_bytes(std::fs::read(&path).unwrap()).unwrap());
|
||||
let want_a: Vec<i64> = (1..=148).collect();
|
||||
let want_b: Vec<f64> = (0..148)
|
||||
.flat_map(|i| (0..4).map(move |j| (10 * i + j + 1) as f64))
|
||||
.collect();
|
||||
let mm = clawhdf5::MmapFile::open(&path).unwrap();
|
||||
assert_eq!(mm.dataset("a").unwrap().read_i64().unwrap(), want_a);
|
||||
assert_eq!(mm.dataset("b").unwrap().read_f64().unwrap(), want_b);
|
||||
let lazy = clawhdf5::LazyFile::open_mmap(&path).unwrap();
|
||||
assert_eq!(lazy.dataset("a").unwrap().read_i64().unwrap(), want_a);
|
||||
assert_eq!(lazy.dataset("b").unwrap().read_f64().unwrap(), want_b);
|
||||
let storage = File::open_storage(Arc::new(std::fs::read(&path).unwrap())).unwrap();
|
||||
check_mid_write_copy(&storage);
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -97,7 +112,10 @@ print("ok")
|
||||
"#,
|
||||
p = path.display()
|
||||
);
|
||||
let out = Command::new(python()).args(["-c", &script]).output().unwrap();
|
||||
let out = Command::new(python())
|
||||
.args(["-c", &script])
|
||||
.output()
|
||||
.unwrap();
|
||||
assert!(
|
||||
out.status.success(),
|
||||
"h5py failed:\n{}",
|
||||
@@ -105,3 +123,459 @@ print("ok")
|
||||
);
|
||||
check_mid_write_copy(&File::open(&path).unwrap());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn open_swmr_reads_the_mid_write_copy() {
|
||||
let f = File::open_swmr(fixture("swmr_mid_write.h5")).unwrap();
|
||||
assert!(f.is_swmr_read());
|
||||
assert_eq!(f.swmr_read_attempts(), clawhdf5::SWMR_READ_ATTEMPTS);
|
||||
// The copy still has the SWMR-write flag: as far as it says, its writer
|
||||
// is still writing.
|
||||
assert!(f.swmr_writer_active().unwrap());
|
||||
check_mid_write_copy(&f);
|
||||
let mut a = f.dataset("a").unwrap();
|
||||
a.refresh().unwrap();
|
||||
assert_eq!(a.shape().unwrap(), vec![148]);
|
||||
assert_eq!(f.swmr_retries(), 0);
|
||||
assert!(
|
||||
!File::open(fixture("swmr_mid_write.h5"))
|
||||
.unwrap()
|
||||
.is_swmr_read()
|
||||
);
|
||||
}
|
||||
|
||||
/// The mid-write copy with its superblock's flags set to `flags` (and the
|
||||
/// superblock checksum updated).
|
||||
fn mid_write_copy_with_flags(flags: u8) -> Vec<u8> {
|
||||
let mut bytes = std::fs::read(fixture("swmr_mid_write.h5")).unwrap();
|
||||
// Superblock v3, 8-byte offsets: 12 bytes, 4 addresses, checksum.
|
||||
bytes[11] = flags;
|
||||
let sum = clawhdf5_format::checksum::jenkins_lookup3(&bytes[..44]);
|
||||
bytes[44..48].copy_from_slice(&sum.to_le_bytes());
|
||||
bytes
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn open_swmr_refuses_a_file_open_for_writing_without_swmr() {
|
||||
// libhdf5 refuses it too: "file is already open for write".
|
||||
let bytes = mid_write_copy_with_flags(0x01);
|
||||
let err = File::open_storage_swmr(Arc::new(bytes)).unwrap_err();
|
||||
assert!(matches!(err, clawhdf5::Error::Locked(_)), "{err}");
|
||||
// Closed (flags clear): opens, and no writer is active.
|
||||
let f = File::open_storage_swmr(Arc::new(mid_write_copy_with_flags(0))).unwrap();
|
||||
assert!(!f.swmr_writer_active().unwrap());
|
||||
}
|
||||
|
||||
/// A storage whose next `torn` reads come back garbled, as a read racing a
|
||||
/// rewrite of the structure can see them: the middle byte changed, or with
|
||||
/// `invert` every byte (signatures included).
|
||||
struct Torn {
|
||||
bytes: Vec<u8>,
|
||||
torn: AtomicU32,
|
||||
invert: AtomicBool,
|
||||
}
|
||||
|
||||
impl clawhdf5::Storage for Torn {
|
||||
fn read_at(
|
||||
&self,
|
||||
offset: u64,
|
||||
len: usize,
|
||||
) -> Result<std::borrow::Cow<'_, [u8]>, clawhdf5_format::error::FormatError> {
|
||||
let got = self.bytes.as_slice().read_at(offset, len)?;
|
||||
let tear = self
|
||||
.torn
|
||||
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |n| n.checked_sub(1))
|
||||
.is_ok();
|
||||
if tear && !got.is_empty() {
|
||||
let mut v = got.into_owned();
|
||||
if self.invert.load(Ordering::SeqCst) {
|
||||
v.iter_mut().for_each(|b| *b = !*b);
|
||||
} else {
|
||||
let mid = v.len() / 2;
|
||||
v[mid] ^= 0x5a;
|
||||
}
|
||||
return Ok(v.into());
|
||||
}
|
||||
Ok(got)
|
||||
}
|
||||
|
||||
fn len(&self) -> u64 {
|
||||
self.bytes.len() as u64
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_torn_metadata_read_is_retried_and_never_returned() {
|
||||
use clawhdf5::Storage as _;
|
||||
let torn = Arc::new(Torn {
|
||||
bytes: std::fs::read(fixture("swmr_mid_write.h5")).unwrap(),
|
||||
torn: AtomicU32::new(0),
|
||||
invert: AtomicBool::new(false),
|
||||
});
|
||||
let f = File::open_storage_swmr(torn.clone()).unwrap();
|
||||
assert_eq!(f.swmr_retries(), 0);
|
||||
assert!(torn.len() > 0);
|
||||
|
||||
// The next 3 reads (object headers) come back garbled: the lookup
|
||||
// fails its checks, is run again, and returns the right dataset. The
|
||||
// same with every byte garbled, signatures included.
|
||||
for invert in [false, true] {
|
||||
torn.invert.store(invert, Ordering::SeqCst);
|
||||
let before = f.swmr_retries();
|
||||
torn.torn.store(3, Ordering::SeqCst);
|
||||
let a = f.dataset("a").unwrap();
|
||||
assert!(f.swmr_retries() > before, "invert {invert}");
|
||||
assert_eq!(a.read_i64().unwrap(), (1..=148).collect::<Vec<i64>>());
|
||||
// Two garbled reads: the prefix of the root group's header (whose
|
||||
// garbled byte may go unused, as it is read again whole) and the
|
||||
// whole header, which fails its checksum.
|
||||
let before = f.swmr_retries();
|
||||
torn.torn.store(2, Ordering::SeqCst);
|
||||
let b = f.dataset("b").unwrap();
|
||||
assert_eq!(b.read_f64().unwrap().len(), 148 * 4);
|
||||
assert!(f.swmr_retries() > before, "invert {invert}");
|
||||
}
|
||||
torn.invert.store(false, Ordering::SeqCst);
|
||||
|
||||
// With one attempt the error is returned instead (both reads of the
|
||||
// header garbled, as above).
|
||||
let mut f = File::open_storage_swmr(torn.clone()).unwrap();
|
||||
f.set_swmr_read_attempts(1);
|
||||
torn.torn.store(2, Ordering::SeqCst);
|
||||
assert!(f.dataset("a").is_err());
|
||||
assert_eq!(f.swmr_retries(), 0);
|
||||
torn.torn.store(0, Ordering::SeqCst);
|
||||
|
||||
// A refresh that keeps failing gives up after the attempts, and the
|
||||
// handle keeps its extent.
|
||||
f.set_swmr_read_attempts(5);
|
||||
let mut b = f.dataset("b").unwrap();
|
||||
torn.torn.store(1000, Ordering::SeqCst);
|
||||
assert!(b.refresh().is_err());
|
||||
assert_eq!(f.swmr_retries(), 4);
|
||||
torn.torn.store(0, Ordering::SeqCst);
|
||||
assert_eq!(b.shape().unwrap(), vec![148, 4]);
|
||||
b.refresh().unwrap();
|
||||
assert_eq!(b.read_f64().unwrap().len(), 148 * 4);
|
||||
|
||||
// A file not opened for SWMR reading does not retry.
|
||||
let plain = File::open_storage(torn.clone()).unwrap();
|
||||
torn.torn.store(2, Ordering::SeqCst);
|
||||
assert!(plain.dataset("a").is_err());
|
||||
assert_eq!(plain.swmr_retries(), 0);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// A live file: an h5py SWMR writer appends while clawhdf5 and h5py read.
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// The h5py SWMR writer. It creates three chunked datasets whose values are
|
||||
/// a function of their position, so a reader can check every value it
|
||||
/// reads without knowing when it was written:
|
||||
///
|
||||
/// - `a`: int64 `(n,)`, chunks of 100, no filter — `a[i] = i + 1`
|
||||
/// (Extensible Array index);
|
||||
/// - `b`: float64 `(n, 4)`, chunks of 16 x 4, gzip — `b[i, j] = 10 i + j + 1`
|
||||
/// (Extensible Array index);
|
||||
/// - `c`: int32 `(r, k)`, both unlimited, chunks of 8 x 8, no filter —
|
||||
/// `c[i, j] = 1000 i + j + 1` (version-2 B-tree index).
|
||||
///
|
||||
/// It switches to SWMR mode, prints `ready`, then for `steps` steps appends a
|
||||
/// random number of rows to each dataset (and every 7th step a column to
|
||||
/// `c`), resizing before writing as SWMR writers must, flushing each
|
||||
/// dataset and pausing 2 ms, and finally closes the file and writes the
|
||||
/// final contents to `<path>.<name>.bin`. `CLAWHDF5_SWMR_STEPS` sets the
|
||||
/// number of steps (default 2500).
|
||||
const WRITER: &str = r#"
|
||||
import sys, time, random, h5py, numpy as np
|
||||
path, steps = sys.argv[1], int(sys.argv[2])
|
||||
rng = random.Random(1234)
|
||||
f = h5py.File(path, "w", libver="latest")
|
||||
a = f.create_dataset("a", shape=(0,), maxshape=(None,), chunks=(100,), dtype="i8")
|
||||
b = f.create_dataset("b", shape=(0, 4), maxshape=(None, 4), chunks=(16, 4), dtype="f8",
|
||||
compression="gzip")
|
||||
c = f.create_dataset("c", shape=(0, 1), maxshape=(None, None), chunks=(8, 8), dtype="i4")
|
||||
f.swmr_mode = True
|
||||
print("ready", flush=True)
|
||||
na = nb = rc = 0
|
||||
kc = 1
|
||||
for step in range(steps):
|
||||
k = rng.randint(1, 60)
|
||||
a.resize((na + k,)); a[na:na + k] = np.arange(na, na + k) + 1; na += k
|
||||
a.flush()
|
||||
k = rng.randint(1, 20)
|
||||
rows = np.arange(nb, nb + k)[:, None] * 10 + np.arange(4) + 1
|
||||
b.resize((nb + k, 4)); b[nb:nb + k] = rows; nb += k
|
||||
b.flush()
|
||||
if step % 7 == 6:
|
||||
c.resize((rc, kc + 1))
|
||||
if rc:
|
||||
c[:, kc] = np.arange(rc) * 1000 + kc + 1
|
||||
kc += 1
|
||||
k = rng.randint(1, 3)
|
||||
c.resize((rc + k, kc))
|
||||
c[rc:rc + k] = np.arange(rc, rc + k)[:, None] * 1000 + np.arange(kc) + 1
|
||||
rc += k
|
||||
c.flush()
|
||||
time.sleep(0.002)
|
||||
f.close()
|
||||
with h5py.File(path, "r") as f:
|
||||
for name in "abc":
|
||||
open(f"{path}.{name}.bin", "wb").write(f[name][()].tobytes())
|
||||
print(name, *f[name].shape, flush=True)
|
||||
"#;
|
||||
|
||||
/// h5py's own SWMR reader, run beside ours as the reference: the same
|
||||
/// checks, until `<path>.stop` exists. Prints its iteration count.
|
||||
const H5PY_READER: &str = r#"
|
||||
import os, sys, h5py, numpy as np
|
||||
path = sys.argv[1]
|
||||
f = h5py.File(path, "r", swmr=True)
|
||||
ds = {n: f[n] for n in "abc"}
|
||||
last = {n: (0,) * ds[n].ndim for n in "abc"}
|
||||
its = 0
|
||||
while True:
|
||||
stop = os.path.exists(path + ".stop")
|
||||
for n, d in ds.items():
|
||||
d.refresh()
|
||||
shape = d.shape
|
||||
assert all(s >= l for s, l in zip(shape, last[n])), (n, shape, last[n])
|
||||
last[n] = shape
|
||||
v = d[()]
|
||||
if n == "a":
|
||||
want = np.arange(shape[0]) + 1
|
||||
elif n == "b":
|
||||
want = np.arange(shape[0])[:, None] * 10 + np.arange(4) + 1
|
||||
else:
|
||||
want = np.arange(shape[0])[:, None] * 1000 + np.arange(shape[1]) + 1
|
||||
bad = np.argwhere(v != want)
|
||||
assert len(bad) == 0, (n, shape, bad[:5], v[tuple(bad[0])], want[tuple(bad[0])])
|
||||
its += 1
|
||||
if stop:
|
||||
break
|
||||
print("iterations", its, *last["a"], *last["b"], *last["c"], flush=True)
|
||||
"#;
|
||||
|
||||
fn want_a(n: u64) -> Vec<i64> {
|
||||
(1..=n as i64).collect()
|
||||
}
|
||||
|
||||
fn want_b(rows: std::ops::Range<u64>) -> Vec<f64> {
|
||||
rows.flat_map(|i| (0..4).map(move |j| (10 * i + j + 1) as f64))
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn want_c(rows: std::ops::Range<u64>, cols: u64) -> Vec<i32> {
|
||||
rows.flat_map(|i| (0..cols).map(move |j| (1000 * i + j + 1) as i32))
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// One pass of the Rust reader: refresh every dataset, check its extent
|
||||
/// did not shrink, and check what it reads (the last rows of each, and all
|
||||
/// of them every `full_every` passes).
|
||||
fn check_pass(
|
||||
datasets: &mut [clawhdf5::Dataset<'_>; 3],
|
||||
last: &mut [Vec<u64>; 3],
|
||||
pass: u64,
|
||||
full_every: u64,
|
||||
) {
|
||||
for (k, d) in datasets.iter_mut().enumerate() {
|
||||
d.refresh().unwrap_or_else(|e| panic!("refresh {k}: {e}"));
|
||||
let shape = d.shape().unwrap();
|
||||
assert!(
|
||||
shape.iter().zip(&last[k]).all(|(s, l)| s >= l),
|
||||
"dataset {k} shrank: {shape:?} after {:?}",
|
||||
last[k]
|
||||
);
|
||||
last[k] = shape;
|
||||
}
|
||||
let full = pass.is_multiple_of(full_every);
|
||||
let na = last[0][0];
|
||||
let from = if full { 0 } else { na.saturating_sub(200) };
|
||||
let got = if full {
|
||||
datasets[0].read_i64()
|
||||
} else {
|
||||
datasets[0].read_i64_selection(&Selection::slice(std::slice::from_ref(&(from..na))))
|
||||
}
|
||||
.unwrap_or_else(|e| panic!("read a: {e}"));
|
||||
let want: Vec<i64> = want_a(na).split_off(from as usize);
|
||||
assert!(got == want, "a {from}..{na}: {:?}", first_diff(&got, &want));
|
||||
|
||||
let nb = last[1][0];
|
||||
let from = if full { 0 } else { nb.saturating_sub(50) };
|
||||
let sel = Selection::slice(&[from..nb, 0..4]);
|
||||
let got = datasets[1]
|
||||
.read_f64_selection(&sel)
|
||||
.unwrap_or_else(|e| panic!("read b: {e}"));
|
||||
let want = want_b(from..nb);
|
||||
assert!(got == want, "b {from}..{nb}: {:?}", first_diff(&got, &want));
|
||||
|
||||
let (rc, kc) = (last[2][0], last[2][1]);
|
||||
let from = if full { 0 } else { rc.saturating_sub(9) };
|
||||
let got = if full {
|
||||
datasets[2].read_i32()
|
||||
} else {
|
||||
datasets[2].read_i32_selection(&Selection::slice(&[from..rc, 0..kc]))
|
||||
}
|
||||
.unwrap_or_else(|e| panic!("read c: {e}"));
|
||||
let want = want_c(from..rc, kc);
|
||||
assert!(
|
||||
got == want,
|
||||
"c {from}..{rc} x {kc}: {:?}",
|
||||
first_diff(&got, &want)
|
||||
);
|
||||
}
|
||||
|
||||
fn first_diff<T: PartialEq + std::fmt::Debug>(got: &[T], want: &[T]) -> String {
|
||||
if got.len() != want.len() {
|
||||
return format!("{} values, want {}", got.len(), want.len());
|
||||
}
|
||||
let i = got.iter().zip(want).position(|(g, w)| g != w).unwrap();
|
||||
format!("first difference at {i}: {:?}, want {:?}", got[i], want[i])
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_live_file_reads_consistently_while_an_h5py_swmr_writer_appends() {
|
||||
skip_if_no_python!();
|
||||
let steps: u32 = std::env::var("CLAWHDF5_SWMR_STEPS")
|
||||
.ok()
|
||||
.and_then(|s| s.parse().ok())
|
||||
.unwrap_or(2500);
|
||||
let dir = tempfile::tempdir_in(env!("CARGO_TARGET_TMPDIR")).unwrap();
|
||||
let path = dir.path().join("live.h5");
|
||||
|
||||
let mut writer = Command::new(python())
|
||||
.args(["-c", WRITER, path.to_str().unwrap(), &steps.to_string()])
|
||||
.stdout(Stdio::piped())
|
||||
.spawn()
|
||||
.unwrap();
|
||||
let mut out = BufReader::new(writer.stdout.take().unwrap());
|
||||
let mut line = String::new();
|
||||
out.read_line(&mut line).unwrap();
|
||||
assert_eq!(line.trim(), "ready", "writer did not start");
|
||||
|
||||
// h5py's SWMR reader on the same file, as the reference.
|
||||
let h5py_reader = Command::new(python())
|
||||
.args(["-c", H5PY_READER, path.to_str().unwrap()])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()
|
||||
.unwrap();
|
||||
|
||||
let file = File::open_swmr(&path).unwrap();
|
||||
assert!(file.is_swmr_read());
|
||||
assert!(file.swmr_writer_active().unwrap());
|
||||
// Two reader threads share the file, each with its own handles; each
|
||||
// makes passes until it has seen the writer exit, then one more.
|
||||
let writer_done = AtomicBool::new(false);
|
||||
let reader = || {
|
||||
let mut datasets = ["a", "b", "c"].map(|n| file.dataset(n).unwrap());
|
||||
let mut last = [vec![0], vec![0, 4], vec![0, 1]];
|
||||
let (mut passes, mut grew) = (0u64, 0u64);
|
||||
loop {
|
||||
let done = writer_done.load(Ordering::Acquire);
|
||||
let before = last[0][0];
|
||||
check_pass(&mut datasets, &mut last, passes, 50);
|
||||
passes += 1;
|
||||
grew += u64::from(last[0][0] > before);
|
||||
if done {
|
||||
break;
|
||||
}
|
||||
}
|
||||
// The writer has closed the file: one more pass reads the final
|
||||
// extents, all of them.
|
||||
check_pass(&mut datasets, &mut last, 0, 50);
|
||||
(passes, grew, last)
|
||||
};
|
||||
let (status, results) = std::thread::scope(|s| {
|
||||
let readers = [s.spawn(reader), s.spawn(reader)];
|
||||
let status = writer.wait().unwrap();
|
||||
writer_done.store(true, Ordering::Release);
|
||||
(status, readers.map(|r| r.join().unwrap()))
|
||||
});
|
||||
assert!(status.success(), "writer failed");
|
||||
assert!(!file.swmr_writer_active().unwrap());
|
||||
let last = results[0].2.clone();
|
||||
assert_eq!(results[1].2, last);
|
||||
let passes: Vec<u64> = results.iter().map(|r| r.0).collect();
|
||||
let grew: u64 = results.iter().map(|r| r.1).min().unwrap();
|
||||
|
||||
std::fs::write(path.with_extension("h5.stop"), b"").unwrap();
|
||||
let reference = h5py_reader.wait_with_output().unwrap();
|
||||
assert!(
|
||||
reference.status.success(),
|
||||
"h5py's SWMR reader failed:\n{}",
|
||||
String::from_utf8_lossy(&reference.stderr)
|
||||
);
|
||||
let reference = String::from_utf8_lossy(&reference.stdout).into_owned();
|
||||
|
||||
// The writer's final contents, as h5py reads the closed file.
|
||||
let mut finals = String::new();
|
||||
out.read_to_string(&mut finals).unwrap();
|
||||
let shapes: std::collections::HashMap<&str, Vec<u64>> = finals
|
||||
.lines()
|
||||
.map(|l| {
|
||||
let mut w = l.split_whitespace();
|
||||
let name = w.next().unwrap();
|
||||
(name, w.map(|v| v.parse().unwrap()).collect())
|
||||
})
|
||||
.collect();
|
||||
assert_eq!(shapes["a"], last[0]);
|
||||
assert_eq!(shapes["b"], last[1]);
|
||||
assert_eq!(shapes["c"], last[2]);
|
||||
let bin = |name: &str| std::fs::read(format!("{}.{name}.bin", path.display())).unwrap();
|
||||
let le = |b: &[u8], n: usize| -> Vec<[u8; 8]> {
|
||||
b.chunks(n)
|
||||
.map(|c| {
|
||||
let mut v = [0u8; 8];
|
||||
v[..n].copy_from_slice(c);
|
||||
v
|
||||
})
|
||||
.collect()
|
||||
};
|
||||
// The live handle and a fresh non-SWMR open of the closed file both
|
||||
// read exactly what h5py reads.
|
||||
for f in [&file, &File::open(&path).unwrap()] {
|
||||
let a: Vec<[u8; 8]> = f
|
||||
.dataset("a")
|
||||
.unwrap()
|
||||
.read_i64()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.map(|v| v.to_le_bytes())
|
||||
.collect();
|
||||
assert_eq!(a, le(&bin("a"), 8));
|
||||
let b: Vec<[u8; 8]> = f
|
||||
.dataset("b")
|
||||
.unwrap()
|
||||
.read_f64()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.map(|v| v.to_le_bytes())
|
||||
.collect();
|
||||
assert_eq!(b, le(&bin("b"), 8));
|
||||
let c: Vec<[u8; 8]> = f
|
||||
.dataset("c")
|
||||
.unwrap()
|
||||
.read_i32()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.map(|v| {
|
||||
let mut x = [0u8; 8];
|
||||
x[..4].copy_from_slice(&v.to_le_bytes());
|
||||
x
|
||||
})
|
||||
.collect();
|
||||
assert_eq!(c, le(&bin("c"), 4));
|
||||
}
|
||||
eprintln!(
|
||||
"SWMR: {steps} writer steps; clawhdf5 readers made {passes:?} passes (each saw `a` \
|
||||
grow at least {grew} times), {} retries; h5py reader: {}",
|
||||
file.swmr_retries(),
|
||||
reference.trim()
|
||||
);
|
||||
assert!(
|
||||
grew >= 2,
|
||||
"a reader never saw the file grow ({passes:?} passes)"
|
||||
);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user