diff --git a/CHANGELOG.md b/CHANGELOG.md index 30f7888..2ee4cca 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -34,8 +34,9 @@ Design: `docs/design/swmr.md`. - **`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, - refresh, reads) that fails with an error a concurrent write can cause — +- **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 @@ -52,6 +53,9 @@ Design: `docs/design/swmr.md`. (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 diff --git a/crates/clawhdf5-format/src/attribute.rs b/crates/clawhdf5-format/src/attribute.rs index 8887bbb..301ebbb 100644 --- a/crates/clawhdf5-format/src/attribute.rs +++ b/crates/clawhdf5-format/src/attribute.rs @@ -576,7 +576,14 @@ pub fn find_attribute_in_file( offset_size: u8, length_size: u8, ) -> Result, 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( ) -> Result, 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( + file_data: &S, + header: &ObjectHeader, + name: &str, + offset_size: u8, + length_size: u8, +) -> Result<(Option, Vec), 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( file_data: &S, header: &ObjectHeader, name: &str, offset_size: u8, length_size: u8, + errors: &mut Vec, ) -> Result, FormatError> { let attr_info = find_attribute_info(header, offset_size)?; let dense = attr_info @@ -612,12 +651,10 @@ fn find_attribute_core( 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( )?; 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( 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) diff --git a/crates/clawhdf5/src/reader.rs b/crates/clawhdf5/src/reader.rs index 0cb4c9c..57b81da 100644 --- a/crates/clawhdf5/src/reader.rs +++ b/crates/clawhdf5/src/reader.rs @@ -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 @@ -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( + &self, + mut op: impl FnMut() -> Result<(T, Vec), Error>, + ) -> Result { + 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 @@ -879,13 +909,15 @@ 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, Error> { - with_bytes!(&self.data, |d| crate::vlen::decode_strings( - d, - datatype, - raw, - self.offset_size(), - self.length_size(), - )) + self.retry(|| { + with_bytes!(&self.data, |d| crate::vlen::decode_strings( + d, + datatype, + raw, + self.offset_size(), + self.length_size(), + )) + }) } /// Like [`decode_strings`](Self::decode_strings) for variable-length @@ -896,13 +928,15 @@ impl File { datatype: &Datatype, raw: &[u8], ) -> Result>, Error> { - with_bytes!(self.data.meta()?, |d| crate::vlen::decode_string_bytes( - d, - datatype, - raw, - self.offset_size(), - self.length_size(), - )) + self.retry(|| { + with_bytes!(self.data.meta()?, |d| crate::vlen::decode_string_bytes( + d, + datatype, + raw, + self.offset_size(), + self.length_size(), + )) + }) } /// Decode the variable-length sequences in `raw`, a buffer of elements @@ -914,13 +948,15 @@ impl File { datatype: &Datatype, raw: &[u8], ) -> Result>, Error> { - with_bytes!(self.data.meta()?, |d| crate::vlen::decode_vlen( - d, - datatype, - raw, - self.offset_size(), - self.length_size(), - )) + self.retry(|| { + with_bytes!(self.data.meta()?, |d| crate::vlen::decode_vlen( + d, + datatype, + raw, + self.offset_size(), + self.length_size(), + )) + }) } fn parse_header(&self, address: u64) -> Result { @@ -964,28 +1000,32 @@ pub struct Group<'f> { impl<'f> Group<'f> { /// List the names of datasets in this group. pub fn datasets(&self) -> Result, Error> { - let entries = self.children()?; - let mut names = Vec::new(); - for entry in &entries { - let hdr = self.file.parse_header(entry.object_header_address)?; - if has_message(&hdr, MessageType::DataLayout) { - names.push(entry.name.clone()); + self.file.retry(|| { + let entries = self.children()?; + let mut names = Vec::new(); + for entry in &entries { + let hdr = self.file.parse_header(entry.object_header_address)?; + if has_message(&hdr, MessageType::DataLayout) { + names.push(entry.name.clone()); + } } - } - Ok(names) + Ok(names) + }) } /// List the names of subgroups in this group. pub fn groups(&self) -> Result, Error> { - let entries = self.children()?; - let mut names = Vec::new(); - for entry in &entries { - let hdr = self.file.parse_header(entry.object_header_address)?; - if is_group(&hdr) { - names.push(entry.name.clone()); + self.file.retry(|| { + let entries = self.children()?; + let mut names = Vec::new(); + for entry in &entries { + let hdr = self.file.parse_header(entry.object_header_address)?; + if is_group(&hdr) { + names.push(entry.name.clone()); + } } - } - Ok(names) + Ok(names) + }) } /// Read all attributes of this group. @@ -1006,13 +1046,14 @@ impl<'f> Group<'f> { pub fn attrs_with_errors( &self, ) -> Result<(HashMap, Vec), Error> { - 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() - )) + self.file.retry_attrs(|| { + let hdr = self.file.parse_header(self.address)?; + 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. @@ -1045,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, Error> { - let hdr = self.file.parse_header(self.address)?; - with_bytes!(&self.file.data, |d| read_attr( - d, - &hdr, - name, - self.file.offset_size(), - self.file.length_size(), - )) + self.file.retry_attrs(|| { + let hdr = self.file.parse_header(self.address)?; + 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 @@ -1187,17 +1230,20 @@ impl<'f> Dataset<'f> { /// Read all data as `f64` values. pub fn read_f64(&self) -> Result, Error> { - 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. - if let Ok(Some(bytes)) = self.contiguous_raw() { - return Ok(data_read::read_as_f64(&bytes, &dt)?); - } - if let Some(values) = self.read_chunked_native::()? { - return Ok(values); - } - let raw = self.read_raw()?; - Ok(data_read::read_as_f64(&raw, &dt)?) + 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. + if let Ok(Some(bytes)) = self.contiguous_raw() { + return Ok(data_read::read_as_f64(&bytes, &dt)?); + } + if let Some(values) = self.read_chunked_native::()? { + return Ok(values); + } + let raw = self.read_raw()?; + Ok(data_read::read_as_f64(&raw, &dt)?) + }) } /// Zero-copy read of contiguous native-endian `f64` data. @@ -1207,62 +1253,74 @@ impl<'f> Dataset<'f> { /// /// Read all data as `f32` values. pub fn read_f32(&self) -> Result, Error> { - 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. - if let Ok(Some(bytes)) = self.contiguous_raw() { - return Ok(data_read::read_as_f32(&bytes, &dt)?); - } - if let Some(values) = self.read_chunked_native::()? { - return Ok(values); - } - let raw = self.read_raw()?; - Ok(data_read::read_as_f32(&raw, &dt)?) + 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. + if let Ok(Some(bytes)) = self.contiguous_raw() { + return Ok(data_read::read_as_f32(&bytes, &dt)?); + } + if let Some(values) = self.read_chunked_native::()? { + return Ok(values); + } + 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, Error> { - 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. - if let Ok(Some(bytes)) = self.contiguous_raw() { - return Ok(data_read::read_as_i32(&bytes, &dt)?); - } - if let Some(values) = self.read_chunked_native::()? { - return Ok(values); - } - let raw = self.read_raw()?; - Ok(data_read::read_as_i32(&raw, &dt)?) + 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. + if let Ok(Some(bytes)) = self.contiguous_raw() { + return Ok(data_read::read_as_i32(&bytes, &dt)?); + } + if let Some(values) = self.read_chunked_native::()? { + return Ok(values); + } + 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, Error> { - 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. - if let Ok(Some(bytes)) = self.contiguous_raw() { - return Ok(data_read::read_as_i64(&bytes, &dt)?); - } - if let Some(values) = self.read_chunked_native::()? { - return Ok(values); - } - let raw = self.read_raw()?; - Ok(data_read::read_as_i64(&raw, &dt)?) + 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. + if let Ok(Some(bytes)) = self.contiguous_raw() { + return Ok(data_read::read_as_i64(&bytes, &dt)?); + } + if let Some(values) = self.read_chunked_native::()? { + return Ok(values); + } + 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, Error> { - 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. - if let Ok(Some(bytes)) = self.contiguous_raw() { - return Ok(data_read::read_as_u64(&bytes, &dt)?); - } - if let Some(values) = self.read_chunked_native::()? { - return Ok(values); - } - let raw = self.read_raw()?; - Ok(data_read::read_as_u64(&raw, &dt)?) + 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. + if let Ok(Some(bytes)) = self.contiguous_raw() { + return Ok(data_read::read_as_u64(&bytes, &dt)?); + } + if let Some(values) = self.read_chunked_native::()? { + return Ok(values); + } + let raw = self.read_raw()?; + Ok(data_read::read_as_u64(&raw, &dt)?) + }) } /// 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 /// [`read_string_bytes`](Self::read_string_bytes) for the exact bytes. pub fn read_string(&self) -> Result, Error> { - let raw = self.read_raw()?; - let dt = self.datatype()?; - self.file.decode_strings(&dt, &raw) + 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>, Error> { - let raw = self.read_raw()?; - let dt = self.datatype()?; - self.file.decode_string_bytes(&dt, &raw) + 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 @@ -1292,9 +1354,11 @@ impl<'f> Dataset<'f> { &self, selection: &clawhdf5_format::selection::Selection, ) -> Result, Error> { - let raw = self.read_selection(selection)?; - let dt = self.datatype()?; - self.file.decode_strings(&dt, &raw) + 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 @@ -1303,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(&self) -> Result>, Error> { - let raw = self.read_raw()?; - let dt = self.datatype()?; - self.file.decode_vlen(&dt, &raw) + 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 @@ -1314,9 +1380,11 @@ impl<'f> Dataset<'f> { &self, selection: &clawhdf5_format::selection::Selection, ) -> Result>, Error> { - let raw = self.read_selection(selection)?; - let dt = self.datatype()?; - self.file.decode_vlen(&dt, &raw) + self.file.retry(|| { + let raw = self.read_selection(selection)?; + let dt = self.datatype()?; + self.file.decode_vlen(&dt, &raw) + }) } // ----- Selection-based read methods ----- @@ -1714,12 +1782,18 @@ impl<'f> Dataset<'f> { pub fn attrs_with_errors( &self, ) -> Result<(HashMap, Vec), Error> { - with_bytes!(&self.file.data, |d| read_attrs( - d, - &self.header, - self.file.offset_size(), - self.file.length_size(), - )) + 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 @@ -1727,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, Error> { - with_bytes!(&self.file.data, |d| read_attr( - d, - &self.header, - name, - self.file.offset_size(), - self.file.length_size(), - )) + 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 @@ -1753,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 { - Ok(clawhdf5_format::provenance::verify_dataset_in( - &self.file.data, - &self.header, - self.file.offset_size(), - self.file.length_size(), - )?) + 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 @@ -1768,11 +1846,10 @@ impl<'f> Dataset<'f> { &self, msg_type: MessageType, ) -> Result>, 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, @@ -1784,6 +1861,7 @@ impl<'f> Dataset<'f> { .map_err(Error::Format) }) .transpose() + }) } fn required_payload(&self, msg_type: MessageType) -> Result, Error> { diff --git a/crates/clawhdf5/src/swmr.rs b/crates/clawhdf5/src/swmr.rs index bf0112d..ef24933 100644 --- a/crates/clawhdf5/src/swmr.rs +++ b/crates/clawhdf5/src/swmr.rs @@ -130,7 +130,7 @@ pub(crate) fn is_transient(e: &Error) -> bool { /// 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. -fn is_transient_format(e: &FormatError) -> bool { +pub(crate) fn is_transient_format(e: &FormatError) -> bool { use FormatError as F; matches!( e, diff --git a/crates/clawhdf5/src/types.rs b/crates/clawhdf5/src/types.rs index b3f7823..77bbe63 100644 --- a/crates/clawhdf5/src/types.rs +++ b/crates/clawhdf5/src/types.rs @@ -170,16 +170,38 @@ pub(crate) fn read_attrs( ), 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, + Vec, + Vec, +); + +/// [`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( + file_data: &S, + header: &clawhdf5_format::object_header::ObjectHeader, + offset_size: u8, + length_size: u8, +) -> Result { 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( offset_size: u8, length_size: u8, ) -> Result, 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( + file_data: &S, + header: &clawhdf5_format::object_header::ObjectHeader, + name: &str, + offset_size: u8, + length_size: u8, +) -> Result<(Option, Vec), 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( +/// 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( attrs: &[clawhdf5_format::attribute::AttributeMessage], file_data: &S, offset_size: u8, length_size: u8, + read_errors: &mut Vec, ) -> HashMap { let mut map = HashMap::new(); for attr in attrs { @@ -224,13 +266,11 @@ pub(crate) fn attrs_to_map( // 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 { - datatype: attr.datatype.clone(), - shape: attr.dataspace.dimensions.clone(), - data: attr.raw_data.clone(), - } + 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( file_data: &S, offset_size: u8, length_size: u8, + read_errors: &mut Vec, ) -> Option { use clawhdf5_format::datatype::Datatype; @@ -310,9 +351,13 @@ fn decode_attr_value( 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 { diff --git a/crates/clawhdf5/tests/fixtures/swmr_strings_attrs.h5 b/crates/clawhdf5/tests/fixtures/swmr_strings_attrs.h5 new file mode 100644 index 0000000..cd9cf29 Binary files /dev/null and b/crates/clawhdf5/tests/fixtures/swmr_strings_attrs.h5 differ diff --git a/crates/clawhdf5/tests/swmr_interop.rs b/crates/clawhdf5/tests/swmr_interop.rs index ca97931..beb3509 100644 --- a/crates/clawhdf5/tests/swmr_interop.rs +++ b/crates/clawhdf5/tests/swmr_interop.rs @@ -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)" ); } + +// --------------------------------------------------------------------------- +// 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, + 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, 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| { + format!( + "{:?}", + m.into_iter().collect::>() + ) + }; + type Op<'a> = Box Result + '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::().map(|n| format!("{n:?}"))), + ), + ( + "v.read_vlen_selection", + Box::new(|| { + v.read_vlen_selection::(&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()); +} diff --git a/docs/design/swmr.md b/docs/design/swmr.md index c797091..dba542a 100644 --- a/docs/design/swmr.md +++ b/docs/design/swmr.md @@ -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 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, dataset lookup, - refresh, every read) that fails with an error a concurrent write can +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 @@ -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 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. @@ -135,6 +148,11 @@ writer cannot add them), and `MmapFile`/`LazyFile`. 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,