Merge branch 'feat/p3-m5-swmr-reader' into feat/p3-wasm-swmr-python
This commit is contained in:
@@ -2,6 +2,76 @@
|
||||
|
||||
## Unreleased
|
||||
|
||||
### Range reads, milestone M5: reading files a SWMR writer is appending to (2026-09-27)
|
||||
Design: `docs/design/swmr.md`.
|
||||
- **Fix: files with the SWMR-write flag were bounded by a stale end of
|
||||
file** (since the end-of-file check of 2026-09-26, unreleased). A
|
||||
libhdf5 SWMR writer (h5py `f.swmr_mode = True`) does not keep
|
||||
the superblock's end-of-file address up to date; a copy h5py made of its
|
||||
own file mid-write records 715 in a 6 030-byte file.
|
||||
`Superblock::data_end` only ignored the recorded end when it lay past the
|
||||
end of the file, so such a file listed, but every chunked read failed
|
||||
("unexpected EOF: need 787 bytes, have 715") and `h5rs check` reported
|
||||
its chunk indexes past the end of the file. For a v3 superblock with the
|
||||
SWMR-write flag the data now ends at the end of the file, as libhdf5's
|
||||
SWMR reader reads it (it skips the end-of-allocation check). Every open
|
||||
path (`File::open`, `open_buffered`, `from_bytes`, `open_storage`,
|
||||
`MmapFile`, `LazyFile`, `h5rs`) reads such a copy now.
|
||||
- **`File::open_swmr(path)` / `File::open_storage_swmr(storage)`**: open a
|
||||
file a SWMR writer may still be appending to, as libhdf5's SWMR reader
|
||||
(h5py `File(path, "r", swmr=True)`) does. The file is read with
|
||||
positioned reads (the new `clawhdf5::FileStorage`: `pread`/`seek_read`,
|
||||
never mapped, `len()` the file's current length), reads are bounded by
|
||||
the file's length at the time of each read, and the chunk cache is not
|
||||
used (a cached chunk index would hide new chunks; a cached edge chunk
|
||||
would read as fill where the writer has since written). A file open for
|
||||
writing without SWMR is refused with `Error::Locked`, as libhdf5 refuses
|
||||
it.
|
||||
Only a file whose superblock has the SWMR-write flag when it is opened is
|
||||
read this way; any other file (its writer has closed it) reads exactly as
|
||||
`File::open` reads it, bounded by its recorded end of file, and
|
||||
`is_swmr_read()` is `false`.
|
||||
- **`Dataset::refresh()`** reads the dataset's object header again
|
||||
(`H5Drefresh`, h5py `Dataset.refresh()`), so `shape()` and later reads
|
||||
see the writer's appends; a handle keeps its extent until refreshed.
|
||||
- **Bounded retries.** On a SWMR-read file, an operation (open, lookups
|
||||
and listings, refresh, every dataset read including strings and
|
||||
variable-length data, attributes, `File::decode_*`) that fails with an error a concurrent write can cause —
|
||||
the failures libhdf5's SWMR reader retries: a checksum mismatch, a read
|
||||
past the file's current end, an object header prefix that does not
|
||||
decode — is run again from the start, up to
|
||||
`File::swmr_read_attempts()` times (`SWMR_READ_ATTEMPTS` = 100,
|
||||
libhdf5's default metadata read attempts for SWMR readers,
|
||||
`H5Pset_metadata_read_attempts`; `set_swmr_read_attempts` changes it),
|
||||
pausing 1 µs doubling to 10 ms between attempts. Results come only from
|
||||
an attempt in which every structure verified, so a torn read is at worst
|
||||
an error, never data. Any other error is returned at once.
|
||||
`File::swmr_retries()` counts the retries.
|
||||
`File::swmr_writer_active()` reads the superblock flags again to tell
|
||||
when the writer has closed the file.
|
||||
- Tests (`crates/clawhdf5/tests/swmr_interop.rs`): the mid-write copy
|
||||
(fixture `tests/fixtures/swmr_mid_write.h5`) through every open path and
|
||||
against h5py's SWMR reader; garbled reads (a storage that corrupts the
|
||||
next reads) retried, never returned, and given up after the attempts;
|
||||
every read path (listings, attributes compact and dense, strings,
|
||||
variable-length data; fixture `tests/fixtures/swmr_strings_attrs.h5`)
|
||||
with each of its reads failing once in turn returns the same result;
|
||||
and a live test: an h5py writer appends to a 1-D and a 2-D dataset with
|
||||
one unlimited dimension (Extensible Array index, one of them gzip) and a
|
||||
2-D dataset with two (v2 B-tree index) for 2 500 steps
|
||||
(`CLAWHDF5_SWMR_STEPS`), flushing after each, while two Rust reader
|
||||
threads refresh and read all three in a loop — every value must be the
|
||||
one the writer wrote at its position, extents never shrink — beside
|
||||
h5py's own SWMR reader doing the same checks; after the writer closes,
|
||||
the live handle and a new `File::open` read exactly what h5py reads.
|
||||
Run on tank, 2026-09-27, with h5py 3.16 / HDF5 2.0:
|
||||
`CLAWHDF5_REQUIRE_INTEROP=1 CLAWHDF5_PYTHON=.venv/bin/python cargo test
|
||||
-p clawhdf5 --test swmr_interop`; also at 20 000 steps in a release
|
||||
build. A test with the chunk cache left on in live mode fails it.
|
||||
- Not covered: SWMR writing, remote SWMR (`BlockCache` caches blocks and
|
||||
`HttpStorage` pins the length), `MmapFile`/`LazyFile`, and refreshing
|
||||
groups or attributes (a SWMR writer cannot add objects or attributes).
|
||||
|
||||
### Range reads, milestone M3: remote files (2026-09-26)
|
||||
- **New crate `clawhdf5-remote`.** `open_url("http://host/file.h5")` gives
|
||||
a `clawhdf5::File` (through `File::open_storage`) that reads the file by
|
||||
|
||||
@@ -101,6 +101,10 @@ breaking change, are in [CHANGELOG.md](CHANGELOG.md).
|
||||
HTTP server (or in S3/GCS/Azure, opt-in) by range requests through a
|
||||
block cache, without downloading it; `h5rs` takes URLs with its `remote`
|
||||
feature. See [Reading remote files](#reading-remote-files).
|
||||
- `File::open_swmr` follows a file an h5py/libhdf5 SWMR writer is still
|
||||
appending to (`Dataset::refresh`, bounded retries); copies of such files
|
||||
taken mid-write read with every open path. See
|
||||
[Following a file a SWMR writer is appending to](#following-a-file-a-swmr-writer-is-appending-to).
|
||||
|
||||
**Tooling**
|
||||
- CI now runs the h5py/netCDF4 interop suites for real (they had been skipping
|
||||
@@ -504,6 +508,43 @@ built with `--features remote` takes the same URLs:
|
||||
`https://` is the `https` feature (rustls with ring, which compiles C).
|
||||
Limits are in [known issues](docs/known-issues.md).
|
||||
|
||||
### Following a file a SWMR writer is appending to
|
||||
|
||||
`File::open_swmr` reads a file that a libhdf5 writer in SWMR mode (h5py
|
||||
`f.swmr_mode = True`) is still appending to, as h5py's
|
||||
`File(path, "r", swmr=True)` does: `Dataset::refresh()` picks up the new
|
||||
extent, every read reads the chunk index as it is now, and a read that
|
||||
races the writer (a checksum that fails mid-flush) is retried, up to 100
|
||||
attempts as in libhdf5, and never returned torn.
|
||||
|
||||
```rust
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
let file = clawhdf5::File::open_swmr("live.h5")?;
|
||||
let mut ds = file.dataset("samples")?;
|
||||
let (mut seen, mut last_growth) = (0, Instant::now());
|
||||
// Stop when the writer closes the file, or when the dataset has not grown
|
||||
// for a minute: a writer that crashed or was killed never clears the
|
||||
// SWMR-write flag, so `swmr_writer_active()` alone can stay true forever.
|
||||
while file.swmr_writer_active()? && last_growth.elapsed() < Duration::from_secs(60) {
|
||||
ds.refresh()?;
|
||||
let n = ds.shape()?[0];
|
||||
if n > seen {
|
||||
// read rows seen..n ...
|
||||
(seen, last_growth) = (n, Instant::now());
|
||||
}
|
||||
std::thread::sleep(Duration::from_millis(100));
|
||||
}
|
||||
ds.refresh()?; // the final extent
|
||||
```
|
||||
|
||||
`swmr_writer_active()` reads the superblock's SWMR-write flag, which libhdf5
|
||||
clears only when the writer closes the file; a file whose writer died keeps
|
||||
it set (as the mid-write copy in `tests/fixtures/swmr_mid_write.h5` does),
|
||||
so a follower needs its own stop condition, like the idle timeout above.
|
||||
|
||||
Design and limits: [docs/design/swmr.md](docs/design/swmr.md).
|
||||
|
||||
### Python
|
||||
|
||||
`crates/clawhdf5-py` is a Python package (PyO3 + numpy) that reads HDF5 with
|
||||
|
||||
@@ -576,7 +576,14 @@ pub fn find_attribute_in_file(
|
||||
offset_size: u8,
|
||||
length_size: u8,
|
||||
) -> Result<Option<AttributeMessage>, FormatError> {
|
||||
find_attribute_core(file_data, header, name, offset_size, length_size)
|
||||
find_attribute_core(
|
||||
file_data,
|
||||
header,
|
||||
name,
|
||||
offset_size,
|
||||
length_size,
|
||||
&mut Vec::new(),
|
||||
)
|
||||
}
|
||||
|
||||
/// [`find_attribute_in_file`] over any [`Storage`] (see
|
||||
@@ -594,16 +601,48 @@ pub fn find_attribute_in<S: Storage + ?Sized>(
|
||||
) -> Result<Option<AttributeMessage>, FormatError> {
|
||||
match file_data.as_contiguous() {
|
||||
Some(all) => find_attribute_in_file(all, header, name, offset_size, length_size),
|
||||
None => find_attribute_core(file_data, header, name, offset_size, length_size),
|
||||
None => find_attribute_core(
|
||||
file_data,
|
||||
header,
|
||||
name,
|
||||
offset_size,
|
||||
length_size,
|
||||
&mut Vec::new(),
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
/// [`find_attribute_in`], also returning the errors of the attributes it
|
||||
/// could not read on the way (which it leaves out rather than failing
|
||||
/// the call): the attribute asked for may be one of them. A reader of a
|
||||
/// file that is being written uses them to tell a read that raced the
|
||||
/// writer from an absent attribute.
|
||||
pub fn find_attribute_reporting_in<S: Storage + ?Sized>(
|
||||
file_data: &S,
|
||||
header: &ObjectHeader,
|
||||
name: &str,
|
||||
offset_size: u8,
|
||||
length_size: u8,
|
||||
) -> Result<(Option<AttributeMessage>, Vec<FormatError>), FormatError> {
|
||||
let mut errors = Vec::new();
|
||||
let found = find_attribute_core(
|
||||
file_data,
|
||||
header,
|
||||
name,
|
||||
offset_size,
|
||||
length_size,
|
||||
&mut errors,
|
||||
)?;
|
||||
Ok((found, errors))
|
||||
}
|
||||
|
||||
fn find_attribute_core<S: Storage + ?Sized>(
|
||||
file_data: &S,
|
||||
header: &ObjectHeader,
|
||||
name: &str,
|
||||
offset_size: u8,
|
||||
length_size: u8,
|
||||
errors: &mut Vec<FormatError>,
|
||||
) -> Result<Option<AttributeMessage>, FormatError> {
|
||||
let attr_info = find_attribute_info(header, offset_size)?;
|
||||
let dense = attr_info
|
||||
@@ -612,12 +651,10 @@ fn find_attribute_core<S: Storage + ?Sized>(
|
||||
let Some((fh_addr, btree_addr)) = dense else {
|
||||
// Compact only (or dense storage without a name index, which a
|
||||
// listing reports): as a listing finds it.
|
||||
return Ok(
|
||||
extract_attributes_tolerant_in(file_data, header, offset_size, length_size)?
|
||||
.0
|
||||
.into_iter()
|
||||
.find(|a| a.name == name),
|
||||
);
|
||||
let (attrs, errs) =
|
||||
extract_attributes_tolerant_in(file_data, header, offset_size, length_size)?;
|
||||
errors.extend(errs);
|
||||
return Ok(attrs.into_iter().find(|a| a.name == name));
|
||||
};
|
||||
let btree_hdr = BTreeV2Header::parse_in(
|
||||
file_data,
|
||||
@@ -627,12 +664,10 @@ fn find_attribute_core<S: Storage + ?Sized>(
|
||||
)?;
|
||||
let fh = FractalHeapHeader::parse_in(file_data, fh_addr, offset_size, length_size)?;
|
||||
if btree_hdr.tree_type != ATTRIBUTE_NAME_INDEX || btree_hdr.record_size < 4 {
|
||||
return Ok(
|
||||
extract_attributes_tolerant_in(file_data, header, offset_size, length_size)?
|
||||
.0
|
||||
.into_iter()
|
||||
.find(|a| a.name == name),
|
||||
);
|
||||
let (attrs, errs) =
|
||||
extract_attributes_tolerant_in(file_data, header, offset_size, length_size)?;
|
||||
errors.extend(errs);
|
||||
return Ok(attrs.into_iter().find(|a| a.name == name));
|
||||
}
|
||||
|
||||
// A listing has the compact attributes first.
|
||||
@@ -671,10 +706,10 @@ fn find_attribute_core<S: Storage + ?Sized>(
|
||||
AttributeMessage::parse_in_storage(&d, file_data, offset_size, length_size)
|
||||
});
|
||||
// One that cannot be read is left out, as from a listing.
|
||||
if let Ok(attr) = attr
|
||||
&& attr.name == name
|
||||
{
|
||||
return Ok(Some(attr));
|
||||
match attr {
|
||||
Ok(attr) if attr.name == name => return Ok(Some(attr)),
|
||||
Ok(_) => {}
|
||||
Err(e) => errors.push(e),
|
||||
}
|
||||
}
|
||||
Ok(None)
|
||||
|
||||
@@ -116,22 +116,27 @@ impl Superblock {
|
||||
/// [`FormatError::TruncatedFile`]. Bytes past that address are not part
|
||||
/// of the file: libhdf5 fails any read of them ("addr overflow" /
|
||||
/// "address plus size exceeds file eoa"), so a reader should parse only
|
||||
/// the data up to the returned end. As libhdf5 does for a SWMR reader,
|
||||
/// the check is skipped for a version-3 superblock whose writer is still
|
||||
/// writing it in SWMR mode (it extends the file as it goes); the data
|
||||
/// then ends at the end of the file.
|
||||
/// the data up to the returned end.
|
||||
///
|
||||
/// A version-3 superblock with the SWMR-write flag set belongs to a file
|
||||
/// a SWMR writer has open (or had, and did not close). That writer does
|
||||
/// not keep the recorded end of file up to date — a copy taken mid-write
|
||||
/// can record an end of a few hundred bytes in a file of tens of
|
||||
/// kilobytes — and libhdf5's SWMR reader skips its end-of-allocation
|
||||
/// check for every read (`H5FD_read`). For such a superblock the data
|
||||
/// ends at the end of the file, whatever end it records.
|
||||
///
|
||||
/// When the superblock's recorded base address differs from where the
|
||||
/// superblock actually is (a user block added or removed after the file
|
||||
/// was written), libhdf5 moves the recorded end of file by the same
|
||||
/// amount, and so does this.
|
||||
pub fn data_end(&self, user_block: u64, file_len: u64) -> Result<u64, FormatError> {
|
||||
let eof =
|
||||
i128::from(self.eof_address) - i128::from(self.base_address) + i128::from(user_block);
|
||||
if eof < 0 || eof > i128::from(file_len) {
|
||||
if self.version >= 3 && self.is_swmr_write() {
|
||||
return Ok(file_len.saturating_sub(user_block));
|
||||
}
|
||||
let eof =
|
||||
i128::from(self.eof_address) - i128::from(self.base_address) + i128::from(user_block);
|
||||
if eof < 0 || eof > i128::from(file_len) {
|
||||
return Err(FormatError::TruncatedFile {
|
||||
stored_eof: u64::try_from(eof).unwrap_or(self.eof_address),
|
||||
actual_len: file_len,
|
||||
@@ -617,6 +622,13 @@ mod tests {
|
||||
let mut swmr = Superblock::parse(&build_v2_bytes(8, 3), 0).unwrap();
|
||||
swmr.consistency_flags = swmr_flags::WRITE_ACCESS | swmr_flags::SWMR_WRITE;
|
||||
assert_eq!(swmr.data_end(0, 1000), Ok(1000));
|
||||
// ... nor bounded by its recorded end, which the writer does not
|
||||
// keep up to date (2048 here).
|
||||
assert_eq!(swmr.data_end(0, 17_857), Ok(17_857));
|
||||
assert_eq!(swmr.data_end(512, 17_857), Ok(17_345));
|
||||
// Without the SWMR-write flag the recorded end bounds the data.
|
||||
swmr.consistency_flags = swmr_flags::WRITE_ACCESS;
|
||||
assert_eq!(swmr.data_end(0, 17_857), Ok(2048));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -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;
|
||||
|
||||
+394
-37
@@ -28,7 +28,9 @@ use clawhdf5_format::superblock_ext::{self, CacheImageState};
|
||||
|
||||
use crate::cache_image::{self, ImageView};
|
||||
use crate::error::Error;
|
||||
use crate::types::{AttrValue, DType, classify_datatype, read_attr, read_attrs};
|
||||
use crate::types::{
|
||||
AttrValue, DType, classify_datatype, read_attr_reporting, read_attrs_reporting,
|
||||
};
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// FileData — internal storage for owned bytes, an mmap, or any Storage
|
||||
@@ -115,6 +117,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 +144,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, &mut None);
|
||||
}
|
||||
let whole = backing.whole_file().unwrap_or_default();
|
||||
let (user_block, hdf5) = signature::split_user_block(whole)?;
|
||||
@@ -173,6 +179,7 @@ impl FileData {
|
||||
overlay: Vec::new(),
|
||||
image_error,
|
||||
contiguous: None,
|
||||
live: false,
|
||||
};
|
||||
data.contiguous = data.find_contiguous();
|
||||
Ok((data, superblock))
|
||||
@@ -180,7 +187,18 @@ impl FileData {
|
||||
|
||||
/// [`Self::new`] for a [`Storage`] backend: the same checks, through
|
||||
/// reads of the storage.
|
||||
fn new_storage(storage: SharedStorage) -> Result<(Self, Superblock), Error> {
|
||||
///
|
||||
/// With `swmr` ([`File::open_storage_swmr`]) the file is read live when
|
||||
/// its superblock has the SWMR-write flag (version 3), as libhdf5's SWMR
|
||||
/// reader reads it: reads end where the file ends at the time of each
|
||||
/// read, and nothing is read as one slice. Any other file is read as
|
||||
/// without `swmr`, bounded by its recorded end of file. Once the
|
||||
/// superblock has been read, `flagged` says which it was.
|
||||
fn new_storage(
|
||||
storage: SharedStorage,
|
||||
swmr: bool,
|
||||
flagged: &mut Option<bool>,
|
||||
) -> Result<(Self, Superblock), Error> {
|
||||
let file_len = storage.len();
|
||||
let base = signature::find_signature_in(&*storage)?;
|
||||
let mut data = Self {
|
||||
@@ -191,11 +209,16 @@ impl FileData {
|
||||
overlay: Vec::new(),
|
||||
image_error: None,
|
||||
contiguous: None,
|
||||
// Until the superblock says otherwise: its own reads are not
|
||||
// bounded by an end of file it has not read yet.
|
||||
live: swmr,
|
||||
};
|
||||
// Worked out again below, once the end of file and any cache image
|
||||
// are known.
|
||||
data.contiguous = data.find_contiguous();
|
||||
let superblock = Superblock::parse_in(&data, 0)?;
|
||||
data.live = swmr && superblock.version >= 3 && superblock.is_swmr_write();
|
||||
*flagged = Some(data.live);
|
||||
data.end = base + superblock.data_end(base, file_len)?;
|
||||
match superblock_ext::cache_image_state_in(&data, &superblock)? {
|
||||
CacheImageState::Absent => {}
|
||||
@@ -230,6 +253,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 +321,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 +337,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 +349,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 +412,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 +434,7 @@ impl File {
|
||||
chunk_cache: ChunkCache::new(),
|
||||
base_dir,
|
||||
vds_resolver: None,
|
||||
swmr: crate::swmr::Retries::default(),
|
||||
})
|
||||
}
|
||||
#[cfg(not(feature = "mmap"))]
|
||||
@@ -428,6 +469,7 @@ impl File {
|
||||
chunk_cache: ChunkCache::new(),
|
||||
base_dir: None,
|
||||
vds_resolver: None,
|
||||
swmr: crate::swmr::Retries::default(),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -463,9 +505,205 @@ 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
|
||||
/// — the failures libhdf5's SWMR reader retries: a checksum mismatch, a
|
||||
/// read past the file's current end, an object header prefix (signature,
|
||||
/// version) that does not decode — 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. Every other error is
|
||||
/// returned at once.
|
||||
///
|
||||
/// All of this applies only to a file whose superblock (version 3) has
|
||||
/// the SWMR-write flag set when it is opened: a file a SWMR writer has
|
||||
/// open, or had and did not close. libhdf5's SWMR writer does not keep
|
||||
/// the superblock's recorded end of file current, so for such a file it
|
||||
/// is ignored, as libhdf5's SWMR reader ignores it. Any other file (one
|
||||
/// whose writer has closed it, or that was never written in SWMR mode)
|
||||
/// is read exactly as [`File::open`] reads it — bounded by its recorded
|
||||
/// end of file, through the chunk cache, each operation tried once —
|
||||
/// and [`is_swmr_read`](Self::is_swmr_read) is `false`. Which of the two
|
||||
/// a handle is does not change after it is opened.
|
||||
///
|
||||
/// 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.
|
||||
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();
|
||||
// Retried only while the superblock cannot be read, or says a SWMR
|
||||
// writer has the file: an error opening any other file is final.
|
||||
let (data, superblock) = swmr.retry(|| {
|
||||
let mut flagged = None;
|
||||
match FileData::new_storage(storage.clone(), true, &mut flagged) {
|
||||
Ok(opened) => Ok(Ok(opened)),
|
||||
Err(e) if flagged == Some(false) => Ok(Err(e)),
|
||||
Err(e) => Err(e),
|
||||
}
|
||||
})??;
|
||||
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 is read live: opened with [`File::open_swmr`] or
|
||||
/// [`File::open_storage_swmr`] while its superblock had the SWMR-write
|
||||
/// flag set. `false` for every other file, including one opened with
|
||||
/// `open_swmr` whose writer had already closed it (it reads as
|
||||
/// [`File::open`] reads it).
|
||||
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 read live
|
||||
/// ([`is_swmr_read`](Self::is_swmr_read) `false`) 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, as far as the file
|
||||
/// says: 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). A writer that crashed or was killed never clears
|
||||
/// it, so the flag alone is not a stop condition: give the loop another
|
||||
/// one, such as a time without growth.
|
||||
///
|
||||
/// ```no_run
|
||||
/// # fn main() -> Result<(), clawhdf5::Error> {
|
||||
/// use std::time::{Duration, Instant};
|
||||
///
|
||||
/// let file = clawhdf5::File::open_swmr("live.h5")?;
|
||||
/// let mut ds = file.dataset("samples")?;
|
||||
/// let (mut seen, mut last_growth) = (0, Instant::now());
|
||||
/// while file.swmr_writer_active()? && last_growth.elapsed() < Duration::from_secs(60) {
|
||||
/// ds.refresh()?;
|
||||
/// let n = ds.shape()?[0];
|
||||
/// if n > seen {
|
||||
/// // read rows seen..n ...
|
||||
/// (seen, last_growth) = (n, Instant::now());
|
||||
/// }
|
||||
/// std::thread::sleep(Duration::from_millis(100));
|
||||
/// }
|
||||
/// ds.refresh()?; // the final extent
|
||||
/// # Ok(())
|
||||
/// # }
|
||||
/// ```
|
||||
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()
|
||||
}
|
||||
}
|
||||
|
||||
/// An attribute read under [`Self::retry`]. Attribute reads leave out
|
||||
/// (or return as [`AttrValue::Raw`]) an attribute they cannot read
|
||||
/// rather than fail, so `op` also returns the errors behind those: on a
|
||||
/// live file one a concurrent write can cause runs the read again too.
|
||||
/// If every attempt meets one, the last attempt's result is returned,
|
||||
/// as it is for a file that is not live.
|
||||
fn retry_attrs<T>(
|
||||
&self,
|
||||
mut op: impl FnMut() -> Result<(T, Vec<FormatError>), Error>,
|
||||
) -> Result<T, Error> {
|
||||
let mut last = None;
|
||||
let result = self.retry(|| {
|
||||
last = None;
|
||||
let (value, errors) = op()?;
|
||||
match errors.into_iter().find(crate::swmr::is_transient_format) {
|
||||
Some(e) if self.data.live => {
|
||||
last = Some(value);
|
||||
Err(Error::Format(e))
|
||||
}
|
||||
_ => Ok(value),
|
||||
}
|
||||
});
|
||||
match (result, last) {
|
||||
(Err(_), Some(value)) => Ok(value),
|
||||
(result, _) => result,
|
||||
}
|
||||
}
|
||||
|
||||
/// 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 +726,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 +738,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 +752,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 +782,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,
|
||||
@@ -665,6 +909,7 @@ impl File {
|
||||
/// [`AttrValue::Raw`] attribute. Variable-length strings are resolved in
|
||||
/// this file's global heap; see [`Dataset::read_string`] for the values.
|
||||
pub fn decode_strings(&self, datatype: &Datatype, raw: &[u8]) -> Result<Vec<String>, Error> {
|
||||
self.retry(|| {
|
||||
with_bytes!(&self.data, |d| crate::vlen::decode_strings(
|
||||
d,
|
||||
datatype,
|
||||
@@ -672,6 +917,7 @@ impl File {
|
||||
self.offset_size(),
|
||||
self.length_size(),
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
/// Like [`decode_strings`](Self::decode_strings) for variable-length
|
||||
@@ -682,6 +928,7 @@ impl File {
|
||||
datatype: &Datatype,
|
||||
raw: &[u8],
|
||||
) -> Result<Vec<Vec<u8>>, Error> {
|
||||
self.retry(|| {
|
||||
with_bytes!(self.data.meta()?, |d| crate::vlen::decode_string_bytes(
|
||||
d,
|
||||
datatype,
|
||||
@@ -689,6 +936,7 @@ impl File {
|
||||
self.offset_size(),
|
||||
self.length_size(),
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
/// Decode the variable-length sequences in `raw`, a buffer of elements
|
||||
@@ -700,6 +948,7 @@ impl File {
|
||||
datatype: &Datatype,
|
||||
raw: &[u8],
|
||||
) -> Result<Vec<Vec<T>>, Error> {
|
||||
self.retry(|| {
|
||||
with_bytes!(self.data.meta()?, |d| crate::vlen::decode_vlen(
|
||||
d,
|
||||
datatype,
|
||||
@@ -707,6 +956,7 @@ impl File {
|
||||
self.offset_size(),
|
||||
self.length_size(),
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
fn parse_header(&self, address: u64) -> Result<ObjectHeader, FormatError> {
|
||||
@@ -750,6 +1000,7 @@ pub struct Group<'f> {
|
||||
impl<'f> Group<'f> {
|
||||
/// List the names of datasets in this group.
|
||||
pub fn datasets(&self) -> Result<Vec<String>, Error> {
|
||||
self.file.retry(|| {
|
||||
let entries = self.children()?;
|
||||
let mut names = Vec::new();
|
||||
for entry in &entries {
|
||||
@@ -759,10 +1010,12 @@ impl<'f> Group<'f> {
|
||||
}
|
||||
}
|
||||
Ok(names)
|
||||
})
|
||||
}
|
||||
|
||||
/// List the names of subgroups in this group.
|
||||
pub fn groups(&self) -> Result<Vec<String>, Error> {
|
||||
self.file.retry(|| {
|
||||
let entries = self.children()?;
|
||||
let mut names = Vec::new();
|
||||
for entry in &entries {
|
||||
@@ -772,6 +1025,7 @@ impl<'f> Group<'f> {
|
||||
}
|
||||
}
|
||||
Ok(names)
|
||||
})
|
||||
}
|
||||
|
||||
/// Read all attributes of this group.
|
||||
@@ -792,26 +1046,31 @@ impl<'f> Group<'f> {
|
||||
pub fn attrs_with_errors(
|
||||
&self,
|
||||
) -> Result<(HashMap<String, AttrValue>, Vec<FormatError>), Error> {
|
||||
self.file.retry_attrs(|| {
|
||||
let hdr = self.file.parse_header(self.address)?;
|
||||
with_bytes!(&self.file.data, |d| read_attrs(
|
||||
d,
|
||||
&hdr,
|
||||
self.file.offset_size(),
|
||||
self.file.length_size()
|
||||
))
|
||||
let (attrs, errors, read_errors) = with_bytes!(&self.file.data, |d| {
|
||||
read_attrs_reporting(d, &hdr, self.file.offset_size(), self.file.length_size())
|
||||
})?;
|
||||
let transient = errors.iter().chain(&read_errors).cloned().collect();
|
||||
Ok(((attrs, errors), transient))
|
||||
})
|
||||
}
|
||||
|
||||
/// 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.
|
||||
@@ -827,14 +1086,16 @@ impl<'f> Group<'f> {
|
||||
/// that name, found without reading the other attributes when they are
|
||||
/// stored densely.
|
||||
pub fn attr(&self, name: &str) -> Result<Option<AttrValue>, Error> {
|
||||
self.file.retry_attrs(|| {
|
||||
let hdr = self.file.parse_header(self.address)?;
|
||||
with_bytes!(&self.file.data, |d| read_attr(
|
||||
with_bytes!(&self.file.data, |d| read_attr_reporting(
|
||||
d,
|
||||
&hdr,
|
||||
name,
|
||||
self.file.offset_size(),
|
||||
self.file.length_size(),
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
/// The object header address of the child called `name`: the entry of
|
||||
@@ -842,6 +1103,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 +1111,7 @@ impl<'f> Group<'f> {
|
||||
name
|
||||
))
|
||||
.map_err(Error::Format)
|
||||
})
|
||||
}
|
||||
|
||||
/// This group's children that can be opened, as `(name, object header
|
||||
@@ -869,10 +1132,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 +1149,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 +1173,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()?;
|
||||
@@ -938,9 +1230,11 @@ impl<'f> Dataset<'f> {
|
||||
|
||||
/// Read all data as `f64` values.
|
||||
pub fn read_f64(&self) -> Result<Vec<f64>, Error> {
|
||||
self.file.retry(|| {
|
||||
let dt = self.datatype()?;
|
||||
// A contiguous dataset is converted straight from the file bytes; going
|
||||
// through `read_raw` first copied the whole dataset an extra time.
|
||||
// A contiguous dataset is converted straight from the file bytes;
|
||||
// going through `read_raw` first copied the whole dataset an
|
||||
// extra time.
|
||||
if let Ok(Some(bytes)) = self.contiguous_raw() {
|
||||
return Ok(data_read::read_as_f64(&bytes, &dt)?);
|
||||
}
|
||||
@@ -949,6 +1243,7 @@ impl<'f> Dataset<'f> {
|
||||
}
|
||||
let raw = self.read_raw()?;
|
||||
Ok(data_read::read_as_f64(&raw, &dt)?)
|
||||
})
|
||||
}
|
||||
|
||||
/// Zero-copy read of contiguous native-endian `f64` data.
|
||||
@@ -958,9 +1253,11 @@ impl<'f> Dataset<'f> {
|
||||
///
|
||||
/// Read all data as `f32` values.
|
||||
pub fn read_f32(&self) -> Result<Vec<f32>, Error> {
|
||||
self.file.retry(|| {
|
||||
let dt = self.datatype()?;
|
||||
// A contiguous dataset is converted straight from the file bytes; going
|
||||
// through `read_raw` first copied the whole dataset an extra time.
|
||||
// A contiguous dataset is converted straight from the file bytes;
|
||||
// going through `read_raw` first copied the whole dataset an
|
||||
// extra time.
|
||||
if let Ok(Some(bytes)) = self.contiguous_raw() {
|
||||
return Ok(data_read::read_as_f32(&bytes, &dt)?);
|
||||
}
|
||||
@@ -969,13 +1266,16 @@ impl<'f> Dataset<'f> {
|
||||
}
|
||||
let raw = self.read_raw()?;
|
||||
Ok(data_read::read_as_f32(&raw, &dt)?)
|
||||
})
|
||||
}
|
||||
|
||||
/// Read all data as `i32` values.
|
||||
pub fn read_i32(&self) -> Result<Vec<i32>, Error> {
|
||||
self.file.retry(|| {
|
||||
let dt = self.datatype()?;
|
||||
// A contiguous dataset is converted straight from the file bytes; going
|
||||
// through `read_raw` first copied the whole dataset an extra time.
|
||||
// A contiguous dataset is converted straight from the file bytes;
|
||||
// going through `read_raw` first copied the whole dataset an
|
||||
// extra time.
|
||||
if let Ok(Some(bytes)) = self.contiguous_raw() {
|
||||
return Ok(data_read::read_as_i32(&bytes, &dt)?);
|
||||
}
|
||||
@@ -984,13 +1284,16 @@ impl<'f> Dataset<'f> {
|
||||
}
|
||||
let raw = self.read_raw()?;
|
||||
Ok(data_read::read_as_i32(&raw, &dt)?)
|
||||
})
|
||||
}
|
||||
|
||||
/// Read all data as `i64` values.
|
||||
pub fn read_i64(&self) -> Result<Vec<i64>, Error> {
|
||||
self.file.retry(|| {
|
||||
let dt = self.datatype()?;
|
||||
// A contiguous dataset is converted straight from the file bytes; going
|
||||
// through `read_raw` first copied the whole dataset an extra time.
|
||||
// A contiguous dataset is converted straight from the file bytes;
|
||||
// going through `read_raw` first copied the whole dataset an
|
||||
// extra time.
|
||||
if let Ok(Some(bytes)) = self.contiguous_raw() {
|
||||
return Ok(data_read::read_as_i64(&bytes, &dt)?);
|
||||
}
|
||||
@@ -999,13 +1302,16 @@ impl<'f> Dataset<'f> {
|
||||
}
|
||||
let raw = self.read_raw()?;
|
||||
Ok(data_read::read_as_i64(&raw, &dt)?)
|
||||
})
|
||||
}
|
||||
|
||||
/// Read all data as `u64` values.
|
||||
pub fn read_u64(&self) -> Result<Vec<u64>, Error> {
|
||||
self.file.retry(|| {
|
||||
let dt = self.datatype()?;
|
||||
// A contiguous dataset is converted straight from the file bytes; going
|
||||
// through `read_raw` first copied the whole dataset an extra time.
|
||||
// A contiguous dataset is converted straight from the file bytes;
|
||||
// going through `read_raw` first copied the whole dataset an
|
||||
// extra time.
|
||||
if let Ok(Some(bytes)) = self.contiguous_raw() {
|
||||
return Ok(data_read::read_as_u64(&bytes, &dt)?);
|
||||
}
|
||||
@@ -1014,6 +1320,7 @@ impl<'f> Dataset<'f> {
|
||||
}
|
||||
let raw = self.read_raw()?;
|
||||
Ok(data_read::read_as_u64(&raw, &dt)?)
|
||||
})
|
||||
}
|
||||
|
||||
/// Read all data as `String` values, in row-major order.
|
||||
@@ -1024,17 +1331,21 @@ impl<'f> Dataset<'f> {
|
||||
/// them; bytes that are not valid UTF-8 are replaced with U+FFFD — use
|
||||
/// [`read_string_bytes`](Self::read_string_bytes) for the exact bytes.
|
||||
pub fn read_string(&self) -> Result<Vec<String>, Error> {
|
||||
self.file.retry(|| {
|
||||
let raw = self.read_raw()?;
|
||||
let dt = self.datatype()?;
|
||||
self.file.decode_strings(&dt, &raw)
|
||||
})
|
||||
}
|
||||
|
||||
/// Read a variable-length string dataset as the exact bytes of each
|
||||
/// string (what h5py's `Dataset[()]` returns), in row-major order.
|
||||
pub fn read_string_bytes(&self) -> Result<Vec<Vec<u8>>, Error> {
|
||||
self.file.retry(|| {
|
||||
let raw = self.read_raw()?;
|
||||
let dt = self.datatype()?;
|
||||
self.file.decode_string_bytes(&dt, &raw)
|
||||
})
|
||||
}
|
||||
|
||||
/// Read the selected elements of a fixed- or variable-length string
|
||||
@@ -1043,9 +1354,11 @@ impl<'f> Dataset<'f> {
|
||||
&self,
|
||||
selection: &clawhdf5_format::selection::Selection,
|
||||
) -> Result<Vec<String>, Error> {
|
||||
self.file.retry(|| {
|
||||
let raw = self.read_selection(selection)?;
|
||||
let dt = self.datatype()?;
|
||||
self.file.decode_strings(&dt, &raw)
|
||||
})
|
||||
}
|
||||
|
||||
/// Read a variable-length sequence dataset (h5py
|
||||
@@ -1054,9 +1367,11 @@ impl<'f> Dataset<'f> {
|
||||
/// is converted to `T` as [`read_f64`](Self::read_f64) and the other
|
||||
/// typed readers convert. A null element is an empty sequence.
|
||||
pub fn read_vlen<T: crate::vlen::VlenValue>(&self) -> Result<Vec<Vec<T>>, Error> {
|
||||
self.file.retry(|| {
|
||||
let raw = self.read_raw()?;
|
||||
let dt = self.datatype()?;
|
||||
self.file.decode_vlen(&dt, &raw)
|
||||
})
|
||||
}
|
||||
|
||||
/// Read the selected elements of a variable-length sequence dataset
|
||||
@@ -1065,9 +1380,11 @@ impl<'f> Dataset<'f> {
|
||||
&self,
|
||||
selection: &clawhdf5_format::selection::Selection,
|
||||
) -> Result<Vec<Vec<T>>, Error> {
|
||||
self.file.retry(|| {
|
||||
let raw = self.read_selection(selection)?;
|
||||
let dt = self.datatype()?;
|
||||
self.file.decode_vlen(&dt, &raw)
|
||||
})
|
||||
}
|
||||
|
||||
// ----- Selection-based read methods -----
|
||||
@@ -1096,6 +1413,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 +1508,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() {
|
||||
@@ -1445,12 +1782,18 @@ impl<'f> Dataset<'f> {
|
||||
pub fn attrs_with_errors(
|
||||
&self,
|
||||
) -> Result<(HashMap<String, AttrValue>, Vec<FormatError>), Error> {
|
||||
with_bytes!(&self.file.data, |d| read_attrs(
|
||||
self.file.retry_attrs(|| {
|
||||
let (attrs, errors, read_errors) = with_bytes!(&self.file.data, |d| {
|
||||
read_attrs_reporting(
|
||||
d,
|
||||
&self.header,
|
||||
self.file.offset_size(),
|
||||
self.file.length_size(),
|
||||
))
|
||||
)
|
||||
})?;
|
||||
let transient = errors.iter().chain(&read_errors).cloned().collect();
|
||||
Ok(((attrs, errors), transient))
|
||||
})
|
||||
}
|
||||
|
||||
/// The attribute called `name`, or `None` if it has none by that name
|
||||
@@ -1458,13 +1801,15 @@ impl<'f> Dataset<'f> {
|
||||
/// that name, found without reading the other attributes when they are
|
||||
/// stored densely.
|
||||
pub fn attr(&self, name: &str) -> Result<Option<AttrValue>, Error> {
|
||||
with_bytes!(&self.file.data, |d| read_attr(
|
||||
self.file.retry_attrs(|| {
|
||||
with_bytes!(&self.file.data, |d| read_attr_reporting(
|
||||
d,
|
||||
&self.header,
|
||||
name,
|
||||
self.file.offset_size(),
|
||||
self.file.length_size(),
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
/// Verify this dataset's content against its stored provenance hash
|
||||
@@ -1484,12 +1829,14 @@ impl<'f> Dataset<'f> {
|
||||
/// result is not a tamper-evidence or authenticity guarantee.
|
||||
#[cfg(feature = "provenance")]
|
||||
pub fn verify_provenance(&self) -> Result<clawhdf5_format::provenance::VerifyResult, Error> {
|
||||
self.file.retry(|| {
|
||||
Ok(clawhdf5_format::provenance::verify_dataset_in(
|
||||
&self.file.data,
|
||||
&self.header,
|
||||
self.file.offset_size(),
|
||||
self.file.length_size(),
|
||||
)?)
|
||||
})
|
||||
}
|
||||
|
||||
/// A header message's payload, resolved through the shared-message
|
||||
@@ -1499,11 +1846,10 @@ impl<'f> Dataset<'f> {
|
||||
&self,
|
||||
msg_type: MessageType,
|
||||
) -> Result<Option<std::borrow::Cow<'_, [u8]>>, Error> {
|
||||
self.header
|
||||
.messages
|
||||
.iter()
|
||||
.find(|m| m.msg_type == msg_type)
|
||||
.map(|msg| {
|
||||
let msg = self.header.messages.iter().find(|m| m.msg_type == msg_type);
|
||||
// A shared message is read from the header it lives in.
|
||||
self.file.retry(|| {
|
||||
msg.map(|msg| {
|
||||
with_bytes!(&self.file.data, |d| {
|
||||
clawhdf5_format::shared_message::message_data_in(
|
||||
d,
|
||||
@@ -1515,6 +1861,7 @@ impl<'f> Dataset<'f> {
|
||||
.map_err(Error::Format)
|
||||
})
|
||||
.transpose()
|
||||
})
|
||||
}
|
||||
|
||||
fn required_payload(&self, msg_type: MessageType) -> Result<std::borrow::Cow<'_, [u8]>, Error> {
|
||||
@@ -1587,6 +1934,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 +1945,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 +1977,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 +1986,9 @@ impl<'f> Dataset<'f> {
|
||||
pipeline.as_ref(),
|
||||
self.file.offset_size(),
|
||||
self.file.length_size(),
|
||||
&self.file.chunk_cache,
|
||||
cache,
|
||||
)?)
|
||||
})
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,362 @@
|
||||
//! 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 it or has not yet written it, so that the operation is worth
|
||||
/// running again (see [`is_transient_format`]).
|
||||
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,
|
||||
}
|
||||
}
|
||||
|
||||
/// The failures libhdf5's SWMR reader retries, and nothing else. libhdf5
|
||||
/// (`H5C__load_entry`) reads a metadata structure again when its checksum
|
||||
/// fails, and when the prefix it decodes before the checksum to learn the
|
||||
/// structure's size does not decode (for an object header, its signature
|
||||
/// and version: a header whose every byte is garbled fails there). A read
|
||||
/// past the file's current end is short here; libhdf5 reads zeros there,
|
||||
/// which then fail the checksum. So:
|
||||
///
|
||||
/// - [`FormatError::ChecksumMismatch`], of any checksummed structure;
|
||||
/// - [`FormatError::UnexpectedEof`], a read past the current end;
|
||||
/// - [`FormatError::InvalidObjectHeaderSignature`] and
|
||||
/// [`FormatError::InvalidObjectHeaderVersion`], the object header prefix.
|
||||
///
|
||||
/// Every other error — an unsupported version or message, a file that is not
|
||||
/// HDF5, a structure that is corrupt behind a valid checksum — is returned
|
||||
/// at once: a concurrent write does not cause it, and retrying it only
|
||||
/// costs up to a second of pauses.
|
||||
pub(crate) fn is_transient_format(e: &FormatError) -> bool {
|
||||
use FormatError as F;
|
||||
matches!(
|
||||
e,
|
||||
F::ChecksumMismatch { .. }
|
||||
| F::UnexpectedEof { .. }
|
||||
| F::InvalidObjectHeaderSignature
|
||||
| F::InvalidObjectHeaderVersion(_)
|
||||
)
|
||||
}
|
||||
|
||||
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 permanent_errors_are_returned_at_once() {
|
||||
// Errors a concurrent write does not cause: none is retried, so
|
||||
// none costs the pauses (about 0.9 s for 100 attempts).
|
||||
let permanent = [
|
||||
FormatError::SignatureNotFound,
|
||||
FormatError::UnsupportedVersion(9),
|
||||
FormatError::UnsupportedMessage(0x99),
|
||||
FormatError::TruncatedFile {
|
||||
stored_eof: 10,
|
||||
actual_len: 5,
|
||||
},
|
||||
FormatError::InvalidDatatypeClass(15),
|
||||
FormatError::InvalidLayoutVersion(9),
|
||||
FormatError::InvalidBTreeSignature,
|
||||
FormatError::ChunkedReadError("x".into()),
|
||||
FormatError::DecompressionError("x".into()),
|
||||
FormatError::Fletcher32Mismatch {
|
||||
expected: 1,
|
||||
computed: 2,
|
||||
},
|
||||
];
|
||||
let n = AtomicU64::new(0);
|
||||
let started = std::time::Instant::now();
|
||||
for e in permanent {
|
||||
assert!(!is_transient_format(&e), "{e:?}");
|
||||
let mut calls = 0;
|
||||
let r: Result<(), Error> = retry(SWMR_READ_ATTEMPTS, &n, || {
|
||||
calls += 1;
|
||||
Err(Error::Format(e.clone()))
|
||||
});
|
||||
assert!(r.is_err());
|
||||
assert_eq!(calls, 1, "{e:?}");
|
||||
}
|
||||
assert_eq!(n.load(Ordering::Relaxed), 0);
|
||||
assert!(started.elapsed() < Duration::from_millis(50));
|
||||
assert!(!is_transient(&Error::Io(std::io::Error::other("x"))));
|
||||
|
||||
// The transient ones are retried to the limit.
|
||||
for e in [
|
||||
FormatError::ChecksumMismatch {
|
||||
expected: 1,
|
||||
computed: 2,
|
||||
},
|
||||
FormatError::UnexpectedEof {
|
||||
expected: 8,
|
||||
available: 0,
|
||||
},
|
||||
FormatError::InvalidObjectHeaderSignature,
|
||||
FormatError::InvalidObjectHeaderVersion(0x4f),
|
||||
] {
|
||||
let mut calls = 0;
|
||||
let _: Result<(), Error> = retry(3, &n, || {
|
||||
calls += 1;
|
||||
Err(Error::Format(e.clone()))
|
||||
});
|
||||
assert_eq!(calls, 3, "{e:?}");
|
||||
}
|
||||
assert!(is_transient(&Error::Io(
|
||||
std::io::ErrorKind::UnexpectedEof.into()
|
||||
)));
|
||||
}
|
||||
|
||||
#[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::InvalidObjectHeaderVersion(0x76)))
|
||||
})
|
||||
});
|
||||
assert!(r.is_err());
|
||||
assert_eq!((outer, inner), (3, 3));
|
||||
// The flag is reset afterwards.
|
||||
assert!(!RETRYING.with(Cell::get));
|
||||
}
|
||||
}
|
||||
@@ -170,16 +170,38 @@ pub(crate) fn read_attrs<S: clawhdf5_format::storage::Storage + ?Sized>(
|
||||
),
|
||||
crate::Error,
|
||||
> {
|
||||
read_attrs_reporting(file_data, header, offset_size, length_size)
|
||||
.map(|(attrs, errors, _)| (attrs, errors))
|
||||
}
|
||||
|
||||
/// What [`read_attrs_reporting`] returns: the attributes, the errors of
|
||||
/// those left out, and the errors behind values returned as
|
||||
/// [`AttrValue::Raw`].
|
||||
pub(crate) type AttrsReport = (
|
||||
HashMap<String, AttrValue>,
|
||||
Vec<clawhdf5_format::error::FormatError>,
|
||||
Vec<clawhdf5_format::error::FormatError>,
|
||||
);
|
||||
|
||||
/// [`read_attrs`], also returning the errors of reading the file for a
|
||||
/// value (a variable-length string's global heap) that was returned as
|
||||
/// [`AttrValue::Raw`] instead: a live file retries on them (see
|
||||
/// `File::open_swmr`).
|
||||
pub(crate) fn read_attrs_reporting<S: clawhdf5_format::storage::Storage + ?Sized>(
|
||||
file_data: &S,
|
||||
header: &clawhdf5_format::object_header::ObjectHeader,
|
||||
offset_size: u8,
|
||||
length_size: u8,
|
||||
) -> Result<AttrsReport, crate::Error> {
|
||||
let (msgs, errors) = clawhdf5_format::attribute::extract_attributes_tolerant_in(
|
||||
file_data,
|
||||
header,
|
||||
offset_size,
|
||||
length_size,
|
||||
)?;
|
||||
Ok((
|
||||
attrs_to_map(&msgs, file_data, offset_size, length_size),
|
||||
errors,
|
||||
))
|
||||
let mut read_errors = Vec::new();
|
||||
let map = attrs_to_map_reporting(&msgs, file_data, offset_size, length_size, &mut read_errors);
|
||||
Ok((map, errors, read_errors))
|
||||
}
|
||||
|
||||
/// The attribute called `name` on the object with header `header`, decoded
|
||||
@@ -192,30 +214,50 @@ pub(crate) fn read_attr<S: clawhdf5_format::storage::Storage + ?Sized>(
|
||||
offset_size: u8,
|
||||
length_size: u8,
|
||||
) -> Result<Option<AttrValue>, crate::Error> {
|
||||
let Some(msg) = clawhdf5_format::attribute::find_attribute_in(
|
||||
read_attr_reporting(file_data, header, name, offset_size, length_size).map(|(v, _)| v)
|
||||
}
|
||||
|
||||
/// [`read_attr`], also returning the errors of the attributes it could not
|
||||
/// read on the way (the one asked for may be among them) and of reading the
|
||||
/// file for its value, if one made it [`AttrValue::Raw`] (see
|
||||
/// [`read_attrs_reporting`]).
|
||||
pub(crate) fn read_attr_reporting<S: clawhdf5_format::storage::Storage + ?Sized>(
|
||||
file_data: &S,
|
||||
header: &clawhdf5_format::object_header::ObjectHeader,
|
||||
name: &str,
|
||||
offset_size: u8,
|
||||
length_size: u8,
|
||||
) -> Result<(Option<AttrValue>, Vec<clawhdf5_format::error::FormatError>), crate::Error> {
|
||||
let (found, mut read_errors) = clawhdf5_format::attribute::find_attribute_reporting_in(
|
||||
file_data,
|
||||
header,
|
||||
name,
|
||||
offset_size,
|
||||
length_size,
|
||||
)?
|
||||
else {
|
||||
return Ok(None);
|
||||
)?;
|
||||
let Some(msg) = found else {
|
||||
return Ok((None, read_errors));
|
||||
};
|
||||
Ok(attrs_to_map(
|
||||
let value = attrs_to_map_reporting(
|
||||
std::slice::from_ref(&msg),
|
||||
file_data,
|
||||
offset_size,
|
||||
length_size,
|
||||
&mut read_errors,
|
||||
)
|
||||
.remove(name))
|
||||
.remove(name);
|
||||
Ok((value, read_errors))
|
||||
}
|
||||
|
||||
pub(crate) fn attrs_to_map<S: clawhdf5_format::storage::Storage + ?Sized>(
|
||||
/// The attributes `attrs` by name, decoded; an error reading the file for
|
||||
/// a value that is returned as [`AttrValue::Raw`] because of it is added
|
||||
/// to `read_errors`.
|
||||
pub(crate) fn attrs_to_map_reporting<S: clawhdf5_format::storage::Storage + ?Sized>(
|
||||
attrs: &[clawhdf5_format::attribute::AttributeMessage],
|
||||
file_data: &S,
|
||||
offset_size: u8,
|
||||
length_size: u8,
|
||||
read_errors: &mut Vec<clawhdf5_format::error::FormatError>,
|
||||
) -> HashMap<String, AttrValue> {
|
||||
let mut map = HashMap::new();
|
||||
for attr in attrs {
|
||||
@@ -224,13 +266,11 @@ pub(crate) fn attrs_to_map<S: clawhdf5_format::storage::Storage + ?Sized>(
|
||||
// verbatim as `AttrValue::Raw` rather than dropped — a partial
|
||||
// attribute list with no indication anything is missing is worse than
|
||||
// an undecoded value.
|
||||
let val =
|
||||
decode_attr_value(attr, file_data, offset_size, length_size).unwrap_or_else(|| {
|
||||
AttrValue::Raw {
|
||||
let val = decode_attr_value(attr, file_data, offset_size, length_size, read_errors)
|
||||
.unwrap_or_else(|| AttrValue::Raw {
|
||||
datatype: attr.datatype.clone(),
|
||||
shape: attr.dataspace.dimensions.clone(),
|
||||
data: attr.raw_data.clone(),
|
||||
}
|
||||
});
|
||||
map.insert(attr.name.clone(), val);
|
||||
}
|
||||
@@ -271,6 +311,7 @@ fn decode_attr_value<S: clawhdf5_format::storage::Storage + ?Sized>(
|
||||
file_data: &S,
|
||||
offset_size: u8,
|
||||
length_size: u8,
|
||||
read_errors: &mut Vec<clawhdf5_format::error::FormatError>,
|
||||
) -> Option<AttrValue> {
|
||||
use clawhdf5_format::datatype::Datatype;
|
||||
|
||||
@@ -310,9 +351,13 @@ fn decode_attr_value<S: clawhdf5_format::storage::Storage + ?Sized>(
|
||||
Datatype::VariableLength {
|
||||
is_string: true, ..
|
||||
} => {
|
||||
let strings = attr
|
||||
.read_vl_strings_in(file_data, offset_size, length_size)
|
||||
.ok()?;
|
||||
let strings = match attr.read_vl_strings_in(file_data, offset_size, length_size) {
|
||||
Ok(strings) => strings,
|
||||
Err(e) => {
|
||||
read_errors.push(e);
|
||||
return None;
|
||||
}
|
||||
};
|
||||
if strings.len() == 1 {
|
||||
Some(AttrValue::String(strings[0].clone()))
|
||||
} else {
|
||||
|
||||
Binary file not shown.
Binary file not shown.
@@ -0,0 +1,847 @@
|
||||
//! Files a libhdf5 SWMR writer (h5py `f.swmr_mode = True`) has open: a copy
|
||||
//! taken mid-write (fixture), and a live file appended to by an h5py writer
|
||||
//! process while clawhdf5 and h5py's own SWMR reader read it (see
|
||||
//! `docs/design/swmr.md`).
|
||||
//!
|
||||
//! 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, Stdio};
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
|
||||
|
||||
use clawhdf5::{File, Selection};
|
||||
|
||||
fn python() -> String {
|
||||
std::env::var("CLAWHDF5_PYTHON").unwrap_or_else(|_| "python3".to_string())
|
||||
}
|
||||
|
||||
fn interop_required() -> bool {
|
||||
std::env::var("CLAWHDF5_REQUIRE_INTEROP").is_ok_and(|v| v == "1")
|
||||
}
|
||||
|
||||
fn python_available() -> bool {
|
||||
Command::new(python())
|
||||
.args(["-c", "import h5py, numpy"])
|
||||
.output()
|
||||
.map(|o| o.status.success())
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
macro_rules! skip_if_no_python {
|
||||
() => {
|
||||
if !python_available() {
|
||||
assert!(
|
||||
!interop_required(),
|
||||
"CLAWHDF5_REQUIRE_INTEROP=1 but python3 with h5py is not available"
|
||||
);
|
||||
eprintln!("SKIP: python3 with h5py not available");
|
||||
return;
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
fn fixture(name: &str) -> PathBuf {
|
||||
Path::new(env!("CARGO_MANIFEST_DIR"))
|
||||
.join("tests/fixtures")
|
||||
.join(name)
|
||||
}
|
||||
|
||||
/// `swmr_mid_write.h5`: a copy h5py 3.16 (HDF5 2.0) made of its own file
|
||||
/// while writing it in SWMR mode, after 4 appends of 37 rows to
|
||||
/// `/a` (int64, chunks of 100, no filter: `a[i] = i + 1`) and `/b` (float64
|
||||
/// `(n, 4)`, chunks of 16 x 4, gzip: `b[i, j] = 10 i + j + 1`), each
|
||||
/// followed by a flush. Its superblock (v3) still has the SWMR-write flag
|
||||
/// set and records an end of file of 715 in a 6 030-byte file.
|
||||
fn check_mid_write_copy(f: &File) {
|
||||
let sb = f.superblock();
|
||||
assert_eq!(sb.version, 3);
|
||||
assert!(sb.is_swmr_write());
|
||||
let a = f.dataset("a").unwrap();
|
||||
assert_eq!(a.shape().unwrap(), vec![148]);
|
||||
let want_a: Vec<i64> = (1..=148).collect();
|
||||
assert_eq!(a.read_i64().unwrap(), want_a);
|
||||
let b = f.dataset("b").unwrap();
|
||||
assert_eq!(b.shape().unwrap(), vec![148, 4]);
|
||||
let want_b: Vec<f64> = (0..148)
|
||||
.flat_map(|i| (0..4).map(move |j| (10 * i + j + 1) as f64))
|
||||
.collect();
|
||||
assert_eq!(b.read_f64().unwrap(), want_b);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_copy_taken_mid_write_reads_past_its_recorded_end_of_file() {
|
||||
let path = fixture("swmr_mid_write.h5");
|
||||
// The recorded end of file (715) is far below the file's length; the
|
||||
// chunk indexes and chunks lie past it. libhdf5's SWMR reader does not
|
||||
// bound reads by it, and neither does any open path here.
|
||||
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]
|
||||
fn a_copy_taken_mid_write_reads_as_h5py_swmr_reader_reads_it() {
|
||||
skip_if_no_python!();
|
||||
let path = fixture("swmr_mid_write.h5");
|
||||
// libhdf5 refuses a non-SWMR open of this file ("file is already open
|
||||
// for write"); its SWMR reader reads the values checked above.
|
||||
let script = format!(
|
||||
r#"
|
||||
import h5py, numpy as np
|
||||
with h5py.File("{p}", "r", swmr=True, locking=False) as f:
|
||||
a = f["a"][()]
|
||||
b = f["b"][()]
|
||||
assert np.array_equal(a, np.arange(148) + 1), a
|
||||
assert np.array_equal(b, (np.arange(148)[:, None] * 10 + np.arange(4) + 1).astype("f8")), b
|
||||
print("ok")
|
||||
"#,
|
||||
p = path.display()
|
||||
);
|
||||
let out = Command::new(python())
|
||||
.args(["-c", &script])
|
||||
.output()
|
||||
.unwrap();
|
||||
assert!(
|
||||
out.status.success(),
|
||||
"h5py failed:\n{}",
|
||||
String::from_utf8_lossy(&out.stderr)
|
||||
);
|
||||
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());
|
||||
}
|
||||
|
||||
/// The outcome of reading every dataset of `f` in full, as text: the
|
||||
/// values, or the error.
|
||||
fn read_all(f: &File) -> Vec<String> {
|
||||
["a", "b"]
|
||||
.iter()
|
||||
.map(|n| match f.dataset(n).and_then(|d| d.read_f64()) {
|
||||
Ok(v) => format!("{n}: {} values", v.len()),
|
||||
Err(e) => format!("{n}: error {e}"),
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn open_swmr_reads_a_file_without_the_swmr_write_flag_as_file_open_does() {
|
||||
// The mid-write copy with its flags cleared: a closed file that
|
||||
// records an end of file of 715 in 6 030 bytes, its chunk indexes past
|
||||
// that end. File::open (and libhdf5's plain reader) refuse what lies
|
||||
// past the recorded end; open_swmr must too, and must not retry (the
|
||||
// file is not live). See the h5py test below for libhdf5's SWMR reader.
|
||||
let bytes = mid_write_copy_with_flags(0);
|
||||
let dir = tempfile::tempdir_in(env!("CARGO_TARGET_TMPDIR")).unwrap();
|
||||
let path = dir.path().join("closed_short_eof.h5");
|
||||
std::fs::write(&path, &bytes).unwrap();
|
||||
|
||||
let plain = File::open(&path).unwrap();
|
||||
let swmr = File::open_swmr(&path).unwrap();
|
||||
let swmr_storage = File::open_storage_swmr(Arc::new(bytes)).unwrap();
|
||||
let want = read_all(&plain);
|
||||
assert!(
|
||||
want.iter().all(|r| r.contains("error")),
|
||||
"File::open reads past the recorded end: {want:?}"
|
||||
);
|
||||
for f in [&swmr, &swmr_storage] {
|
||||
assert!(!f.is_swmr_read());
|
||||
assert!(!f.swmr_writer_active().unwrap());
|
||||
let started = std::time::Instant::now();
|
||||
assert_eq!(read_all(f), want);
|
||||
assert!(started.elapsed() < std::time::Duration::from_millis(100));
|
||||
assert_eq!(f.swmr_retries(), 0);
|
||||
}
|
||||
|
||||
// The same bytes with the SWMR-write flag (and write access) set are
|
||||
// read live, past the recorded end.
|
||||
let live = File::open_storage_swmr(Arc::new(mid_write_copy_with_flags(0x05))).unwrap();
|
||||
assert!(live.is_swmr_read());
|
||||
check_mid_write_copy(&live);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn how_h5py_reads_the_closed_copy_past_its_recorded_end_of_file() {
|
||||
skip_if_no_python!();
|
||||
let dir = tempfile::tempdir_in(env!("CARGO_TARGET_TMPDIR")).unwrap();
|
||||
let path = dir.path().join("closed_short_eof.h5");
|
||||
std::fs::write(&path, mid_write_copy_with_flags(0)).unwrap();
|
||||
// libhdf5's plain reader refuses `a` (its chunk index is past the
|
||||
// recorded end of file), as File::open and open_swmr do above.
|
||||
// libhdf5's SWMR reader skips its end-of-allocation check for every
|
||||
// file it opens (H5FD_read), and reads it; it still refuses an object
|
||||
// *header* past the end ("address of object past end of allocation",
|
||||
// H5O_protect). open_swmr does not copy that half-way rule for files
|
||||
// without the SWMR-write flag: they read as File::open reads them.
|
||||
let script = format!(
|
||||
r#"
|
||||
import h5py
|
||||
for swmr in (False, True):
|
||||
try:
|
||||
with h5py.File("{p}", "r", swmr=swmr, locking=False) as f:
|
||||
f["a"][()]
|
||||
except Exception as e:
|
||||
print("refused", swmr)
|
||||
else:
|
||||
print("read", swmr)
|
||||
"#,
|
||||
p = path.display()
|
||||
);
|
||||
let out = Command::new(python())
|
||||
.args(["-c", &script])
|
||||
.output()
|
||||
.unwrap();
|
||||
let stdout = String::from_utf8_lossy(&out.stdout);
|
||||
assert_eq!(
|
||||
stdout.split_whitespace().collect::<Vec<_>>(),
|
||||
["refused", "False", "read", "True"],
|
||||
"{}",
|
||||
String::from_utf8_lossy(&out.stderr)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn open_swmr_returns_a_permanent_error_at_once() {
|
||||
// Not HDF5 at all: SignatureNotFound is not something a writer
|
||||
// causes, so it is not retried (100 attempts would pause about 0.9 s).
|
||||
let started = std::time::Instant::now();
|
||||
let err = File::open_storage_swmr(Arc::new(vec![7u8; 4096])).unwrap_err();
|
||||
assert!(
|
||||
matches!(
|
||||
err,
|
||||
clawhdf5::Error::Format(clawhdf5_format::error::FormatError::SignatureNotFound)
|
||||
),
|
||||
"{err}"
|
||||
);
|
||||
assert!(started.elapsed() < std::time::Duration::from_millis(100));
|
||||
|
||||
// A live file: a lookup of a name it does not have fails at once, and
|
||||
// a torn object header read (below) is retried.
|
||||
let f = File::open_storage_swmr(Arc::new(mid_write_copy_with_flags(0x05))).unwrap();
|
||||
assert!(f.is_swmr_read());
|
||||
let started = std::time::Instant::now();
|
||||
assert!(f.dataset("no_such").is_err());
|
||||
assert!(started.elapsed() < std::time::Duration::from_millis(100));
|
||||
assert_eq!(f.swmr_retries(), 0);
|
||||
}
|
||||
|
||||
/// 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)"
|
||||
);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Every read path retries: strings, variable-length data, attributes.
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// A storage over `bytes` whose read number `fail_at` (counted from the
|
||||
/// last [`Flaky::arm`]) fails with a checksum mismatch, as a read racing a
|
||||
/// SWMR writer's flush can; it counts the reads.
|
||||
struct Flaky {
|
||||
bytes: Vec<u8>,
|
||||
reads: AtomicU32,
|
||||
fail_at: AtomicU32,
|
||||
}
|
||||
|
||||
impl Flaky {
|
||||
fn arm(&self, fail_at: u32) {
|
||||
self.reads.store(0, Ordering::SeqCst);
|
||||
self.fail_at.store(fail_at, Ordering::SeqCst);
|
||||
}
|
||||
}
|
||||
|
||||
impl clawhdf5::Storage for Flaky {
|
||||
fn read_at(
|
||||
&self,
|
||||
offset: u64,
|
||||
len: usize,
|
||||
) -> Result<std::borrow::Cow<'_, [u8]>, clawhdf5_format::error::FormatError> {
|
||||
let n = self.reads.fetch_add(1, Ordering::SeqCst);
|
||||
if n == self.fail_at.load(Ordering::SeqCst) {
|
||||
return Err(clawhdf5_format::error::FormatError::ChecksumMismatch {
|
||||
expected: 0,
|
||||
computed: 1,
|
||||
});
|
||||
}
|
||||
self.bytes.as_slice().read_at(offset, len)
|
||||
}
|
||||
|
||||
fn len(&self) -> u64 {
|
||||
self.bytes.len() as u64
|
||||
}
|
||||
}
|
||||
|
||||
/// `swmr_strings_attrs.h5`: a copy h5py 3.16 (HDF5 2.0) made of its own
|
||||
/// file right after switching it to SWMR mode (the SWMR-write flag is set),
|
||||
/// holding a variable-length string dataset `s` (`alpha beta gamma
|
||||
/// delta`), a variable-length int32 dataset `v` (`[1] [2 3] [4 5 6]`), an
|
||||
/// int64 dataset `d` (`0..10`) with 12 string attributes `k00`..`k11`
|
||||
/// (`value 0`..; dense storage: a fractal heap and a v2 B-tree), a group
|
||||
/// `g` with attribute `note`, and root attributes `title` (a
|
||||
/// variable-length string) and `n`.
|
||||
#[test]
|
||||
fn every_read_path_of_a_live_file_retries_a_failed_read() {
|
||||
let flaky = Arc::new(Flaky {
|
||||
bytes: std::fs::read(fixture("swmr_strings_attrs.h5")).unwrap(),
|
||||
reads: AtomicU32::new(0),
|
||||
fail_at: AtomicU32::new(u32::MAX),
|
||||
});
|
||||
let f = File::open_storage_swmr(flaky.clone()).unwrap();
|
||||
assert!(f.is_swmr_read());
|
||||
let root = f.root();
|
||||
let g = f.group("g").unwrap();
|
||||
let s = f.dataset("s").unwrap();
|
||||
let v = f.dataset("v").unwrap();
|
||||
let d = f.dataset("d").unwrap();
|
||||
let sorted = |m: std::collections::HashMap<String, clawhdf5::AttrValue>| {
|
||||
format!(
|
||||
"{:?}",
|
||||
m.into_iter().collect::<std::collections::BTreeMap<_, _>>()
|
||||
)
|
||||
};
|
||||
type Op<'a> = Box<dyn Fn() -> Result<String, clawhdf5::Error> + 'a>;
|
||||
let ops: Vec<(&str, Op<'_>)> = vec![
|
||||
("root.attrs", Box::new(|| root.attrs().map(sorted))),
|
||||
(
|
||||
"root.attr",
|
||||
Box::new(|| root.attr("title").map(|a| format!("{a:?}"))),
|
||||
),
|
||||
("g.attrs", Box::new(|| g.attrs().map(sorted))),
|
||||
("d.attrs", Box::new(|| d.attrs().map(sorted))),
|
||||
(
|
||||
"d.attr",
|
||||
Box::new(|| d.attr("k05").map(|a| format!("{a:?}"))),
|
||||
),
|
||||
(
|
||||
"root.datasets",
|
||||
Box::new(|| root.datasets().map(|n| format!("{n:?}"))),
|
||||
),
|
||||
(
|
||||
"root.groups",
|
||||
Box::new(|| root.groups().map(|n| format!("{n:?}"))),
|
||||
),
|
||||
("d.shape", Box::new(|| d.shape().map(|n| format!("{n:?}")))),
|
||||
(
|
||||
"d.read_i64",
|
||||
Box::new(|| d.read_i64().map(|n| format!("{n:?}"))),
|
||||
),
|
||||
(
|
||||
"s.read_string",
|
||||
Box::new(|| s.read_string().map(|n| format!("{n:?}"))),
|
||||
),
|
||||
(
|
||||
"s.read_string_bytes",
|
||||
Box::new(|| s.read_string_bytes().map(|n| format!("{n:?}"))),
|
||||
),
|
||||
(
|
||||
"s.read_string_selection",
|
||||
Box::new(|| {
|
||||
s.read_string_selection(&Selection::slice(std::slice::from_ref(&(1..3))))
|
||||
.map(|n| format!("{n:?}"))
|
||||
}),
|
||||
),
|
||||
(
|
||||
"v.read_vlen",
|
||||
Box::new(|| v.read_vlen::<i32>().map(|n| format!("{n:?}"))),
|
||||
),
|
||||
(
|
||||
"v.read_vlen_selection",
|
||||
Box::new(|| {
|
||||
v.read_vlen_selection::<i32>(&Selection::slice(std::slice::from_ref(&(1..3))))
|
||||
.map(|n| format!("{n:?}"))
|
||||
}),
|
||||
),
|
||||
];
|
||||
for (name, op) in &ops {
|
||||
flaky.arm(u32::MAX);
|
||||
let want = op().unwrap_or_else(|e| panic!("{name}: {e}"));
|
||||
let reads = flaky.reads.load(Ordering::SeqCst);
|
||||
// Each read the operation makes fails once in turn: the operation
|
||||
// is run again and returns the same result.
|
||||
for k in 0..reads {
|
||||
flaky.arm(k);
|
||||
let before = f.swmr_retries();
|
||||
let got = op().unwrap_or_else(|e| panic!("{name}, read {k} of {reads} failing: {e}"));
|
||||
assert_eq!(got, want, "{name}, read {k} of {reads} failing");
|
||||
assert_eq!(f.swmr_retries(), before + 1, "{name}, read {k} of {reads}");
|
||||
}
|
||||
}
|
||||
flaky.arm(u32::MAX);
|
||||
assert!(ops[0].1().unwrap().contains("live strings"));
|
||||
assert!(ops[3].1().unwrap().contains("value 11"));
|
||||
assert_eq!(
|
||||
ops[9].1().unwrap(),
|
||||
r#"["alpha", "beta", "gamma", "delta"]"#
|
||||
);
|
||||
assert_eq!(ops[12].1().unwrap(), "[[1], [2, 3], [4, 5, 6]]");
|
||||
|
||||
// With one attempt the failure reaches the caller.
|
||||
let mut f = File::open_storage_swmr(flaky.clone()).unwrap();
|
||||
f.set_swmr_read_attempts(1);
|
||||
let s = f.dataset("s").unwrap();
|
||||
flaky.arm(0);
|
||||
assert!(s.read_string().is_err());
|
||||
}
|
||||
@@ -7,7 +7,9 @@ change. Progress: M0 and M1 are done, and so is M2 (branch
|
||||
`File::open_storage` gives the facade's read API over any `Storage` (see
|
||||
`CHANGELOG.md`, "Range reads, milestone M2"). M3 is done on branch
|
||||
`feat/p3-m3-remote`: the `clawhdf5-remote` crate (block cache, HTTP(S),
|
||||
object stores) and URLs in `h5rs` (see the M3 status below). M4 (wasm) is
|
||||
object stores) and URLs in `h5rs` (see the M3 status below). M5 (SWMR) is
|
||||
done on branch `feat/p3-m5-swmr-reader`, with its own design in
|
||||
[`swmr.md`](swmr.md) (see the M5 status below). M4 (wasm) is
|
||||
next. Every count in §1–§2 was
|
||||
|
||||
change. Progress: M1, first part (the `Storage` trait and the metadata
|
||||
@@ -523,6 +525,19 @@ fast path within benchmark noise.
|
||||
**M5 — SWMR and growth (later, separate design).** `Storage::len()` may grow;
|
||||
add `File::refresh()` that re-reads the superblock/EOF and invalidates cached
|
||||
blocks past the old end. Needs libhdf5 SWMR semantics research first.
|
||||
- *Status 2026-09-27:* done on branch `feat/p3-m5-swmr-reader`; design and
|
||||
libhdf5 research in [`swmr.md`](swmr.md). Differences from the sketch
|
||||
above: the refresh is per dataset (`Dataset::refresh`, as libhdf5's
|
||||
`H5Drefresh`), not per file — a SWMR writer only grows datasets, and the
|
||||
superblock's EOF is not kept up to date by it, so there is nothing to
|
||||
re-read there (`File::swmr_writer_active` re-reads its flags). A live
|
||||
file (`File::open_swmr`) reads through `FileStorage` (positioned reads,
|
||||
`len()` the current length) with no block or chunk cache, rather than
|
||||
invalidating cached blocks: `BlockCache`/`HttpStorage` stay snapshot
|
||||
readers. Operations that fail with an error a racing write can cause are
|
||||
retried up to 100 times (libhdf5's default for SWMR readers).
|
||||
Tested against a live h5py writer and h5py's SWMR reader
|
||||
(`crates/clawhdf5/tests/swmr_interop.rs`).
|
||||
|
||||
Total: roughly 6–10 engineer-weeks for M0–M4 (estimate, not measured).
|
||||
|
||||
|
||||
@@ -0,0 +1,178 @@
|
||||
# Design: reading files a SWMR writer is still appending to (range-read M5)
|
||||
|
||||
Status: design 2026-09-27, implemented on branch `feat/p3-m5-swmr-reader`
|
||||
(see "Status" at the end). This is milestone M5 of
|
||||
[`range-reads.md`](range-reads.md): "`Storage::len()` may grow; add a
|
||||
refresh". It covers the reader only; clawhdf5 does not write SWMR files.
|
||||
|
||||
## What libhdf5 does
|
||||
|
||||
A SWMR ("single writer, multiple readers") writer is a libhdf5 process that
|
||||
opened a file with `libver='latest'` and switched to SWMR mode (h5py
|
||||
`f.swmr_mode = True`, `H5Fstart_swmr_write`). Readers open the same file
|
||||
with `H5F_ACC_SWMR_READ` (h5py `File(path, 'r', swmr=True)`) while the
|
||||
writer keeps appending. What the format and the library guarantee:
|
||||
|
||||
- **Superblock v3, flags set.** The writer sets the superblock's
|
||||
file-consistency flags to write access + SWMR write (`0x05`) and clears
|
||||
them on close. The superblock's end-of-file address is *not* kept up to
|
||||
date while writing: a copy of a file taken mid-write records an EOF of a
|
||||
few hundred bytes while the file is tens of kilobytes (checked on tank,
|
||||
2026-09-27, h5py 3.16 / HDF5 2.0: EOF 715 in a 17 857-byte file). A SWMR
|
||||
reader therefore skips libhdf5's end-of-allocation check for every read
|
||||
(`H5FD_read`: "allow access to data past the end of the allocated space
|
||||
… for SWMR read access"), and bounds reads by the file's real length.
|
||||
- **A non-SWMR open of such a file fails** in libhdf5: "file is already
|
||||
open for write (may use <h5clear file> to clear file consistency flags)".
|
||||
- **The writer only appends.** New objects and attributes cannot be created
|
||||
in SWMR mode; datasets grow with `H5Dset_extent` and are written. Chunked
|
||||
datasets with one unlimited dimension use an Extensible Array index, with
|
||||
more than one a version-2 B-tree; both are updated in a SWMR-safe way.
|
||||
Fixed Array and single-chunk indexes are for datasets that cannot grow.
|
||||
- **Flush ordering.** Every metadata structure the writer uses is
|
||||
checksummed, and flush dependencies order the writes: a chunk's data is
|
||||
written before the index entry that points to it, and index blocks before
|
||||
the object header whose dataspace announces the new extent. A reader
|
||||
that reads the object header first and the index after sees an index at
|
||||
least as new as the extent, so every chunk inside the extent it read is
|
||||
either in the index or never written (then it reads as the fill value, as
|
||||
it does for libhdf5's reader).
|
||||
- **Refresh.** A reader sees a dataset's new extent only when it refreshes
|
||||
it (`H5Drefresh`, h5py `Dataset.refresh()`), which evicts the dataset's
|
||||
cached metadata and reads the object header again.
|
||||
- **Retries.** Reads are not atomic against writes on every system, so a
|
||||
checksum can fail when a structure is read while the writer rewrites it.
|
||||
A SWMR reader reads checksummed metadata up to 100 times before failing
|
||||
(`H5Pset_metadata_read_attempts`; default 100 for SWMR access, 1
|
||||
otherwise — `H5Ppublic.h` of HDF5 1.14.6).
|
||||
|
||||
## What clawhdf5 did before
|
||||
|
||||
- `File::open` of a SWMR-flagged file bounded every read by the recorded
|
||||
EOF (`Superblock::data_end` only tolerated an EOF *past* the end of the
|
||||
file). A file copied or read mid-write therefore listed, but every
|
||||
chunked read failed ("unexpected EOF: need 787 bytes, have 715"), and
|
||||
`h5rs check` reported chunk indexes "past the end of the file".
|
||||
- A `File` is a snapshot: an mmap (or a buffer) of the length at open, and a
|
||||
per-file chunk cache that keeps each dataset's chunk index and decoded
|
||||
chunks for the life of the `File`. A reader could not see growth at all,
|
||||
and a mapping of a file that is being rewritten can change under a read.
|
||||
|
||||
## Design
|
||||
|
||||
1. **Bound SWMR files by their length.** `Superblock::data_end` returns the
|
||||
file's length for a version-3 superblock with the SWMR-write flag,
|
||||
whatever EOF it records, as libhdf5's SWMR reader does. Every existing
|
||||
open path (`File::open`, `open_storage`, `h5rs`) then reads a finished
|
||||
copy of a live file. We keep opening such files without a SWMR flag
|
||||
(libhdf5 refuses): it is read-only and the alternative is an error.
|
||||
2. **A live open: `File::open_swmr(path)` / `File::open_storage_swmr`.**
|
||||
- The file is read with positioned reads (`FileStorage`, `pread` on
|
||||
Unix, `seek_read` on Windows), never mapped, and `len()` is the file's
|
||||
current length, so reads past the length seen at open work.
|
||||
- The facade's view of the file (`FileData`) is *live*: reads are not
|
||||
clamped to an end fixed at open, only by the storage's current length.
|
||||
- The chunk cache is not used: every read reads the chunk index and the
|
||||
chunks it needs again. A cached index would hide new chunks, and a
|
||||
cached partial edge chunk would read as fill where the writer has since
|
||||
written data (unfiltered edge chunks are rewritten in place).
|
||||
- Only a file whose superblock has the SWMR-write flag when it is opened
|
||||
is read live. Any other file is read exactly as `File::open` reads it
|
||||
(bounded by its recorded EOF, through the chunk cache, no retries), and
|
||||
`is_swmr_read()` is `false`. libhdf5's SWMR reader is looser: it skips
|
||||
the end-of-allocation check in `H5FD_read` for every file it opens,
|
||||
flagged or not, yet still refuses an object header past the EOF
|
||||
(`H5O_protect`, "address of object past end of allocation"). On a
|
||||
closed file whose EOF is below its length (checked 2026-09-27, h5py
|
||||
3.16 / HDF5 2.0, the mid-write fixture with its flags cleared) it
|
||||
therefore reads a dataset whose chunk index lies past the EOF, which
|
||||
its plain reader refuses. We follow the plain reader there: the EOF of
|
||||
a file no SWMR writer has open is what the file says it is.
|
||||
- A storage that caches blocks (`clawhdf5-remote`'s `BlockCache`) would
|
||||
serve stale bytes; `HttpStorage` also pins a file by ETag and length and
|
||||
refuses a changed file. Remote SWMR is out of scope.
|
||||
3. **`Dataset::refresh()`** reads the dataset's object header again (same
|
||||
address) and replaces the handle's copy, so `shape()` and every later
|
||||
read use the new extent. Like h5py, a handle that is not refreshed keeps
|
||||
its extent; its reads still read the index as it is now, and only return
|
||||
elements inside that extent.
|
||||
4. **Bounded retries.** In a live file, an operation (open, lookups and
|
||||
listings, refresh, every dataset read — typed, raw, selections, strings
|
||||
and variable-length data with their global-heap decoding — attributes,
|
||||
`File::decode_*`, `verify_provenance`) that fails with an error a
|
||||
concurrent write can
|
||||
cause is run again from the start, up to `File::swmr_read_attempts()`
|
||||
times (default 100, libhdf5's default; `set_swmr_read_attempts` changes
|
||||
it), sleeping 1 µs, 2 µs, … up to 10 ms between attempts (under a second
|
||||
in all). libhdf5 retries the one structure whose checksum failed; we
|
||||
retry the whole operation, because the parsers are pure functions of the
|
||||
bytes they read. Which errors: those libhdf5 retries (`H5C__load_entry`)
|
||||
— a checksum mismatch, and a failure to decode the prefix it reads
|
||||
before the checksum to size the structure (for an object header its
|
||||
signature and version: a test that garbles every byte of a header gets
|
||||
`InvalidObjectHeaderVersion`) — and a read past the file's current end
|
||||
(short here; libhdf5 reads zeros there, which fail the checksum). Every
|
||||
other error (an unsupported version or message, a file that is not HDF5,
|
||||
a structure corrupt behind a valid checksum) is returned at once: an
|
||||
earlier version retried nearly every format error, so a permanent one
|
||||
cost about 0.9 s of pauses. Non-live files never retry.
|
||||
Results are only returned from a run where every structure verified, so
|
||||
a torn metadata read is an error, never data. `File::swmr_retries()`
|
||||
counts the retries (libhdf5: `H5Fget_metadata_read_retry_info`).
|
||||
Attribute reads leave out an attribute they cannot read (or return a
|
||||
variable-length string one as `AttrValue::Raw`) instead of failing;
|
||||
on a live file such an error of a retried kind runs the read again too,
|
||||
and after the last attempt the last result is returned as before. The
|
||||
zero-copy reads (`read_raw_ref`, `read_as_slice`, `read_*_zerocopy`)
|
||||
need the file in memory, which a live file never is: they report
|
||||
`None` / `ContiguousStorageRequired` without reading data.
|
||||
Global heap collections (variable-length data) have no checksum in
|
||||
HDF5, like raw data, so a torn read of one is only caught when it fails
|
||||
a check (its signature, a bound).
|
||||
Raw data has no checksum in HDF5 (unless Fletcher-32 is on), in libhdf5
|
||||
as here: correctness rests on the writer's ordering (chunk data before
|
||||
the index entry, only appends), as for libhdf5's reader.
|
||||
5. **Writer state.** `File::swmr_writer_active()` reads the superblock
|
||||
flags again, so a reader can tell when the writer has closed the file.
|
||||
A writer that crashed or was killed never clears the flag, so a reader
|
||||
that follows a file needs a second stop condition (the README example
|
||||
stops after a minute without growth).
|
||||
|
||||
What stays out: SWMR writing, VFD SWMR (HDF5 1.13's page-buffer protocol,
|
||||
not in 1.14 or 2.0), remote SWMR, refresh of groups/attributes (a SWMR
|
||||
writer cannot add them), and `MmapFile`/`LazyFile`.
|
||||
|
||||
## Tests
|
||||
|
||||
- `crates/clawhdf5-format`: `data_end` of a SWMR-flagged v3 superblock whose
|
||||
EOF is below the file length.
|
||||
- `crates/clawhdf5/tests/swmr_interop.rs`:
|
||||
- a copy of a file taken mid-write (fixture) reads like h5py's SWMR reader;
|
||||
- the same bytes with the flags cleared read as `File::open` reads them
|
||||
(and how h5py's two readers read them);
|
||||
- a permanent error returns at once; every read path (listings,
|
||||
attributes, strings, variable-length data) with each of its reads
|
||||
failing once in turn returns the same result;
|
||||
- a live test: an h5py writer (`swmr_mode = True`) appends to a 1-D and a
|
||||
2-D dataset with one unlimited dimension (Extensible Array, one of them
|
||||
gzip) and a 2-D dataset with two (v2 B-tree), flushing after every step,
|
||||
while a Rust reader refreshes and reads them in a loop, and an h5py SWMR
|
||||
reader does the same as the reference. Every value read must be the value
|
||||
the writer wrote (a deterministic function of its position), extents
|
||||
never shrink, and after the writer closes both readers must read the
|
||||
same data as h5py.
|
||||
|
||||
## Status
|
||||
|
||||
Implemented 2026-09-27 on branch `feat/p3-m5-swmr-reader` as designed
|
||||
above (`CHANGELOG.md`, "Range reads, milestone M5"). Observed on tank the
|
||||
same day (h5py 3.16 / HDF5 2.0, `cargo test -p clawhdf5 --test
|
||||
swmr_interop`, and once with `CLAWHDF5_SWMR_STEPS=20000` in a release
|
||||
build): no read returned a value the writer had not written at that
|
||||
position, h5py's reader agreed, and retries were needed but rare (the
|
||||
failures seen were checksum mismatches, each cured by one retry). A
|
||||
variant of the test with the chunk cache left on in live mode fails it
|
||||
(stale chunk index / edge chunk), which is why live files do not use it.
|
||||
|
||||
Also found: `File::open` of such a file had been failing since the
|
||||
end-of-file check of 2026-09-26 (item 1; `docs/known-issues.md`).
|
||||
@@ -7,6 +7,29 @@ deleting it.
|
||||
|
||||
---
|
||||
|
||||
## Files a SWMR writer had open could not be read past a stale end of file
|
||||
|
||||
**Status:** fixed 2026-09-27 (branch `feat/p3-m5-swmr-reader`), before any
|
||||
release: reads have been bounded by the recorded end of file since
|
||||
`7d7a7e7` (2026-09-26), which no release contains.
|
||||
|
||||
A libhdf5 writer in SWMR mode (h5py `f.swmr_mode = True`) sets the
|
||||
superblock's SWMR-write flag and does not keep its end-of-file address up to
|
||||
date: a copy h5py made of its own file mid-write records 715 in a
|
||||
6 030-byte file (h5py 3.16 / HDF5 2.0, tank). `Superblock::data_end` only
|
||||
ignored the recorded end when it lay *past* the end of the file, so every
|
||||
reader bounded such a file at 715 bytes: it listed, but every chunked read
|
||||
failed ("unexpected EOF: need 787 bytes, have 715") and `h5rs check`
|
||||
reported the chunk indexes past the end of the file. Never wrong data.
|
||||
|
||||
**Fix:** for a v3 superblock with the SWMR-write flag, the data ends at the
|
||||
end of the file, as libhdf5's SWMR reader reads it. **Test:**
|
||||
`crates/clawhdf5/tests/swmr_interop.rs` (the mid-write copy,
|
||||
`tests/fixtures/swmr_mid_write.h5`, through every open path and against
|
||||
h5py's SWMR reader). A file still being written is read with
|
||||
`File::open_swmr` (see `docs/design/swmr.md`); `File::open` maps the file
|
||||
at its length at open and is not meant for files that change while open.
|
||||
|
||||
## Fletcher-32 checksums disagreed with libhdf5 on about 1 chunk in 32768
|
||||
|
||||
**Status:** fixed 2026-09-26, after v2.7.0. **Every release (v2.1.0 to
|
||||
|
||||
Reference in New Issue
Block a user