diff --git a/crates/clawhdf5/src/lib.rs b/crates/clawhdf5/src/lib.rs index 5ceec87..1d763d9 100644 --- a/crates/clawhdf5/src/lib.rs +++ b/crates/clawhdf5/src/lib.rs @@ -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; diff --git a/crates/clawhdf5/src/reader.rs b/crates/clawhdf5/src/reader.rs index 1786f17..47788e7 100644 --- a/crates/clawhdf5/src/reader.rs +++ b/crates/clawhdf5/src/reader.rs @@ -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, + /// 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 { + 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> = 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, + /// 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>(path: P) -> Result { + 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 { + 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 { + 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(&self, op: impl FnMut() -> Result) -> Result { + 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(&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,20 +634,23 @@ impl File { /// /// The path uses `/` separators (e.g., `"group1/values"`). pub fn dataset(&self, path: &str) -> Result, Error> { - let addr = with_bytes!(self.data.meta()?, |d| group_v2::resolve_path_any_in( - d, - &self.superblock, - path - ))?; - let hdr = self.parse_header(addr)?; - if !has_message(&hdr, MessageType::DataLayout) { - return Err(Error::NotADataset(path.to_string())); - } - Dataset { - file: self, - header: hdr, - } - .check_open() + self.retry(|| { + let addr = with_bytes!(self.data.meta()?, |d| group_v2::resolve_path_any_in( + d, + &self.superblock, + path + ))?; + let hdr = self.parse_header(addr)?; + if !has_message(&hdr, MessageType::DataLayout) { + return Err(Error::NotADataset(path.to_string())); + } + 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, Error> { - let hdr = self.parse_header(address)?; - if !has_message(&hdr, MessageType::DataLayout) { - return Err(Error::NotADataset(format!("object at address {address}"))); - } - Dataset { - file: self, - header: hdr, - } - .check_open() + 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, 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, Error> { - let hdr = self.file.parse_header(self.child_address(name)?)?; - if !has_message(&hdr, MessageType::DataLayout) { - return Err(Error::NotADataset(name.to_string())); - } - Dataset { - file: self.file, - header: hdr, - } - .check_open() + 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,13 +998,15 @@ impl<'f> Group<'f> { /// name index rather than by listing the group (see /// [`group_v2::resolve_child`]). fn child_address(&self, name: &str) -> Result { - with_bytes!(self.file.data.meta()?, |d| group_v2::resolve_child_in( - d, - &self.file.superblock, - self.address, - name - )) - .map_err(Error::Format) + self.file.retry(|| { + with_bytes!(self.file.data.meta()?, |d| group_v2::resolve_child_in( + d, + &self.file.superblock, + self.address, + 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, Error> { - with_bytes!(self.file.data.meta()?, |d| { - group_v2::resolve_group_children_in(d, &self.file.superblock, self.address) + 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) }) - .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, 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, 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( + &self, + selection: &clawhdf5_format::selection::Selection, + convert: fn(&[u8], &Datatype) -> Result, FormatError>, + ) -> Result, 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,20 +1794,28 @@ impl<'f> Dataset<'f> { } let ds = self.dataspace()?; let pipeline = self.filter_pipeline()?; - Ok(data_read::read_chunked_native_in::( - &self.header.messages, - &self.file.data, - &dl, - &ds, - &dt, - pipeline.as_ref(), - self.file.offset_size(), - self.file.length_size(), - Some(&self.file.chunk_cache), - )?) + self.file.retry(|| { + self.file.with_chunk_cache(|cache| { + Ok(data_read::read_chunked_native_in::( + &self.header.messages, + &self.file.data, + &dl, + &ds, + &dt, + pipeline.as_ref(), + self.file.offset_size(), + self.file.length_size(), + Some(cache), + )?) + }) + }) } fn read_raw(&self) -> Result, Error> { + self.file.retry(|| self.read_raw_once()) + } + + fn read_raw_once(&self) -> Result, Error> { let dt = self.datatype()?; let ds = self.dataspace()?; let dl = self.data_layout()?; @@ -1622,16 +1837,18 @@ impl<'f> Dataset<'f> { self.file.offset_size(), self.file.length_size(), || { - Ok(data_read::read_raw_data_cached_in( - &self.file.data, - &dl, - &ds, - &dt, - pipeline.as_ref(), - self.file.offset_size(), - self.file.length_size(), - &self.file.chunk_cache, - )?) + self.file.with_chunk_cache(|cache| { + Ok(data_read::read_raw_data_cached_in( + &self.file.data, + &dl, + &ds, + &dt, + pipeline.as_ref(), + self.file.offset_size(), + self.file.length_size(), + cache, + )?) + }) }, ) } diff --git a/crates/clawhdf5/src/swmr.rs b/crates/clawhdf5/src/swmr.rs new file mode 100644 index 0000000..6d6ba76 --- /dev/null +++ b/crates/clawhdf5/src/swmr.rs @@ -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>(path: P) -> std::io::Result { + 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 { + 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 { + 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 { + 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 { + 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, 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 = 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(&self, op: impl FnMut() -> Result) -> Result { + 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( + attempts: u32, + retried: &AtomicU64, + mut op: impl FnMut() -> Result, +) -> Result { + 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 = 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)); + } +} diff --git a/crates/clawhdf5/tests/swmr_interop.rs b/crates/clawhdf5/tests/swmr_interop.rs index f22e7cd..58b1fac 100644 --- a/crates/clawhdf5/tests/swmr_interop.rs +++ b/crates/clawhdf5/tests/swmr_interop.rs @@ -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 = (1..=148).collect(); + let want_b: Vec = (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 { + 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, + torn: AtomicU32, + invert: AtomicBool, +} + +impl clawhdf5::Storage for Torn { + fn read_at( + &self, + offset: u64, + len: usize, + ) -> Result, 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::>()); + // 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 `..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 `.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 { + (1..=n as i64).collect() +} + +fn want_b(rows: std::ops::Range) -> Vec { + rows.flat_map(|i| (0..4).map(move |j| (10 * i + j + 1) as f64)) + .collect() +} + +fn want_c(rows: std::ops::Range, cols: u64) -> Vec { + 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; 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 = 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(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 = 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> = 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)" + ); +}