#![cfg(feature = "threadsafe")]
use rust_hdf5::H5File;
fn tmp(name: &str) -> std::path::PathBuf {
use std::sync::atomic::{AtomicU64, Ordering};
static COUNTER: AtomicU64 = AtomicU64::new(0);
let n = COUNTER.fetch_add(1, Ordering::Relaxed);
std::env::temp_dir().join(format!(
"rust_hdf5_same_ds_{}_{}_{}.h5",
name,
std::process::id(),
n
))
}
fn tag(k: usize, j: usize, f: usize) -> i32 {
(k * 1_000_000 + j * 100 + f) as i32
}
#[test]
fn concurrent_same_dataset_appends_serialize_wholly() {
const N: usize = 4; const M: usize = 40; const FRAMES: usize = 3; const W: usize = 8; const CHUNK0: usize = 4; const ITERS: usize = 10;
for iter in 0..ITERS {
let path = tmp(&format!("iter{iter}"));
{
let file = H5File::create(&path).unwrap();
let ds = file
.new_dataset::<i32>()
.shape([0, W])
.chunk(&[CHUNK0, W])
.max_shape(&[None, Some(W)])
.create("d")
.unwrap();
std::thread::scope(|s| {
for k in 0..N {
let ds = &ds;
s.spawn(move || {
for j in 0..M {
let data: Vec<i32> = (0..FRAMES)
.flat_map(|f| std::iter::repeat_n(tag(k, j, f), W))
.collect();
ds.append(&data)
.unwrap_or_else(|e| panic!("append k={k} j={j}: {e}"));
}
});
}
});
file.close().unwrap();
}
let file = H5File::open(&path).unwrap();
let ds = file.dataset("d").unwrap();
assert_eq!(ds.shape(), vec![N * M * FRAMES, W], "iter {iter}");
let got = ds.read_raw::<i32>().unwrap();
let rows: Vec<i32> = got
.chunks_exact(W)
.enumerate()
.map(|(r, row)| {
assert!(
row.iter().all(|&v| v == row[0]),
"iter {iter}: torn frame at row {r}: {row:?}"
);
row[0]
})
.collect();
let mut seen = std::collections::HashSet::new();
let mut r = 0;
while r < rows.len() {
let first = rows[r];
let (k, rem) = ((first / 1_000_000) as usize, first % 1_000_000);
let (j, f) = ((rem / 100) as usize, (rem % 100) as usize);
assert_eq!(f, 0, "iter {iter}: run at row {r} starts mid-call: {first}");
for f in 0..FRAMES {
assert_eq!(
rows[r + f],
tag(k, j, f),
"iter {iter}: call (k={k}, j={j}) split at row {}",
r + f
);
}
assert!(
seen.insert((k, j)),
"iter {iter}: call (k={k}, j={j}) appended twice"
);
r += FRAMES;
}
assert_eq!(seen.len(), N * M, "iter {iter}");
std::fs::remove_file(&path).ok();
}
}
#[test]
fn concurrent_same_dataset_vlen_appends_serialize_wholly() {
const N: usize = 4; const M: usize = 30; const STRINGS: usize = 3; const CHUNK: usize = 4; const ITERS: usize = 10;
for iter in 0..ITERS {
let path = tmp(&format!("vlen{iter}"));
{
let file = H5File::create(&path).unwrap();
file.create_appendable_vlen_dataset("strs", CHUNK, None)
.unwrap();
std::thread::scope(|s| {
for k in 0..N {
let file = &file;
s.spawn(move || {
for j in 0..M {
let batch: Vec<String> =
(0..STRINGS).map(|f| format!("{k}:{j}:{f}")).collect();
let refs: Vec<&str> = batch.iter().map(String::as_str).collect();
file.append_vlen_strings("strs", &refs)
.unwrap_or_else(|e| panic!("append k={k} j={j}: {e}"));
}
});
}
});
file.close().unwrap();
}
let file = H5File::open(&path).unwrap();
let ds = file.dataset("strs").unwrap();
let got = ds.read_vlen_strings().unwrap();
assert_eq!(got.len(), N * M * STRINGS, "iter {iter}");
let mut seen = std::collections::HashSet::new();
let mut r = 0;
while r < got.len() {
let parts: Vec<usize> = got[r].split(':').map(|p| p.parse().unwrap()).collect();
let (k, j, f) = (parts[0], parts[1], parts[2]);
assert_eq!(
f, 0,
"iter {iter}: run at element {r} starts mid-call: {}",
got[r]
);
for f in 0..STRINGS {
assert_eq!(
got[r + f],
format!("{k}:{j}:{f}"),
"iter {iter}: call (k={k}, j={j}) split at element {}",
r + f
);
}
assert!(
seen.insert((k, j)),
"iter {iter}: call (k={k}, j={j}) appended twice"
);
r += STRINGS;
}
assert_eq!(seen.len(), N * M, "iter {iter}");
std::fs::remove_file(&path).ok();
}
}