Lazy remote files in the browser (M4), SWMR reader (M5), Python remote reads and editing #19

Merged
osobh merged 33 commits from feat/p3-wasm-swmr-python into main 2026-09-27 14:47:26 +00:00
8 changed files with 524 additions and 191 deletions
Showing only changes of commit 2e906ffe3a - Show all commits
+6 -2
View File
@@ -34,8 +34,9 @@ Design: `docs/design/swmr.md`.
- **`Dataset::refresh()`** reads the dataset's object header again - **`Dataset::refresh()`** reads the dataset's object header again
(`H5Drefresh`, h5py `Dataset.refresh()`), so `shape()` and later reads (`H5Drefresh`, h5py `Dataset.refresh()`), so `shape()` and later reads
see the writer's appends; a handle keeps its extent until refreshed. see the writer's appends; a handle keeps its extent until refreshed.
- **Bounded retries.** On a SWMR-read file, an operation (open, lookups, - **Bounded retries.** On a SWMR-read file, an operation (open, lookups
refresh, reads) that fails with an error a concurrent write can cause — 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 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 past the file's current end, an object header prefix that does not
decode — is run again from the start, up to decode — is run again from the start, up to
@@ -52,6 +53,9 @@ Design: `docs/design/swmr.md`.
(fixture `tests/fixtures/swmr_mid_write.h5`) through every open path and (fixture `tests/fixtures/swmr_mid_write.h5`) through every open path and
against h5py's SWMR reader; garbled reads (a storage that corrupts the against h5py's SWMR reader; garbled reads (a storage that corrupts the
next reads) retried, never returned, and given up after the attempts; 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 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 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 2-D dataset with two (v2 B-tree index) for 2 500 steps
+53 -18
View File
@@ -576,7 +576,14 @@ pub fn find_attribute_in_file(
offset_size: u8, offset_size: u8,
length_size: u8, length_size: u8,
) -> Result<Option<AttributeMessage>, FormatError> { ) -> 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 /// [`find_attribute_in_file`] over any [`Storage`] (see
@@ -594,16 +601,48 @@ pub fn find_attribute_in<S: Storage + ?Sized>(
) -> Result<Option<AttributeMessage>, FormatError> { ) -> Result<Option<AttributeMessage>, FormatError> {
match file_data.as_contiguous() { match file_data.as_contiguous() {
Some(all) => find_attribute_in_file(all, header, name, offset_size, length_size), 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>( fn find_attribute_core<S: Storage + ?Sized>(
file_data: &S, file_data: &S,
header: &ObjectHeader, header: &ObjectHeader,
name: &str, name: &str,
offset_size: u8, offset_size: u8,
length_size: u8, length_size: u8,
errors: &mut Vec<FormatError>,
) -> Result<Option<AttributeMessage>, FormatError> { ) -> Result<Option<AttributeMessage>, FormatError> {
let attr_info = find_attribute_info(header, offset_size)?; let attr_info = find_attribute_info(header, offset_size)?;
let dense = attr_info let dense = attr_info
@@ -612,12 +651,10 @@ fn find_attribute_core<S: Storage + ?Sized>(
let Some((fh_addr, btree_addr)) = dense else { let Some((fh_addr, btree_addr)) = dense else {
// Compact only (or dense storage without a name index, which a // Compact only (or dense storage without a name index, which a
// listing reports): as a listing finds it. // listing reports): as a listing finds it.
return Ok( let (attrs, errs) =
extract_attributes_tolerant_in(file_data, header, offset_size, length_size)? extract_attributes_tolerant_in(file_data, header, offset_size, length_size)?;
.0 errors.extend(errs);
.into_iter() return Ok(attrs.into_iter().find(|a| a.name == name));
.find(|a| a.name == name),
);
}; };
let btree_hdr = BTreeV2Header::parse_in( let btree_hdr = BTreeV2Header::parse_in(
file_data, 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)?; 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 { if btree_hdr.tree_type != ATTRIBUTE_NAME_INDEX || btree_hdr.record_size < 4 {
return Ok( let (attrs, errs) =
extract_attributes_tolerant_in(file_data, header, offset_size, length_size)? extract_attributes_tolerant_in(file_data, header, offset_size, length_size)?;
.0 errors.extend(errs);
.into_iter() return Ok(attrs.into_iter().find(|a| a.name == name));
.find(|a| a.name == name),
);
} }
// A listing has the compact attributes first. // 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) AttributeMessage::parse_in_storage(&d, file_data, offset_size, length_size)
}); });
// One that cannot be read is left out, as from a listing. // One that cannot be read is left out, as from a listing.
if let Ok(attr) = attr match attr {
&& attr.name == name Ok(attr) if attr.name == name => return Ok(Some(attr)),
{ Ok(_) => {}
return Ok(Some(attr)); Err(e) => errors.push(e),
} }
} }
Ok(None) Ok(None)
+225 -147
View File
@@ -28,7 +28,9 @@ use clawhdf5_format::superblock_ext::{self, CacheImageState};
use crate::cache_image::{self, ImageView}; use crate::cache_image::{self, ImageView};
use crate::error::Error; 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 // FileData — internal storage for owned bytes, an mmap, or any Storage
@@ -662,6 +664,34 @@ impl File {
} }
} }
/// 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 /// 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 /// 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 /// would hide the writer's new chunks and a cached edge chunk would
@@ -879,13 +909,15 @@ impl File {
/// [`AttrValue::Raw`] attribute. Variable-length strings are resolved in /// [`AttrValue::Raw`] attribute. Variable-length strings are resolved in
/// this file's global heap; see [`Dataset::read_string`] for the values. /// 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> { pub fn decode_strings(&self, datatype: &Datatype, raw: &[u8]) -> Result<Vec<String>, Error> {
with_bytes!(&self.data, |d| crate::vlen::decode_strings( self.retry(|| {
d, with_bytes!(&self.data, |d| crate::vlen::decode_strings(
datatype, d,
raw, datatype,
self.offset_size(), raw,
self.length_size(), self.offset_size(),
)) self.length_size(),
))
})
} }
/// Like [`decode_strings`](Self::decode_strings) for variable-length /// Like [`decode_strings`](Self::decode_strings) for variable-length
@@ -896,13 +928,15 @@ impl File {
datatype: &Datatype, datatype: &Datatype,
raw: &[u8], raw: &[u8],
) -> Result<Vec<Vec<u8>>, Error> { ) -> Result<Vec<Vec<u8>>, Error> {
with_bytes!(self.data.meta()?, |d| crate::vlen::decode_string_bytes( self.retry(|| {
d, with_bytes!(self.data.meta()?, |d| crate::vlen::decode_string_bytes(
datatype, d,
raw, datatype,
self.offset_size(), raw,
self.length_size(), self.offset_size(),
)) self.length_size(),
))
})
} }
/// Decode the variable-length sequences in `raw`, a buffer of elements /// Decode the variable-length sequences in `raw`, a buffer of elements
@@ -914,13 +948,15 @@ impl File {
datatype: &Datatype, datatype: &Datatype,
raw: &[u8], raw: &[u8],
) -> Result<Vec<Vec<T>>, Error> { ) -> Result<Vec<Vec<T>>, Error> {
with_bytes!(self.data.meta()?, |d| crate::vlen::decode_vlen( self.retry(|| {
d, with_bytes!(self.data.meta()?, |d| crate::vlen::decode_vlen(
datatype, d,
raw, datatype,
self.offset_size(), raw,
self.length_size(), self.offset_size(),
)) self.length_size(),
))
})
} }
fn parse_header(&self, address: u64) -> Result<ObjectHeader, FormatError> { fn parse_header(&self, address: u64) -> Result<ObjectHeader, FormatError> {
@@ -964,28 +1000,32 @@ pub struct Group<'f> {
impl<'f> Group<'f> { impl<'f> Group<'f> {
/// List the names of datasets in this group. /// List the names of datasets in this group.
pub fn datasets(&self) -> Result<Vec<String>, Error> { pub fn datasets(&self) -> Result<Vec<String>, Error> {
let entries = self.children()?; self.file.retry(|| {
let mut names = Vec::new(); let entries = self.children()?;
for entry in &entries { let mut names = Vec::new();
let hdr = self.file.parse_header(entry.object_header_address)?; for entry in &entries {
if has_message(&hdr, MessageType::DataLayout) { let hdr = self.file.parse_header(entry.object_header_address)?;
names.push(entry.name.clone()); if has_message(&hdr, MessageType::DataLayout) {
names.push(entry.name.clone());
}
} }
} Ok(names)
Ok(names) })
} }
/// List the names of subgroups in this group. /// List the names of subgroups in this group.
pub fn groups(&self) -> Result<Vec<String>, Error> { pub fn groups(&self) -> Result<Vec<String>, Error> {
let entries = self.children()?; self.file.retry(|| {
let mut names = Vec::new(); let entries = self.children()?;
for entry in &entries { let mut names = Vec::new();
let hdr = self.file.parse_header(entry.object_header_address)?; for entry in &entries {
if is_group(&hdr) { let hdr = self.file.parse_header(entry.object_header_address)?;
names.push(entry.name.clone()); if is_group(&hdr) {
names.push(entry.name.clone());
}
} }
} Ok(names)
Ok(names) })
} }
/// Read all attributes of this group. /// Read all attributes of this group.
@@ -1006,13 +1046,14 @@ impl<'f> Group<'f> {
pub fn attrs_with_errors( pub fn attrs_with_errors(
&self, &self,
) -> Result<(HashMap<String, AttrValue>, Vec<FormatError>), Error> { ) -> Result<(HashMap<String, AttrValue>, Vec<FormatError>), Error> {
let hdr = self.file.parse_header(self.address)?; self.file.retry_attrs(|| {
with_bytes!(&self.file.data, |d| read_attrs( let hdr = self.file.parse_header(self.address)?;
d, let (attrs, errors, read_errors) = with_bytes!(&self.file.data, |d| {
&hdr, read_attrs_reporting(d, &hdr, self.file.offset_size(), self.file.length_size())
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. /// Get a dataset within this group by name.
@@ -1045,14 +1086,16 @@ impl<'f> Group<'f> {
/// that name, found without reading the other attributes when they are /// that name, found without reading the other attributes when they are
/// stored densely. /// stored densely.
pub fn attr(&self, name: &str) -> Result<Option<AttrValue>, Error> { pub fn attr(&self, name: &str) -> Result<Option<AttrValue>, Error> {
let hdr = self.file.parse_header(self.address)?; self.file.retry_attrs(|| {
with_bytes!(&self.file.data, |d| read_attr( let hdr = self.file.parse_header(self.address)?;
d, with_bytes!(&self.file.data, |d| read_attr_reporting(
&hdr, d,
name, &hdr,
self.file.offset_size(), name,
self.file.length_size(), self.file.offset_size(),
)) self.file.length_size(),
))
})
} }
/// The object header address of the child called `name`: the entry of /// The object header address of the child called `name`: the entry of
@@ -1187,17 +1230,20 @@ impl<'f> Dataset<'f> {
/// Read all data as `f64` values. /// Read all data as `f64` values.
pub fn read_f64(&self) -> Result<Vec<f64>, Error> { pub fn read_f64(&self) -> Result<Vec<f64>, Error> {
let dt = self.datatype()?; self.file.retry(|| {
// A contiguous dataset is converted straight from the file bytes; going let dt = self.datatype()?;
// through `read_raw` first copied the whole dataset an extra time. // A contiguous dataset is converted straight from the file bytes;
if let Ok(Some(bytes)) = self.contiguous_raw() { // going through `read_raw` first copied the whole dataset an
return Ok(data_read::read_as_f64(&bytes, &dt)?); // extra time.
} if let Ok(Some(bytes)) = self.contiguous_raw() {
if let Some(values) = self.read_chunked_native::<f64>()? { return Ok(data_read::read_as_f64(&bytes, &dt)?);
return Ok(values); }
} if let Some(values) = self.read_chunked_native::<f64>()? {
let raw = self.read_raw()?; return Ok(values);
Ok(data_read::read_as_f64(&raw, &dt)?) }
let raw = self.read_raw()?;
Ok(data_read::read_as_f64(&raw, &dt)?)
})
} }
/// Zero-copy read of contiguous native-endian `f64` data. /// Zero-copy read of contiguous native-endian `f64` data.
@@ -1207,62 +1253,74 @@ impl<'f> Dataset<'f> {
/// ///
/// Read all data as `f32` values. /// Read all data as `f32` values.
pub fn read_f32(&self) -> Result<Vec<f32>, Error> { pub fn read_f32(&self) -> Result<Vec<f32>, Error> {
let dt = self.datatype()?; self.file.retry(|| {
// A contiguous dataset is converted straight from the file bytes; going let dt = self.datatype()?;
// through `read_raw` first copied the whole dataset an extra time. // A contiguous dataset is converted straight from the file bytes;
if let Ok(Some(bytes)) = self.contiguous_raw() { // going through `read_raw` first copied the whole dataset an
return Ok(data_read::read_as_f32(&bytes, &dt)?); // extra time.
} if let Ok(Some(bytes)) = self.contiguous_raw() {
if let Some(values) = self.read_chunked_native::<f32>()? { return Ok(data_read::read_as_f32(&bytes, &dt)?);
return Ok(values); }
} if let Some(values) = self.read_chunked_native::<f32>()? {
let raw = self.read_raw()?; return Ok(values);
Ok(data_read::read_as_f32(&raw, &dt)?) }
let raw = self.read_raw()?;
Ok(data_read::read_as_f32(&raw, &dt)?)
})
} }
/// Read all data as `i32` values. /// Read all data as `i32` values.
pub fn read_i32(&self) -> Result<Vec<i32>, Error> { pub fn read_i32(&self) -> Result<Vec<i32>, Error> {
let dt = self.datatype()?; self.file.retry(|| {
// A contiguous dataset is converted straight from the file bytes; going let dt = self.datatype()?;
// through `read_raw` first copied the whole dataset an extra time. // A contiguous dataset is converted straight from the file bytes;
if let Ok(Some(bytes)) = self.contiguous_raw() { // going through `read_raw` first copied the whole dataset an
return Ok(data_read::read_as_i32(&bytes, &dt)?); // extra time.
} if let Ok(Some(bytes)) = self.contiguous_raw() {
if let Some(values) = self.read_chunked_native::<i32>()? { return Ok(data_read::read_as_i32(&bytes, &dt)?);
return Ok(values); }
} if let Some(values) = self.read_chunked_native::<i32>()? {
let raw = self.read_raw()?; return Ok(values);
Ok(data_read::read_as_i32(&raw, &dt)?) }
let raw = self.read_raw()?;
Ok(data_read::read_as_i32(&raw, &dt)?)
})
} }
/// Read all data as `i64` values. /// Read all data as `i64` values.
pub fn read_i64(&self) -> Result<Vec<i64>, Error> { pub fn read_i64(&self) -> Result<Vec<i64>, Error> {
let dt = self.datatype()?; self.file.retry(|| {
// A contiguous dataset is converted straight from the file bytes; going let dt = self.datatype()?;
// through `read_raw` first copied the whole dataset an extra time. // A contiguous dataset is converted straight from the file bytes;
if let Ok(Some(bytes)) = self.contiguous_raw() { // going through `read_raw` first copied the whole dataset an
return Ok(data_read::read_as_i64(&bytes, &dt)?); // extra time.
} if let Ok(Some(bytes)) = self.contiguous_raw() {
if let Some(values) = self.read_chunked_native::<i64>()? { return Ok(data_read::read_as_i64(&bytes, &dt)?);
return Ok(values); }
} if let Some(values) = self.read_chunked_native::<i64>()? {
let raw = self.read_raw()?; return Ok(values);
Ok(data_read::read_as_i64(&raw, &dt)?) }
let raw = self.read_raw()?;
Ok(data_read::read_as_i64(&raw, &dt)?)
})
} }
/// Read all data as `u64` values. /// Read all data as `u64` values.
pub fn read_u64(&self) -> Result<Vec<u64>, Error> { pub fn read_u64(&self) -> Result<Vec<u64>, Error> {
let dt = self.datatype()?; self.file.retry(|| {
// A contiguous dataset is converted straight from the file bytes; going let dt = self.datatype()?;
// through `read_raw` first copied the whole dataset an extra time. // A contiguous dataset is converted straight from the file bytes;
if let Ok(Some(bytes)) = self.contiguous_raw() { // going through `read_raw` first copied the whole dataset an
return Ok(data_read::read_as_u64(&bytes, &dt)?); // extra time.
} if let Ok(Some(bytes)) = self.contiguous_raw() {
if let Some(values) = self.read_chunked_native::<u64>()? { return Ok(data_read::read_as_u64(&bytes, &dt)?);
return Ok(values); }
} if let Some(values) = self.read_chunked_native::<u64>()? {
let raw = self.read_raw()?; return Ok(values);
Ok(data_read::read_as_u64(&raw, &dt)?) }
let raw = self.read_raw()?;
Ok(data_read::read_as_u64(&raw, &dt)?)
})
} }
/// Read all data as `String` values, in row-major order. /// Read all data as `String` values, in row-major order.
@@ -1273,17 +1331,21 @@ impl<'f> Dataset<'f> {
/// them; bytes that are not valid UTF-8 are replaced with U+FFFD — use /// 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. /// [`read_string_bytes`](Self::read_string_bytes) for the exact bytes.
pub fn read_string(&self) -> Result<Vec<String>, Error> { pub fn read_string(&self) -> Result<Vec<String>, Error> {
let raw = self.read_raw()?; self.file.retry(|| {
let dt = self.datatype()?; let raw = self.read_raw()?;
self.file.decode_strings(&dt, &raw) let dt = self.datatype()?;
self.file.decode_strings(&dt, &raw)
})
} }
/// Read a variable-length string dataset as the exact bytes of each /// Read a variable-length string dataset as the exact bytes of each
/// string (what h5py's `Dataset[()]` returns), in row-major order. /// string (what h5py's `Dataset[()]` returns), in row-major order.
pub fn read_string_bytes(&self) -> Result<Vec<Vec<u8>>, Error> { pub fn read_string_bytes(&self) -> Result<Vec<Vec<u8>>, Error> {
let raw = self.read_raw()?; self.file.retry(|| {
let dt = self.datatype()?; let raw = self.read_raw()?;
self.file.decode_string_bytes(&dt, &raw) let dt = self.datatype()?;
self.file.decode_string_bytes(&dt, &raw)
})
} }
/// Read the selected elements of a fixed- or variable-length string /// Read the selected elements of a fixed- or variable-length string
@@ -1292,9 +1354,11 @@ impl<'f> Dataset<'f> {
&self, &self,
selection: &clawhdf5_format::selection::Selection, selection: &clawhdf5_format::selection::Selection,
) -> Result<Vec<String>, Error> { ) -> Result<Vec<String>, Error> {
let raw = self.read_selection(selection)?; self.file.retry(|| {
let dt = self.datatype()?; let raw = self.read_selection(selection)?;
self.file.decode_strings(&dt, &raw) let dt = self.datatype()?;
self.file.decode_strings(&dt, &raw)
})
} }
/// Read a variable-length sequence dataset (h5py /// Read a variable-length sequence dataset (h5py
@@ -1303,9 +1367,11 @@ impl<'f> Dataset<'f> {
/// is converted to `T` as [`read_f64`](Self::read_f64) and the other /// is converted to `T` as [`read_f64`](Self::read_f64) and the other
/// typed readers convert. A null element is an empty sequence. /// typed readers convert. A null element is an empty sequence.
pub fn read_vlen<T: crate::vlen::VlenValue>(&self) -> Result<Vec<Vec<T>>, Error> { pub fn read_vlen<T: crate::vlen::VlenValue>(&self) -> Result<Vec<Vec<T>>, Error> {
let raw = self.read_raw()?; self.file.retry(|| {
let dt = self.datatype()?; let raw = self.read_raw()?;
self.file.decode_vlen(&dt, &raw) let dt = self.datatype()?;
self.file.decode_vlen(&dt, &raw)
})
} }
/// Read the selected elements of a variable-length sequence dataset /// Read the selected elements of a variable-length sequence dataset
@@ -1314,9 +1380,11 @@ impl<'f> Dataset<'f> {
&self, &self,
selection: &clawhdf5_format::selection::Selection, selection: &clawhdf5_format::selection::Selection,
) -> Result<Vec<Vec<T>>, Error> { ) -> Result<Vec<Vec<T>>, Error> {
let raw = self.read_selection(selection)?; self.file.retry(|| {
let dt = self.datatype()?; let raw = self.read_selection(selection)?;
self.file.decode_vlen(&dt, &raw) let dt = self.datatype()?;
self.file.decode_vlen(&dt, &raw)
})
} }
// ----- Selection-based read methods ----- // ----- Selection-based read methods -----
@@ -1714,12 +1782,18 @@ impl<'f> Dataset<'f> {
pub fn attrs_with_errors( pub fn attrs_with_errors(
&self, &self,
) -> Result<(HashMap<String, AttrValue>, Vec<FormatError>), Error> { ) -> Result<(HashMap<String, AttrValue>, Vec<FormatError>), Error> {
with_bytes!(&self.file.data, |d| read_attrs( self.file.retry_attrs(|| {
d, let (attrs, errors, read_errors) = with_bytes!(&self.file.data, |d| {
&self.header, read_attrs_reporting(
self.file.offset_size(), d,
self.file.length_size(), &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 /// The attribute called `name`, or `None` if it has none by that name
@@ -1727,13 +1801,15 @@ impl<'f> Dataset<'f> {
/// that name, found without reading the other attributes when they are /// that name, found without reading the other attributes when they are
/// stored densely. /// stored densely.
pub fn attr(&self, name: &str) -> Result<Option<AttrValue>, Error> { pub fn attr(&self, name: &str) -> Result<Option<AttrValue>, Error> {
with_bytes!(&self.file.data, |d| read_attr( self.file.retry_attrs(|| {
d, with_bytes!(&self.file.data, |d| read_attr_reporting(
&self.header, d,
name, &self.header,
self.file.offset_size(), name,
self.file.length_size(), self.file.offset_size(),
)) self.file.length_size(),
))
})
} }
/// Verify this dataset's content against its stored provenance hash /// Verify this dataset's content against its stored provenance hash
@@ -1753,12 +1829,14 @@ impl<'f> Dataset<'f> {
/// result is not a tamper-evidence or authenticity guarantee. /// result is not a tamper-evidence or authenticity guarantee.
#[cfg(feature = "provenance")] #[cfg(feature = "provenance")]
pub fn verify_provenance(&self) -> Result<clawhdf5_format::provenance::VerifyResult, Error> { pub fn verify_provenance(&self) -> Result<clawhdf5_format::provenance::VerifyResult, Error> {
Ok(clawhdf5_format::provenance::verify_dataset_in( self.file.retry(|| {
&self.file.data, Ok(clawhdf5_format::provenance::verify_dataset_in(
&self.header, &self.file.data,
self.file.offset_size(), &self.header,
self.file.length_size(), self.file.offset_size(),
)?) self.file.length_size(),
)?)
})
} }
/// A header message's payload, resolved through the shared-message /// A header message's payload, resolved through the shared-message
@@ -1768,11 +1846,10 @@ impl<'f> Dataset<'f> {
&self, &self,
msg_type: MessageType, msg_type: MessageType,
) -> Result<Option<std::borrow::Cow<'_, [u8]>>, Error> { ) -> Result<Option<std::borrow::Cow<'_, [u8]>>, Error> {
self.header let msg = self.header.messages.iter().find(|m| m.msg_type == msg_type);
.messages // A shared message is read from the header it lives in.
.iter() self.file.retry(|| {
.find(|m| m.msg_type == msg_type) msg.map(|msg| {
.map(|msg| {
with_bytes!(&self.file.data, |d| { with_bytes!(&self.file.data, |d| {
clawhdf5_format::shared_message::message_data_in( clawhdf5_format::shared_message::message_data_in(
d, d,
@@ -1784,6 +1861,7 @@ impl<'f> Dataset<'f> {
.map_err(Error::Format) .map_err(Error::Format)
}) })
.transpose() .transpose()
})
} }
fn required_payload(&self, msg_type: MessageType) -> Result<std::borrow::Cow<'_, [u8]>, Error> { fn required_payload(&self, msg_type: MessageType) -> Result<std::borrow::Cow<'_, [u8]>, Error> {
+1 -1
View File
@@ -130,7 +130,7 @@ pub(crate) fn is_transient(e: &Error) -> bool {
/// HDF5, a structure that is corrupt behind a valid checksum — is returned /// 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 /// at once: a concurrent write does not cause it, and retrying it only
/// costs up to a second of pauses. /// costs up to a second of pauses.
fn is_transient_format(e: &FormatError) -> bool { pub(crate) fn is_transient_format(e: &FormatError) -> bool {
use FormatError as F; use FormatError as F;
matches!( matches!(
e, e,
+66 -21
View File
@@ -170,16 +170,38 @@ pub(crate) fn read_attrs<S: clawhdf5_format::storage::Storage + ?Sized>(
), ),
crate::Error, 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( let (msgs, errors) = clawhdf5_format::attribute::extract_attributes_tolerant_in(
file_data, file_data,
header, header,
offset_size, offset_size,
length_size, length_size,
)?; )?;
Ok(( let mut read_errors = Vec::new();
attrs_to_map(&msgs, file_data, offset_size, length_size), let map = attrs_to_map_reporting(&msgs, file_data, offset_size, length_size, &mut read_errors);
errors, Ok((map, errors, read_errors))
))
} }
/// The attribute called `name` on the object with header `header`, decoded /// 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, offset_size: u8,
length_size: u8, length_size: u8,
) -> Result<Option<AttrValue>, crate::Error> { ) -> 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, file_data,
header, header,
name, name,
offset_size, offset_size,
length_size, length_size,
)? )?;
else { let Some(msg) = found else {
return Ok(None); return Ok((None, read_errors));
}; };
Ok(attrs_to_map( let value = attrs_to_map_reporting(
std::slice::from_ref(&msg), std::slice::from_ref(&msg),
file_data, file_data,
offset_size, offset_size,
length_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], attrs: &[clawhdf5_format::attribute::AttributeMessage],
file_data: &S, file_data: &S,
offset_size: u8, offset_size: u8,
length_size: u8, length_size: u8,
read_errors: &mut Vec<clawhdf5_format::error::FormatError>,
) -> HashMap<String, AttrValue> { ) -> HashMap<String, AttrValue> {
let mut map = HashMap::new(); let mut map = HashMap::new();
for attr in attrs { 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 // verbatim as `AttrValue::Raw` rather than dropped — a partial
// attribute list with no indication anything is missing is worse than // attribute list with no indication anything is missing is worse than
// an undecoded value. // an undecoded value.
let val = let val = decode_attr_value(attr, file_data, offset_size, length_size, read_errors)
decode_attr_value(attr, file_data, offset_size, length_size).unwrap_or_else(|| { .unwrap_or_else(|| AttrValue::Raw {
AttrValue::Raw { datatype: attr.datatype.clone(),
datatype: attr.datatype.clone(), shape: attr.dataspace.dimensions.clone(),
shape: attr.dataspace.dimensions.clone(), data: attr.raw_data.clone(),
data: attr.raw_data.clone(),
}
}); });
map.insert(attr.name.clone(), val); map.insert(attr.name.clone(), val);
} }
@@ -271,6 +311,7 @@ fn decode_attr_value<S: clawhdf5_format::storage::Storage + ?Sized>(
file_data: &S, file_data: &S,
offset_size: u8, offset_size: u8,
length_size: u8, length_size: u8,
read_errors: &mut Vec<clawhdf5_format::error::FormatError>,
) -> Option<AttrValue> { ) -> Option<AttrValue> {
use clawhdf5_format::datatype::Datatype; use clawhdf5_format::datatype::Datatype;
@@ -310,9 +351,13 @@ fn decode_attr_value<S: clawhdf5_format::storage::Storage + ?Sized>(
Datatype::VariableLength { Datatype::VariableLength {
is_string: true, .. is_string: true, ..
} => { } => {
let strings = attr let strings = match attr.read_vl_strings_in(file_data, offset_size, length_size) {
.read_vl_strings_in(file_data, offset_size, length_size) Ok(strings) => strings,
.ok()?; Err(e) => {
read_errors.push(e);
return None;
}
};
if strings.len() == 1 { if strings.len() == 1 {
Some(AttrValue::String(strings[0].clone())) Some(AttrValue::String(strings[0].clone()))
} else { } else {
Binary file not shown.
+153
View File
@@ -692,3 +692,156 @@ fn a_live_file_reads_consistently_while_an_h5py_swmr_writer_appends() {
"a reader never saw the file grow ({passes:?} passes)" "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());
}
+20 -2
View File
@@ -96,8 +96,11 @@ writer keeps appending. What the format and the library guarantee:
read use the new extent. Like h5py, a handle that is not refreshed keeps 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 its extent; its reads still read the index as it is now, and only return
elements inside that extent. elements inside that extent.
4. **Bounded retries.** In a live file, an operation (open, dataset lookup, 4. **Bounded retries.** In a live file, an operation (open, lookups and
refresh, every read) that fails with an error a concurrent write can 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()` cause is run again from the start, up to `File::swmr_read_attempts()`
times (default 100, libhdf5's default; `set_swmr_read_attempts` changes 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 it), sleeping 1 µs, 2 µs, … up to 10 ms between attempts (under a second
@@ -116,6 +119,16 @@ writer keeps appending. What the format and the library guarantee:
Results are only returned from a run where every structure verified, so Results are only returned from a run where every structure verified, so
a torn metadata read is an error, never data. `File::swmr_retries()` a torn metadata read is an error, never data. `File::swmr_retries()`
counts the retries (libhdf5: `H5Fget_metadata_read_retry_info`). 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 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 as here: correctness rests on the writer's ordering (chunk data before
the index entry, only appends), as for libhdf5's reader. the index entry, only appends), as for libhdf5's reader.
@@ -135,6 +148,11 @@ writer cannot add them), and `MmapFile`/`LazyFile`.
EOF is below the file length. EOF is below the file length.
- `crates/clawhdf5/tests/swmr_interop.rs`: - `crates/clawhdf5/tests/swmr_interop.rs`:
- a copy of a file taken mid-write (fixture) reads like h5py's SWMR reader; - 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 - 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 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, gzip) and a 2-D dataset with two (v2 B-tree), flushing after every step,