use std::collections::HashMap;
use std::path::Path;
use crate::chunk_index_inplace::{InPlaceFile, Located, apply_ea_append, plan_ea_append};
use crate::datatype::Datatype;
use crate::edit::{AppendBuilder, datatype_is_raw_appendable, pipeline_reencodable};
use crate::element::H5Element;
use crate::error::Error;
use crate::file_lock::FileLocking;
use crate::filter_pipeline::FilterPipeline;
use crate::group_v2;
use crate::source::Source;
struct DatasetState {
loc: Located,
datatype: Datatype,
spatial: Vec<u64>,
element_size: usize,
pipeline: Option<FilterPipeline>,
}
#[deprecated(
since = "0.22.0",
note = "use File::open_rw + Dataset::append; see the AppendWriter type docs for migration"
)]
pub struct AppendWriter {
file: InPlaceFile,
datasets: HashMap<String, DatasetState>,
}
#[allow(deprecated)] impl AppendWriter {
pub fn open<P: AsRef<Path>>(path: P) -> Result<Self, Error> {
Self::open_with_locking(path, FileLocking::Enabled)
}
pub fn open_with_locking<P: AsRef<Path>>(path: P, locking: FileLocking) -> Result<Self, Error> {
let file = InPlaceFile::open(path, Some(locking), Error::AppendUnsupported)?;
Ok(Self {
file,
datasets: HashMap::new(),
})
}
pub fn append_raw(&mut self, dataset: &str, bytes: &[u8]) -> Result<(), Error> {
let mut b = AppendBuilder::new();
b.append_raw(bytes);
self.append_gathered(dataset, &b, 4)
}
pub fn append<T: H5Element>(&mut self, dataset: &str, data: &[T]) -> Result<(), Error> {
let mut b = AppendBuilder::new();
b.append(data);
self.append_gathered(dataset, &b, 4)
}
pub fn close(mut self) -> Result<(), Error> {
self.file.sync()
}
fn ensure_located(&mut self, dataset: &str) -> Result<(), Error> {
if self.datasets.contains_key(dataset) {
return Ok(());
}
let oh_addr = group_v2::resolve_path_any(self.file.data(), &self.file.superblock, dataset)?;
let result = Located::locate_at(&self.file, oh_addr, Error::AppendUnsupported)?;
if result.located.chunk_elems == 0 {
return Err(Error::AppendUnsupported(
"append requires a nonzero chunk length",
));
}
let (dt_off, dt_size) = result.spans.datatype;
let dt_bytes = self
.file
.read_metadata_at(dt_off, dt_size)
.map_err(|_| Error::AppendUnsupported("dataset datatype could not be parsed"))?;
let (datatype, _) = Datatype::parse(&dt_bytes)
.map_err(|_| Error::AppendUnsupported("dataset datatype could not be parsed"))?;
let pipeline = match result.spans.filter {
Some((fb, fsize)) => {
let fp_bytes = self.file.read_metadata_at(fb, fsize).map_err(|_| {
Error::AppendUnsupported("dataset filter pipeline could not be parsed")
})?;
let parsed = FilterPipeline::parse(&fp_bytes).map_err(|_| {
Error::AppendUnsupported("dataset filter pipeline could not be parsed")
})?;
if !pipeline_reencodable(&parsed) {
return Err(Error::AppendUnsupported(
"dataset uses a filter this engine cannot re-encode",
));
}
Some(parsed)
}
None => None,
};
let element_size = result.located.elem_bytes;
let spatial = vec![result.located.chunk_elems];
self.datasets.insert(
dataset.to_string(),
DatasetState {
loc: result.located,
datatype,
spatial,
element_size,
pipeline,
},
);
Ok(())
}
fn append_gathered(
&mut self,
dataset: &str,
b: &AppendBuilder,
max_phase: u8,
) -> Result<(), Error> {
if b.dt_conflict() {
return Err(Error::AppendUnsupported(
"append mixes element types in one call; use one element type per append",
));
}
self.ensure_located(dataset)?;
let raw = b.raw();
{
let st = &self.datasets[dataset];
if raw.len() % st.element_size != 0 {
return Err(Error::AppendUnsupported(
"appended byte length is not a whole number of elements",
));
}
match b.elem_dt() {
Some(expected) if *expected != st.datatype => {
return Err(Error::AppendUnsupported(
"append datatype does not match the on-disk dataset (wrong element \
type or byte order)",
));
}
Some(_) => {}
None => {
if !datatype_is_raw_appendable(&st.datatype) {
return Err(Error::AppendUnsupported(
"append_raw onto this dataset's datatype (non-little-endian, \
variable-length, or reference) could misencode the bytes; use a \
typed append",
));
}
}
}
}
let new_elems = (raw.len() / self.datasets[dataset].element_size) as u64;
if new_elems == 0 {
return Ok(());
}
let plan = {
let st = &self.datasets[dataset];
plan_ea_append(
&self.file,
&st.loc,
&st.datatype,
&st.spatial,
st.element_size,
st.pipeline.as_ref(),
raw,
new_elems,
)?
};
let st = self
.datasets
.get_mut(dataset)
.expect("dataset was located above");
apply_ea_append(&mut self.file, &mut st.loc, &plan, max_phase)
}
#[cfg(test)]
fn append_i32_phased(
&mut self,
dataset: &str,
values: &[i32],
max_phase: u8,
) -> Result<(), Error> {
let mut b = AppendBuilder::new();
b.append_i32(values);
self.append_gathered(dataset, &b, max_phase)
}
}
macro_rules! append_typed {
($($method:ident, $ty:ty;)*) => {
#[allow(deprecated)] impl AppendWriter {
$(
#[doc = concat!("Append `", stringify!($ty), "` values to `dataset`. \
A filtered dataset accepts only chunk-aligned appends; see \
[`AppendWriter`] for the full contract.")]
pub fn $method(&mut self, dataset: &str, data: &[$ty]) -> Result<(), Error> {
let mut b = AppendBuilder::new();
b.$method(data);
self.append_gathered(dataset, &b, 4)
}
)*
}
};
}
append_typed! {
append_f64, f64;
append_f32, f32;
append_i8, i8;
append_i16, i16;
append_i32, i32;
append_i64, i64;
append_u8, u8;
append_u16, u16;
append_u32, u32;
append_u64, u64;
}
#[cfg(test)]
#[allow(deprecated)] mod tests {
use super::*;
use crate::reader::File as PureFile;
use crate::writer::FileBuilder;
use tempfile::tempdir;
fn build_unfiltered(path: &std::path::Path, n: i32, chunk: u64) {
let data: Vec<i32> = (0..n).collect();
let mut b = FileBuilder::new();
b.create_dataset("d")
.with_i32_data(&data)
.with_shape(&[n as u64])
.with_maxshape(&[u64::MAX])
.with_chunks(&[chunk]);
b.write(path).unwrap();
}
#[test]
fn crash_consistency_partial_tail_prefix() {
for (n, chunk, add) in [(6i32, 4u64, 5i32), (9, 2, 6)] {
let dir = tempdir().unwrap();
let base = dir.path().join("base.h5");
build_unfiltered(&base, n, chunk);
for max_phase in 1u8..=4 {
let p = dir.path().join(format!("crash_{n}_{chunk}_{max_phase}.h5"));
std::fs::copy(&base, &p).unwrap();
{
let mut w = AppendWriter::open(&p).unwrap();
w.append_i32_phased("d", &(n..n + add).collect::<Vec<_>>(), max_phase)
.unwrap();
}
let expected_len = if max_phase == 4 { n + add } else { n };
let f = PureFile::from_bytes(std::fs::read(&p).unwrap()).unwrap();
assert_eq!(
f.dataset("d").unwrap().read_i32().unwrap(),
(0..expected_len).collect::<Vec<_>>(),
"inconsistent view after crash at phase {max_phase} (n={n}, chunk={chunk})"
);
}
}
}
#[test]
#[cfg(not(target_pointer_width = "32"))]
fn crash_consistency_c_reads_partial_tail_prefix() {
let dir = tempdir().unwrap();
let base = dir.path().join("base.h5");
let n = 6i32;
let add = 5i32; build_unfiltered(&base, n, 4);
for max_phase in 1u8..=4 {
let p = dir.path().join(format!("crash_c_{max_phase}.h5"));
std::fs::copy(&base, &p).unwrap();
{
let mut w = AppendWriter::open(&p).unwrap();
w.append_i32_phased("d", &(n..n + add).collect::<Vec<_>>(), max_phase)
.unwrap();
}
let expected_len = if max_phase == 4 { n + add } else { n };
let f = hdf5::File::open(&p).unwrap();
let v = f.dataset("d").unwrap().read_raw::<i32>().unwrap();
assert_eq!(
v,
(0..expected_len).collect::<Vec<_>>(),
"C library saw an inconsistent view after crash at phase {max_phase}"
);
f.close().unwrap();
}
}
}