use std::path::{Path, PathBuf};
use crate::dataset::{DatasetAccess, VirtualView};
use crate::format::btree_v1::{BTreeV1Config, BTreeV1Node, ChunkBTreeV1Node};
use crate::format::bytes::read_le_uint as read_uint;
use crate::format::creation_order::CreationOrder;
use crate::format::fractal_heap::{self, FractalHeapHeader};
use crate::format::global_heap::{
decode_vlen_reference, vlen_reference_size, GlobalHeapCollection,
};
use crate::format::local_heap::{local_heap_get_string, LocalHeapHeader};
use crate::format::messages::attr_info::AttributeInfoMessage;
use crate::format::messages::attribute::{AttributeEntry, AttributeMessage};
use crate::format::messages::data_layout::{self, DataLayoutMessage};
use crate::format::messages::dataspace::DataspaceMessage;
use crate::format::messages::datatype::{DatatypeMessage, OldReferenceKind, ReferenceEncoding};
use crate::format::messages::external_file_list::ExternalFileListMessage;
use crate::format::messages::fill_value::{
try_tiled_fill, FillValueMessage, ALLOC_TIME_LATE, FILL_TIME_IFSET,
};
use crate::format::messages::filter::{self, FilterPipeline};
use crate::format::messages::link::LinkMessage;
use crate::format::messages::link::LinkTarget;
use crate::format::messages::link_info::LinkInfoMessage;
use crate::format::messages::shared::{MessageStorage, MSG_FLAG_SHARED};
use crate::format::messages::superblock_ext::{
BtreeKMessage, DriverInfoMessage, FileSpaceInfoMessage, SharedMessageTableMessage,
};
use crate::format::messages::virtual_mapping::{
parse_source_name, VirtualMapping, VirtualMappingList,
};
use crate::format::messages::*;
use crate::format::object_header::ObjectHeader;
use crate::format::reference::{
decode_object_element, decode_region_element, decode_region_heap_object, decode_revised_body,
decode_revised_element, DecodedReference, Reference, ReferenceTarget, RevisedElement,
};
use crate::format::selection::{
Hyperslab, PointSelection, RegularHyperslab, ResolvedSelection, Selection,
};
use crate::format::sohm::SohmMasterTable;
use crate::format::storage_kind::{AttributeStorage, LinkStorage};
use crate::format::superblock::{
detect_superblock_version, SuperblockV0V1, SuperblockV2V3, SymbolTableCache,
};
use crate::format::symbol_table::SymbolTableNode;
use crate::format::{BlockReader, FormatContext, UNDEF_ADDR};
use crate::format::selection::check_hyperslab;
use crate::io::file_handle::{FileHandle, ReadDst};
use crate::io::hyperslab::{compute_strides, for_each_contiguous_run};
use crate::io::locking::FileLocking;
use crate::io::{FileMeta, IoResult};
struct ChunkIndexDesc<'a> {
index_type: data_layout::ChunkIndexType,
index_address: u64,
earray_params: Option<&'a data_layout::EarrayParams>,
single_chunk_filter: Option<data_layout::SingleChunkFilter>,
}
type Bt2ChunkEntry = (u64, usize, Vec<u64>, u32);
#[derive(Clone, Copy)]
enum ChunkTarget<'a> {
Full,
Slice {
starts: &'a [u64],
counts: &'a [u64],
},
}
impl<'a> ChunkTarget<'a> {
fn overlaps(&self, coords: &[u64], chunk_dims: &[u64]) -> bool {
match self {
ChunkTarget::Full => true,
ChunkTarget::Slice { starts, counts } => coords.iter().enumerate().all(|(d, &c)| {
let origin = c.saturating_mul(chunk_dims[d]);
let chunk_end = origin.saturating_add(chunk_dims[d]);
let sel_end = starts[d].saturating_add(counts[d]);
origin < sel_end && starts[d] < chunk_end
}),
}
}
}
#[derive(Clone, Copy)]
struct ChunkReadRequest<'a> {
pipeline: Option<&'a FilterPipeline>,
target: ChunkTarget<'a>,
fill_value: Option<&'a [u8]>,
dst: ReadDst,
}
#[derive(Clone, Copy)]
struct ChunkOutputGeometry<'a> {
dims: &'a [u64],
chunk_dims: &'a [u64],
element_size: u64,
}
#[derive(Clone, Copy)]
struct ChunkPlacement<'a> {
geo: ChunkOutputGeometry<'a>,
starts: &'a [u64],
counts: &'a [u64],
}
impl ChunkOutputGeometry<'_> {
fn image_bytes(&self) -> Option<u64> {
self.chunk_dims
.iter()
.copied()
.try_fold(1u64, |a, d| a.checked_mul(d))
.and_then(|elems| elems.checked_mul(self.element_size))
.filter(|b| *b > 0)
}
fn resident_bytes(&self, coords: &[u64]) -> u64 {
let mut elems = 1u64;
for (d, &c) in coords.iter().enumerate().take(self.dims.len()) {
let origin = c.saturating_mul(self.chunk_dims[d]);
let end = origin.saturating_add(self.chunk_dims[d]).min(self.dims[d]);
elems = elems.saturating_mul(end.saturating_sub(origin));
}
elems.saturating_mul(self.element_size)
}
}
impl<'a> ChunkPlacement<'a> {
fn resolve(geo: &ChunkOutputGeometry<'a>, target: ChunkTarget<'a>, zeros: &'a [u64]) -> Self {
let (starts, counts) = match target {
ChunkTarget::Full => (zeros, geo.dims),
ChunkTarget::Slice { starts, counts } => (starts, counts),
};
ChunkPlacement {
geo: *geo,
starts,
counts,
}
}
fn leaves_chunk_unconsumed(&self, coords: &[u64]) -> bool {
match ChunkOverlap::of(self, coords) {
Some(overlap) => overlap.bytes(self.geo.element_size) < self.geo.resident_bytes(coords),
None => false,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ExternalFileSegment {
pub name: String,
pub offset: u64,
pub size: u64,
}
#[cfg(feature = "mmap")]
pub(crate) enum ViewStorage {
Contiguous { offset: u64, len: u64 },
Unallocated,
Elsewhere(&'static str),
}
#[cfg(feature = "mmap")]
pub(crate) struct DatasetViewSource {
pub map: Option<std::sync::Arc<crate::io::file_handle::LockedMap>>,
pub storage: ViewStorage,
pub datatype: DatatypeMessage,
pub dims: Vec<u64>,
}
pub struct DatasetReadInfo {
pub name: String,
pub object_header_address: u64,
pub datatype: DatatypeMessage,
pub dataspace: DataspaceMessage,
pub layout: DataLayoutMessage,
pub filter_pipeline: Option<FilterPipeline>,
pub attributes: ObjectAttributes,
pub fill_value: Option<Vec<u8>>,
pub fill_defined: u8,
pub fill_write_time: u8,
pub alloc_time: u8,
pub external_files: Vec<ExternalFileSegment>,
pub virtual_mappings: Option<VirtualMappingList>,
pub virtual_resolution: Option<Vec<MappingResolution>>,
pub virtual_stored_dims: Option<Vec<u64>>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum MappingResolution {
Bounded,
Unlimited { virtual_clip: u64, source_clip: u64 },
Printf { blocks: u64, present: Vec<u64> },
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum LinkClass {
Hard,
Soft { path: String },
External { file: String, path: String },
UserDefined { link_type: u8 },
}
impl LinkClass {
pub(crate) fn from_target(target: &LinkTarget) -> Self {
match target {
LinkTarget::Hard { .. } => Self::Hard,
LinkTarget::Soft { target } => Self::Soft {
path: target.clone(),
},
LinkTarget::External { file, path } => Self::External {
file: file.clone(),
path: path.clone(),
},
LinkTarget::UserDefined { link_type, .. } => Self::UserDefined {
link_type: *link_type,
},
}
}
}
struct SoftLinkRef {
link: String,
target: String,
}
pub(crate) struct ExternalEdge {
pub link: String,
pub file: String,
pub path: String,
}
impl ExternalEdge {
fn dangling(&self) -> crate::io::IoError {
crate::io::IoError::DanglingLink {
link: self.link.clone(),
target: format!("{}::{}", self.file, self.path),
}
}
}
struct Origin {
path: PathBuf,
locking: crate::io::locking::FileLocking,
}
const MAX_EXTERNAL_HOPS: usize = 16;
enum Traversal {
Path {
path: String,
via: Option<SoftLinkRef>,
},
External {
link: String,
file: String,
path: String,
},
}
#[derive(Clone, Copy)]
enum Rewrite<'a> {
Alias(&'a str),
Soft(&'a str),
External { file: &'a str, path: &'a str },
}
fn resolve_link_value(link_path: &str, value: &str) -> String {
let mut components: Vec<&str> = Vec::new();
if !value.starts_with('/') {
if let Some(parent) = link_path.rsplit_once('/').map(|(p, _)| p) {
components.extend(parent.split('/').filter(|c| !c.is_empty()));
}
}
for component in value.split('/') {
match component {
"" | "." => {}
".." => {
components.pop();
}
c => components.push(c),
}
}
components.join("/")
}
#[derive(Default)]
struct Catalog {
datasets: Vec<DatasetReadInfo>,
unreadable: std::collections::BTreeMap<String, String>,
group_attributes: std::collections::HashMap<String, ObjectAttributes>,
group_link_storage: std::collections::HashMap<String, (LinkStorage, CreationOrder)>,
group_paths: std::collections::BTreeSet<String>,
group_object_paths: std::collections::HashMap<u64, String>,
group_aliases: std::collections::HashMap<String, String>,
links: std::collections::BTreeMap<String, LinkClass>,
datatypes: std::collections::BTreeMap<String, CommittedDatatypeInfo>,
datatype_object_paths: std::collections::HashMap<u64, String>,
}
impl Catalog {
fn object_paths(&self, root_addr: u64) -> std::collections::HashMap<u64, String> {
let mut paths = std::collections::HashMap::new();
paths.insert(root_addr, "/".to_string());
for (addr, path) in &self.group_object_paths {
paths.insert(*addr, absolute_path(path));
}
for ds in &self.datasets {
paths.insert(ds.object_header_address, absolute_path(&ds.name));
}
for (addr, path) in &self.datatype_object_paths {
paths.insert(*addr, absolute_path(path));
}
paths
}
}
struct CatalogWalk<'a> {
handle: &'a mut FileHandle,
meta: &'a FileMeta,
catalog: Catalog,
visited: std::collections::HashMap<u64, String>,
}
impl<'a> CatalogWalk<'a> {
const MAX_DEPTH: usize = 256;
fn new(handle: &'a mut FileHandle, meta: &'a FileMeta, root_addr: u64) -> Self {
let mut visited = std::collections::HashMap::new();
visited.insert(root_addr, String::new());
Self {
handle,
meta,
catalog: Catalog::default(),
visited,
}
}
fn ctx(&self) -> &FormatContext {
&self.meta.ctx
}
fn finish(mut self) -> Catalog {
self.catalog.group_object_paths = self.visited;
self.catalog
}
fn group(
&mut self,
header: Option<&ObjectHeader>,
prefix: &str,
depth: usize,
stab: Option<(u64, u64)>,
) -> IoResult<()> {
if depth > Self::MAX_DEPTH {
return Ok(());
}
let link_storage = header.filter(|h| header_declares_link_storage(h));
if let Some(h) = link_storage {
return self.links(h, prefix, depth);
}
let (btree_addr, heap_addr) = match stab {
Some(pair) if pair.0 != UNDEF_ADDR && pair.1 != UNDEF_ADDR => pair,
_ => header.map_or((UNDEF_ADDR, UNDEF_ADDR), |h| {
Hdf5Reader::stab_from_header(h, self.ctx())
}),
};
if btree_addr != UNDEF_ADDR && heap_addr != UNDEF_ADDR {
self.btree(btree_addr, heap_addr, prefix, depth)?;
}
Ok(())
}
fn links(&mut self, header: &ObjectHeader, prefix: &str, depth: usize) -> IoResult<()> {
let mut links: Vec<LinkMessage> = Vec::new();
for msg in &header.messages {
if msg.msg_type == MSG_LINK {
let (link, _) = LinkMessage::decode(&msg.data, self.ctx())?;
links.push(link);
} else if msg.msg_type == MSG_LINK_INFO {
let (info, _) = LinkInfoMessage::decode(&msg.data, self.ctx())?;
if info.fractal_heap_address != UNDEF_ADDR {
let ctx = self.meta.ctx;
let dense =
Hdf5Reader::read_dense_links(self.handle, &ctx, info.fractal_heap_address)?;
links.extend(dense);
}
}
}
for link in &links {
let full_name = join_path(prefix, &link.name);
self.catalog
.links
.insert(full_name.clone(), LinkClass::from_target(&link.target));
let LinkTarget::Hard { address } = &link.target else {
continue;
};
self.child(full_name, *address, depth, None)?;
}
Ok(())
}
fn btree(
&mut self,
btree_addr: u64,
heap_addr: u64,
prefix: &str,
depth: usize,
) -> IoResult<()> {
let sa = self.ctx().sizeof_addr as usize;
let ss = self.ctx().sizeof_size as usize;
let heap_hdr_buf = self.handle.read_at_most(heap_addr, 64)?;
let heap_hdr = LocalHeapHeader::decode(&heap_hdr_buf, sa, ss)?;
let heap_data = self
.handle
.read_at(heap_hdr.data_addr, heap_hdr.data_size as usize)?;
let mut snod_tree_visited = std::collections::HashSet::new();
let snod_addrs = Hdf5Reader::collect_snod_addresses(
self.handle,
self.meta,
btree_addr,
0,
&mut snod_tree_visited,
)?;
let snod_size = self.meta.btree.symbol_table_node_size(sa, ss);
for snod_addr in snod_addrs {
let snod_buf = self.handle.read_at_most(snod_addr, snod_size)?;
let snod =
SymbolTableNode::decode(&snod_buf, sa, ss, self.meta.btree.sym_leaf_max_entries())?;
for entry in &snod.entries {
let name = local_heap_get_string(&heap_data, entry.name_offset)?;
if name.is_empty() {
continue;
}
let full_name = join_path(prefix, &name);
if let SymbolTableCache::SoftLink { value_offset } = entry.cache {
let target = local_heap_get_string(&heap_data, value_offset as u64)?;
self.catalog
.links
.insert(full_name, LinkClass::Soft { path: target });
continue;
}
self.catalog
.links
.insert(full_name.clone(), LinkClass::Hard);
self.child(
full_name,
entry.obj_header_addr,
depth,
entry.cached_symbol_table(),
)?;
}
}
Ok(())
}
fn child(
&mut self,
full_name: String,
addr: u64,
depth: usize,
stab: Option<(u64, u64)>,
) -> IoResult<()> {
let header = match Hdf5Reader::read_object_header_full(self.handle, self.meta, addr) {
Ok(h) => h,
Err(e) => {
self.catalog
.unreadable
.insert(full_name, format!("its object header does not decode: {e}"));
return Ok(());
}
};
match Hdf5Reader::classify_object(self.handle, &header, self.meta, &full_name, addr) {
ObjectKind::Dataset(info) => {
self.catalog.datasets.push(*info);
return Ok(());
}
ObjectKind::UnreadableDataset(why) => {
self.catalog.unreadable.insert(full_name, why);
return Ok(());
}
ObjectKind::CommittedDatatype(info) => {
self.catalog
.datatype_object_paths
.insert(addr, full_name.clone());
self.catalog.datatypes.insert(full_name, *info);
return Ok(());
}
ObjectKind::Group => {}
}
self.catalog.group_paths.insert(full_name.clone());
let ctx = self.meta.ctx;
let attrs = collect_object_attributes(self.handle, &ctx, &header);
self.catalog
.group_attributes
.insert(full_name.clone(), attrs);
self.catalog.group_link_storage.insert(
full_name.clone(),
describe_link_storage(Some(&header), &ctx, stab),
);
if let Some(first) = self.visited.get(&addr) {
let first = first.clone();
self.catalog.group_aliases.insert(full_name, first);
return Ok(());
}
self.visited.insert(addr, full_name.clone());
self.group(Some(&header), &full_name, depth + 1, stab)
}
}
fn join_path(prefix: &str, name: &str) -> String {
if prefix.is_empty() {
name.to_string()
} else {
format!("{}/{}", prefix, name)
}
}
fn header_declares_link_storage(header: &ObjectHeader) -> bool {
header
.messages
.iter()
.any(|m| m.msg_type == MSG_LINK || m.msg_type == MSG_LINK_INFO)
}
fn describe_link_storage(
header: Option<&ObjectHeader>,
ctx: &FormatContext,
stab: Option<(u64, u64)>,
) -> (LinkStorage, CreationOrder) {
if let Some(h) = header.filter(|h| header_declares_link_storage(h)) {
return h
.messages
.iter()
.find(|m| m.msg_type == MSG_LINK_INFO)
.and_then(|m| LinkInfoMessage::decode(&m.data, ctx).ok())
.map(|(info, _)| {
let storage = if info.is_dense() {
LinkStorage::Dense
} else {
LinkStorage::Compact
};
(storage, info.creation_order())
})
.unwrap_or((LinkStorage::Compact, CreationOrder::Untracked));
}
let (btree_addr, heap_addr) = match stab {
Some(pair) if pair.0 != UNDEF_ADDR && pair.1 != UNDEF_ADDR => pair,
_ => header.map_or((UNDEF_ADDR, UNDEF_ADDR), |h| {
Hdf5Reader::stab_from_header(h, ctx)
}),
};
let storage = if btree_addr != UNDEF_ADDR && heap_addr != UNDEF_ADDR {
LinkStorage::SymbolTable
} else {
LinkStorage::Compact
};
(storage, CreationOrder::Untracked)
}
enum ObjectKind {
Dataset(Box<DatasetReadInfo>),
UnreadableDataset(String),
Group,
CommittedDatatype(Box<CommittedDatatypeInfo>),
}
pub(crate) fn header_is_committed_datatype(header: &ObjectHeader) -> bool {
let present = |t: u8| header.messages.iter().any(|m| m.msg_type == t);
let is_group = present(MSG_LINK)
|| present(MSG_LINK_INFO)
|| present(MSG_SYMBOL_TABLE)
|| present(MSG_GROUP_INFO);
!is_group && present(MSG_DATATYPE) && !(present(MSG_DATASPACE) && present(MSG_DATA_LAYOUT))
}
#[derive(Debug, Clone)]
pub struct CommittedDatatypeInfo {
datatype: Result<DatatypeMessage, String>,
attributes: Vec<AttributeMessage>,
}
impl CommittedDatatypeInfo {
pub fn datatype(&self) -> Result<&DatatypeMessage, &str> {
self.datatype.as_ref().map_err(String::as_str)
}
pub fn attributes(&self) -> &[AttributeMessage] {
&self.attributes
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct SuperblockExtension {
pub shared_message_table: Option<SharedMessageTableMessage>,
pub btree_k: Option<BtreeKMessage>,
pub driver_info: Option<DriverInfoMessage>,
pub file_space_info: Option<FileSpaceInfoMessage>,
}
struct DecodedChunkIndex {
entries: Vec<(u64, u64, u32)>,
coords: Vec<u64>,
images: ChunkImageCache,
}
impl DecodedChunkIndex {
fn new(entries: Vec<(u64, u64, u32)>, coords: Vec<u64>) -> Self {
Self {
entries,
coords,
images: ChunkImageCache::default(),
}
}
}
#[derive(Clone, Copy, PartialEq, Eq, Hash)]
struct ChunkImageKey {
addr: u64,
len: usize,
mask: u32,
}
const CHUNK_CACHE_BYTES: usize = 1024 * 1024;
const CHUNK_CACHE_SLOTS: usize = 521;
#[derive(Default)]
struct ChunkImageCache {
state: std::sync::Mutex<ChunkImageCacheState>,
}
#[derive(Default)]
struct ChunkImageCacheState {
images: std::collections::HashMap<ChunkImageKey, (u64, std::sync::Arc<Vec<u8>>)>,
bytes: usize,
clock: u64,
}
impl ChunkImageCache {
fn get(&self, key: &ChunkImageKey) -> Option<std::sync::Arc<Vec<u8>>> {
let mut state = self.state.lock().ok()?;
state.clock += 1;
let clock = state.clock;
let (stamp, image) = state.images.get_mut(key)?;
*stamp = clock;
Some(std::sync::Arc::clone(image))
}
#[cfg(test)]
fn held(&self) -> usize {
self.state.lock().map(|s| s.images.len()).unwrap_or(0)
}
fn keep(&self, key: ChunkImageKey, image: Vec<u8>) -> std::sync::Arc<Vec<u8>> {
let image = std::sync::Arc::new(image);
if let Ok(mut state) = self.state.lock() {
state.insert(key, std::sync::Arc::clone(&image));
}
image
}
}
impl ChunkImageCacheState {
fn insert(&mut self, key: ChunkImageKey, image: std::sync::Arc<Vec<u8>>) {
let len = image.len();
if len > CHUNK_CACHE_BYTES {
return;
}
self.clock += 1;
if let Some((_, old)) = self.images.insert(key, (self.clock, image)) {
self.bytes -= old.len();
}
self.bytes += len;
while self.bytes > CHUNK_CACHE_BYTES || self.images.len() > CHUNK_CACHE_SLOTS {
let Some(victim) = self
.images
.iter()
.min_by_key(|(_, (stamp, _))| *stamp)
.map(|(k, _)| *k)
else {
break;
};
if let Some((_, old)) = self.images.remove(&victim) {
self.bytes -= old.len();
}
}
}
}
mod dataset_table {
use super::{DatasetReadInfo, DecodedChunkIndex};
use std::sync::Arc;
pub(super) struct DatasetTable {
list: Vec<DatasetReadInfo>,
by_name: std::collections::HashMap<String, usize>,
chunk_index: Vec<Option<(u64, Arc<DecodedChunkIndex>)>>,
}
impl DatasetTable {
pub(super) fn new(list: Vec<DatasetReadInfo>) -> Self {
let mut by_name = std::collections::HashMap::with_capacity(list.len());
for (i, ds) in list.iter().enumerate() {
by_name.entry(ds.name.clone()).or_insert(i);
}
let chunk_index = list.iter().map(|_| None).collect();
Self {
list,
by_name,
chunk_index,
}
}
pub(super) fn position(&self, name: &str) -> Option<usize> {
self.by_name.get(name).copied()
}
pub(super) fn get(&self, name: &str) -> Option<&DatasetReadInfo> {
self.list.get(self.position(name)?)
}
pub(super) fn iter(&self) -> std::slice::Iter<'_, DatasetReadInfo> {
self.list.iter()
}
pub(super) fn entry_mut(&mut self, i: usize) -> &mut DatasetReadInfo {
self.chunk_index[i] = None;
&mut self.list[i]
}
pub(super) fn chunk_index(
&self,
i: usize,
index_address: u64,
) -> Option<&Arc<DecodedChunkIndex>> {
match self.chunk_index.get(i)? {
Some((addr, index)) if *addr == index_address => Some(index),
_ => None,
}
}
pub(super) fn cache_chunk_index(
&mut self,
i: usize,
index_address: u64,
index: DecodedChunkIndex,
) -> Arc<DecodedChunkIndex> {
let index = Arc::new(index);
self.chunk_index[i] = Some((index_address, Arc::clone(&index)));
index
}
}
impl std::ops::Index<usize> for DatasetTable {
type Output = DatasetReadInfo;
fn index(&self, i: usize) -> &DatasetReadInfo {
&self.list[i]
}
}
}
use dataset_table::DatasetTable;
pub struct Hdf5Reader {
handle: FileHandle,
meta: FileMeta,
ext: SuperblockExtension,
_eof: u64,
superblock_version: u8,
datasets: DatasetTable,
unreadable: std::collections::BTreeMap<String, String>,
root_attributes: ObjectAttributes,
root_link_storage: (LinkStorage, CreationOrder),
group_attributes: std::collections::HashMap<String, ObjectAttributes>,
group_link_storage: std::collections::HashMap<String, (LinkStorage, CreationOrder)>,
group_paths: std::collections::BTreeSet<String>,
group_aliases: std::collections::HashMap<String, String>,
links: std::collections::BTreeMap<String, LinkClass>,
path: PathBuf,
source_dir: PathBuf,
locking: crate::io::locking::FileLocking,
elink_prefix: Option<String>,
external: std::collections::BTreeMap<PathBuf, CrossFileEntry>,
vds_resolved: std::collections::BTreeMap<(String, String), PathBuf>,
external_resolved: std::collections::BTreeMap<String, PathBuf>,
datatypes: std::collections::BTreeMap<String, CommittedDatatypeInfo>,
object_paths: std::collections::HashMap<u64, String>,
dataset_access: std::collections::BTreeMap<String, AccessInForce>,
}
pub(crate) type DatasetOpenToken = std::sync::Arc<()>;
struct CrossFileEntry {
reader: Box<Hdf5Reader>,
owner: CrossFileOwner,
}
enum CrossFileOwner {
Reader,
VirtualOpens(std::collections::BTreeSet<String>),
}
impl CrossFileOwner {
fn widen(&mut self, also: CrossFileOwner) {
match (&mut *self, also) {
(CrossFileOwner::Reader, _) => {}
(slot, CrossFileOwner::Reader) => *slot = CrossFileOwner::Reader,
(CrossFileOwner::VirtualOpens(have), CrossFileOwner::VirtualOpens(more)) => {
have.extend(more)
}
}
}
fn virtual_open(vds: &str) -> Self {
CrossFileOwner::VirtualOpens(std::iter::once(vds.to_string()).collect())
}
}
struct AccessInForce {
access: DatasetAccess,
open: std::sync::Weak<()>,
}
fn absolute_path(path: &str) -> String {
format!("/{}", path.trim_start_matches('/'))
}
fn saturating_byte_len(dims: &[u64], element_size: u64) -> u64 {
dims.iter()
.fold(1u64, |acc, &d| acc.saturating_mul(d))
.saturating_mul(element_size)
}
pub(crate) fn fixed_string_attr_value(attr: &AttributeMessage) -> IoResult<String> {
use crate::format::messages::datatype::{fixed_string_content, DatatypeMessage};
let padding = match attr.datatype {
DatatypeMessage::FixedString { padding, .. } => padding,
_ => 0,
};
let content = fixed_string_content(&attr.data, padding).ok_or_else(|| {
crate::io::IoError::InvalidState(format!(
"attribute {:?} uses string padding rule {padding}, which the format reserves",
attr.name
))
})?;
Ok(String::from_utf8_lossy(content).to_string())
}
fn alloc_tiled_fill(total: usize, fill_value: Option<&[u8]>) -> IoResult<Vec<u8>> {
try_tiled_fill(total, fill_value).map_err(|_| {
crate::io::IoError::InvalidState(format!(
"cannot allocate {total} bytes for dataset buffer (file may be corrupt)"
))
})
}
pub(crate) fn read_image_into_new<T, E, F>(count: usize, define: F) -> Result<Vec<T>, E>
where
T: crate::types::H5Type,
F: FnOnce(&mut [u8]) -> Result<(), E>,
E: From<crate::io::IoError>,
{
let too_big = || {
E::from(crate::io::IoError::InvalidState(format!(
"cannot allocate {count} elements of {} bytes for a dataset buffer \
(file may be corrupt)",
std::mem::size_of::<T>()
)))
};
let bytes = count
.checked_mul(std::mem::size_of::<T>())
.ok_or_else(too_big)?;
let mut out: Vec<T> = Vec::new();
out.try_reserve_exact(count).map_err(|_| too_big())?;
let image = unsafe { std::slice::from_raw_parts_mut(out.as_mut_ptr().cast::<u8>(), bytes) };
define(image)?;
unsafe { out.set_len(count) };
Ok(out)
}
fn fill_tiled_into(out: &mut [u8], fill_value: Option<&[u8]>) {
out.fill(0);
if let Some(fv) = fill_value {
if !fv.is_empty() && !out.is_empty() {
for slot in out.chunks_mut(fv.len()) {
let n = slot.len().min(fv.len());
slot[..n].copy_from_slice(&fv[..n]);
}
}
}
}
fn resolve_file_prefix(env_var: &str, prop: Option<&str>, source_dir: &Path) -> Option<PathBuf> {
let env = std::env::var(env_var).ok().filter(|v| !v.is_empty());
let prefix = match env.as_deref() {
Some(v) => v,
None => prop.filter(|v| !v.is_empty())?,
};
if prefix.is_empty() || prefix == "." {
return None;
}
Some(match prefix.strip_prefix("${ORIGIN}") {
Some(rest) => {
let rest = rest.trim_start_matches(['/', '\\']);
if rest.is_empty() {
source_dir.to_path_buf()
} else {
source_dir.join(rest)
}
}
None => PathBuf::from(prefix),
})
}
pub(crate) fn resolve_extfile_prefix(prop: Option<&str>, source_dir: &Path) -> Option<PathBuf> {
resolve_file_prefix("HDF5_EXTFILE_PREFIX", prop, source_dir)
}
fn resolve_vdsfile_prefix(prop: Option<&str>, source_dir: &Path) -> Option<PathBuf> {
resolve_file_prefix("HDF5_VDS_PREFIX", prop, source_dir)
}
pub(crate) fn combine_prefixed_path(prefix: Option<&Path>, name: &str) -> PathBuf {
match prefix {
Some(p) => p.join(name),
None => PathBuf::from(name),
}
}
fn read_external_file_bytes(
external_files: &[ExternalFileSegment],
extfile_prefix: Option<&Path>,
mut skip: u64,
out: &mut [u8],
) -> IoResult<()> {
let mut slot_idx = 0usize;
while slot_idx < external_files.len() && skip >= external_files[slot_idx].size {
skip -= external_files[slot_idx].size;
slot_idx += 1;
}
let mut written = 0usize;
while written < out.len() {
let Some(slot) = external_files.get(slot_idx) else {
return Err(crate::io::IoError::InvalidState(
"read past the logical end of the external file list".into(),
));
};
let full_path = combine_prefixed_path(extfile_prefix, &slot.name);
let ext_handle = FileHandle::open_read_with_locking(&full_path, FileLocking::Disabled)
.map_err(|e| {
crate::io::IoError::InvalidState(format!(
"unable to open external raw data file {}: {e}",
full_path.display()
))
})?;
let avail_in_slot = slot.size.saturating_sub(skip);
let want = (out.len() - written) as u64;
let this_read = avail_in_slot.min(want) as usize;
let dst = &mut out[written..written + this_read];
let at = slot.offset.checked_add(skip).ok_or_else(|| {
crate::io::IoError::InvalidState(format!(
"external file '{}' slot offset {} overflows {skip} bytes into the slot",
slot.name, slot.offset
))
})?;
let got = ext_handle.read_at_most(at, this_read)?;
dst[..got.len()].copy_from_slice(&got);
dst[got.len()..].fill(0);
written += this_read;
skip = 0;
slot_idx += 1;
}
Ok(())
}
const MAX_VIRTUAL_DEPTH: usize = 16;
fn concrete_virtual_mappings(
list: &VirtualMappingList,
resolution: &[MappingResolution],
) -> IoResult<Vec<VirtualMapping>> {
let mut out = Vec::with_capacity(list.mappings.len());
for (i, m) in list.mappings.iter().enumerate() {
match resolution.get(i) {
Some(MappingResolution::Unlimited {
virtual_clip,
source_clip,
}) => out.push(VirtualMapping {
virtual_selection: m.virtual_selection.clip_unlimited(*virtual_clip)?,
source_selection: m.source_selection.clip_unlimited(*source_clip)?,
..built_names(m, 0)?
}),
Some(MappingResolution::Printf { present, .. }) => {
let Some(r) = regular_hyperslab(&m.virtual_selection) else {
continue;
};
let rank = r.start.len();
for &j in present {
out.push(VirtualMapping {
virtual_selection: Selection::Hyperslab {
rank,
form: Hyperslab::Regular(r.unlim_block(j)),
},
..built_names(m, j)?
});
}
}
_ => out.push(built_names(m, 0)?),
}
}
Ok(out)
}
fn built_names(m: &VirtualMapping, blockno: u64) -> IoResult<VirtualMapping> {
Ok(VirtualMapping {
source_file_name: parse_source_name(&m.source_file_name)?.build(blockno),
source_dset_name: parse_source_name(&m.source_dset_name)?.build(blockno),
..m.clone()
})
}
struct VirtualResolveDepth;
thread_local! {
static VIRTUAL_RESOLVE_DEPTH: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
}
impl VirtualResolveDepth {
fn enter<F: FnOnce() -> IoResult<()>>(f: F) -> IoResult<()> {
let depth = VIRTUAL_RESOLVE_DEPTH.with(std::cell::Cell::get);
if depth >= MAX_VIRTUAL_DEPTH {
return Ok(());
}
VIRTUAL_RESOLVE_DEPTH.with(|d| d.set(depth + 1));
let out = f();
VIRTUAL_RESOLVE_DEPTH.with(|d| d.set(depth));
out
}
}
fn regular_hyperslab(sel: &Selection) -> Option<&RegularHyperslab> {
match sel {
Selection::Hyperslab {
form: Hyperslab::Regular(r),
..
} => Some(r),
_ => None,
}
}
fn copy_matched_selections(
mut read_source_box: impl FnMut(&[u64], &[u64], &mut [u8]) -> IoResult<()>,
source: &ResolvedSelection,
target: &ResolvedSelection,
element_size: u64,
out: &mut [u8],
) -> IoResult<()> {
let (n_source, n_target) = (source.n_elements(), target.n_elements());
if n_source != n_target {
return Err(crate::io::IoError::InvalidState(format!(
"virtual dataset mapping's source selection holds {n_source} elements and its virtual selection {n_target}, which H5D_virtual_check_mapping_pre refuses"
)));
}
let mut per_box: Vec<Vec<(u64, u64, u64)>> = vec![Vec::new(); source.boxes.len()];
let (mut si, mut ti) = (0usize, 0usize);
let (mut s_done, mut t_done) = (0u64, 0u64);
while si < source.runs.len() && ti < target.runs.len() {
let (s, t) = (source.runs[si], target.runs[ti]);
let len = (s.len - s_done).min(t.len - t_done);
per_box[s.box_index].push((s.offset_in_box + s_done, t.offset_in_extent + t_done, len));
s_done += len;
t_done += len;
if s_done == s.len {
si += 1;
s_done = 0;
}
if t_done == t.len {
ti += 1;
t_done = 0;
}
}
for (segments, (box_start, box_count)) in per_box.iter().zip(&source.boxes) {
if segments.is_empty() {
continue;
}
let nbytes = saturating_byte_len(box_count, element_size) as usize;
let mut buf = alloc_tiled_fill(nbytes, None)?;
read_source_box(box_start, box_count, &mut buf)?;
for &(from, to, len) in segments {
let (from, to, len) = (
(from * element_size) as usize,
(to * element_size) as usize,
(len * element_size) as usize,
);
let src = buf.get(from..from + len).ok_or_else(|| {
crate::io::IoError::InvalidState(
"virtual dataset mapping's source selection reaches past its source box".into(),
)
})?;
let dst = out.get_mut(to..to + len).ok_or_else(|| {
crate::io::IoError::InvalidState(
"virtual dataset mapping's virtual selection reaches past the dataset".into(),
)
})?;
dst.copy_from_slice(src);
}
}
Ok(())
}
fn read_heap_collection_from(
handle: &mut FileHandle,
ctx: &FormatContext,
addr: u64,
) -> IoResult<GlobalHeapCollection> {
let ss = ctx.sizeof_size as usize;
let header_len = 4 + 1 + 3 + ss;
let header_buf = handle.read_at_most(addr, header_len)?;
if header_buf.len() < header_len || header_buf[0..4] != *b"GCOL" {
return Err(crate::io::IoError::InvalidState(format!(
"bad global heap collection signature at address {addr:#x}"
)));
}
let declared = read_uint(&header_buf[8..], ss) as usize;
if declared < 4096 {
return Err(crate::io::IoError::InvalidState(format!(
"global heap collection at address {addr:#x} declares size {declared}, \
below the 4096-byte minimum"
)));
}
let heap_buf = handle.read_at(addr, declared)?;
let (coll, _) = GlobalHeapCollection::decode(&heap_buf, ctx)?;
Ok(coll)
}
struct ChunkReadJob {
addr: u64,
len: usize,
at_most: bool,
mask: u32,
}
fn read_chunk_raw(handle: &FileHandle, j: &ChunkReadJob) -> IoResult<Vec<u8>> {
if j.at_most {
Ok(handle.read_at_most(j.addr, j.len)?)
} else {
Ok(handle.read_at(j.addr, j.len)?)
}
}
const PREAD_COST_BYTES: u64 = 16 * 1024;
struct ChunkOverlap {
lo: Vec<u64>,
hi: Vec<u64>,
}
impl ChunkOverlap {
fn of(place: &ChunkPlacement, chunk_coords: &[u64]) -> Option<Self> {
let ChunkPlacement {
geo: ChunkOutputGeometry {
dims, chunk_dims, ..
},
starts,
counts,
} = *place;
let ndims = dims.len();
if ndims == 0 {
return None;
}
let mut lo = vec![0u64; ndims];
let mut hi = vec![0u64; ndims];
for d in 0..ndims {
let origin = chunk_coords[d].saturating_mul(chunk_dims[d]);
let chunk_end = origin.saturating_add(chunk_dims[d]).min(dims[d]);
let sel_end = starts[d].saturating_add(counts[d]);
lo[d] = origin.max(starts[d]);
hi[d] = chunk_end.min(sel_end);
if lo[d] >= hi[d] {
return None;
}
}
Some(Self { lo, hi })
}
fn bytes(&self, element_size: u64) -> u64 {
self.lo
.iter()
.zip(&self.hi)
.map(|(&l, &h)| h - l)
.product::<u64>()
.saturating_mul(element_size)
}
}
fn planned_coverage(place: &ChunkPlacement, jobs: &[Option<ChunkReadJob>], coords: &[u64]) -> u64 {
let rank = place.geo.dims.len();
if rank == 0 {
return 0;
}
let mut seen: std::collections::HashSet<&[u64]> = std::collections::HashSet::new();
let mut covered = 0u64;
for (i, job) in jobs.iter().enumerate() {
if job.is_none() {
continue;
}
let c = &coords[i * rank..(i + 1) * rank];
if !seen.insert(c) {
continue;
}
if let Some(overlap) = ChunkOverlap::of(place, c) {
covered = covered.saturating_add(overlap.bytes(place.geo.element_size));
}
}
covered
}
fn for_each_chunk_run(
place: &ChunkPlacement,
chunk_coords: &[u64],
mut f: impl FnMut(u64, u64, usize),
) {
let ChunkPlacement {
geo:
ChunkOutputGeometry {
dims,
chunk_dims,
element_size,
},
starts,
counts,
} = *place;
let ndims = dims.len();
let Some(ChunkOverlap { lo, hi }) = ChunkOverlap::of(place, chunk_coords) else {
return;
};
let chunk_strides = compute_strides(chunk_dims, element_size);
let out_strides = compute_strides(counts, element_size);
let last = ndims - 1;
let run_bytes = ((hi[last] - lo[last]) * element_size) as usize;
let outer_extent: Vec<u64> = (0..last).map(|d| hi[d] - lo[d]).collect();
let n_outer: u64 = outer_extent.iter().product(); let mut oc = vec![0u64; last];
for _ in 0..n_outer {
let mut src_off = 0u64;
let mut dst_off = 0u64;
for d in 0..ndims {
let g = if d < last { lo[d] + oc[d] } else { lo[last] };
let origin = chunk_coords[d].saturating_mul(chunk_dims[d]);
src_off += (g - origin) * chunk_strides[d];
dst_off += (g - starts[d]) * out_strides[d];
}
f(src_off, dst_off, run_bytes);
for d in (0..last).rev() {
oc[d] += 1;
if oc[d] < outer_extent[d] {
break;
}
oc[d] = 0;
}
}
}
fn copy_chunk_runs(
chunk_data: &[u8],
output: &mut [u8],
place: &ChunkPlacement,
chunk_coords: &[u64],
skipped: &mut Vec<(usize, usize)>,
) {
for_each_chunk_run(place, chunk_coords, |src, dst, len| {
let (s, d) = (src as usize, dst as usize);
if s + len <= chunk_data.len() && d + len <= output.len() {
output[d..d + len].copy_from_slice(&chunk_data[s..s + len]);
} else {
push_skipped(skipped, output.len(), d, len);
}
});
}
fn push_skipped(skipped: &mut Vec<(usize, usize)>, out_len: usize, dst: usize, len: usize) {
let start = dst.min(out_len);
let end = dst.saturating_add(len).min(out_len);
if start < end {
skipped.push((start, end - start));
}
}
fn read_chunk_runs_into(
handle: &FileHandle,
job: &ChunkReadJob,
place: &ChunkPlacement,
chunk_coords: &[u64],
output: &mut [u8],
skipped: &mut Vec<(usize, usize)>,
dst: ReadDst,
) -> bool {
let (addr, image_len) = (job.addr, job.len);
let mut runs: Vec<(u64, usize, usize)> = Vec::new();
let mut selected = 0u64;
let out_len = output.len();
for_each_chunk_run(place, chunk_coords, |src, out_off, len| {
let (s, d) = (src as usize, out_off as usize);
if s + len > image_len || d + len > out_len {
push_skipped(skipped, out_len, d, len);
return;
}
selected += len as u64;
let at = addr.saturating_add(src);
if let Some(last) = runs.last_mut() {
if last.0.checked_add(last.2 as u64) == Some(at) && last.1 + last.2 == d {
last.2 += len;
return;
}
}
runs.push((at, d, len));
});
if runs.is_empty() {
return false;
}
if (runs.len() as u64 - 1).saturating_mul(PREAD_COST_BYTES) > image_len as u64 + selected {
return false;
}
for (offset, at, len) in runs {
if handle
.read_exact_at_into(offset, &mut output[at..at + len], dst)
.is_err()
{
return false;
}
}
true
}
fn place_chunk_jobs(
handle: &FileHandle,
mut jobs: Vec<Option<ChunkReadJob>>,
coords: &[u64],
req: ChunkReadRequest,
geo: &ChunkOutputGeometry,
cache: Option<&ChunkImageCache>,
output: &mut [u8],
) -> IoResult<()> {
let ChunkReadRequest {
pipeline,
target,
fill_value,
dst,
} = req;
let rank = geo.dims.len();
let zeros = vec![0u64; rank];
let place = ChunkPlacement::resolve(geo, target, &zeros);
let at = |i: usize| &coords[i * rank..(i + 1) * rank];
let prefilled = planned_coverage(&place, &jobs, coords) != output.len() as u64;
if prefilled {
fill_tiled_into(output, fill_value);
}
let mut skipped: Vec<(usize, usize)> = Vec::new();
if pipeline.is_none() {
for (i, job) in jobs.iter_mut().enumerate() {
let Some(j) = job.as_ref() else { continue };
let mark = skipped.len();
if read_chunk_runs_into(handle, j, &place, at(i), output, &mut skipped, dst) {
*job = None;
} else {
skipped.truncate(mark);
}
}
}
let mut hits: Vec<Option<std::sync::Arc<Vec<u8>>>> = Vec::new();
let mut keys: Vec<Option<ChunkImageKey>> = Vec::new();
let cache = cache.filter(|_| pipeline.is_some());
if let Some(cache) = cache {
hits.resize_with(jobs.len(), || None);
keys.resize(jobs.len(), None);
for (i, job) in jobs.iter_mut().enumerate() {
let Some(j) = job.as_ref() else { continue };
if !place.leaves_chunk_unconsumed(at(i)) {
continue;
}
let key = ChunkImageKey {
addr: j.addr,
len: j.len,
mask: j.mask,
};
match cache.get(&key) {
Some(image) => {
hits[i] = Some(image);
*job = None;
}
None => keys[i] = Some(key),
}
}
for (i, image) in hits.iter().enumerate() {
if let Some(image) = image {
copy_chunk_runs(image, output, &place, at(i), &mut skipped);
}
}
}
if jobs.iter().any(Option::is_some) {
let decoded = {
let sinks = carve_sinks(output, &jobs, coords, &place);
read_and_decompress_chunks(handle, pipeline, jobs, sinks, geo.image_bytes())?
};
for (i, chunk) in decoded.into_iter().enumerate() {
match chunk {
ChunkDecoded::Absent => {}
ChunkDecoded::Image(data) => {
match keys.get_mut(i).and_then(Option::take).zip(cache) {
Some((key, cache)) => {
let image = cache.keep(key, data);
copy_chunk_runs(&image, output, &place, at(i), &mut skipped)
}
None => copy_chunk_runs(&data, output, &place, at(i), &mut skipped),
}
}
ChunkDecoded::InPlace { dst, len, bytes } if bytes < len => {
fill_tiled_into(&mut output[dst..dst + len], fill_value)
}
ChunkDecoded::InPlace { .. } => {}
}
}
}
if !prefilled {
for (dst, len) in skipped {
fill_tiled_into(&mut output[dst..dst + len], fill_value);
}
}
Ok(())
}
enum ChunkSink<'a> {
Direct { dst: usize, out: &'a mut [u8] },
Staged,
}
enum ChunkDecoded {
Absent,
Image(Vec<u8>),
InPlace {
dst: usize,
len: usize,
bytes: usize,
},
}
fn carve_sinks<'a>(
output: &'a mut [u8],
jobs: &[Option<ChunkReadJob>],
coords: &[u64],
place: &ChunkPlacement,
) -> Vec<ChunkSink<'a>> {
let mut sinks = Vec::with_capacity(jobs.len());
sinks.resize_with(jobs.len(), || ChunkSink::Staged);
let rank = place.geo.dims.len();
let out_len = output.len();
if rank == 0 {
return sinks;
}
let Some(image_bytes) = place.geo.image_bytes() else {
return sinks;
};
if place
.counts
.iter()
.zip(place.geo.chunk_dims)
.any(|(c, k)| c < k)
{
return sinks;
}
let mut wanted: Vec<(usize, usize)> = Vec::new();
for (i, job) in jobs.iter().enumerate() {
if job.is_none() {
continue;
}
let mut runs = 0usize;
let mut first = (0u64, 0u64, 0usize);
for_each_chunk_run(place, &coords[i * rank..(i + 1) * rank], |src, dst, len| {
if runs == 0 {
first = (src, dst, len);
}
runs += 1;
});
let (src, dst, len) = first;
if runs == 1
&& src == 0
&& len as u64 == image_bytes
&& dst.saturating_add(len as u64) <= out_len as u64
{
wanted.push((dst as usize, i));
}
}
wanted.sort_unstable();
let len = image_bytes as usize;
let mut rest: &'a mut [u8] = output;
let mut base = 0usize;
for (dst, i) in wanted {
if dst < base {
continue;
}
let (_, tail) = std::mem::take(&mut rest).split_at_mut(dst - base);
let (mine, tail) = tail.split_at_mut(len);
sinks[i] = ChunkSink::Direct { dst, out: mine };
rest = tail;
base = dst + len;
}
sinks
}
fn decompress_chunk_into(
pipeline: Option<&FilterPipeline>,
raw: &[u8],
mask: u32,
out: &mut [u8],
) -> IoResult<usize> {
match pipeline {
Some(pl) => Ok(filter::reverse_filters_masked_into(pl, raw, mask, out)?),
None => {
let n = raw.len().min(out.len());
out[..n].copy_from_slice(&raw[..n]);
Ok(raw.len())
}
}
}
fn decompress_chunk(
pipeline: Option<&FilterPipeline>,
raw: Vec<u8>,
mask: u32,
image_bytes: Option<u64>,
) -> IoResult<Vec<u8>> {
let Some(pl) = pipeline else { return Ok(raw) };
let Some(image_bytes) = image_bytes.and_then(|b| usize::try_from(b).ok()) else {
return Ok(filter::reverse_filters_masked(pl, &raw, mask)?);
};
let mut image = vec![0u8; image_bytes];
let produced = filter::reverse_filters_masked_into(pl, &raw, mask, &mut image)?;
image.truncate(produced.min(image_bytes));
Ok(image)
}
#[cfg(feature = "parallel")]
const PARALLEL_MIN_JOB_BYTES: u64 = 256 * 1024;
#[cfg(feature = "parallel")]
fn worth_parallel(jobs: &[Option<ChunkReadJob>]) -> bool {
let mut n = 0usize;
let mut bytes = 0u64;
for j in jobs.iter().flatten() {
n += 1;
bytes = bytes.saturating_add(j.len as u64);
if n > 1 && bytes >= PARALLEL_MIN_JOB_BYTES {
return true;
}
}
false
}
fn read_and_decompress_chunks(
handle: &FileHandle,
pipeline: Option<&FilterPipeline>,
jobs: Vec<Option<ChunkReadJob>>,
sinks: Vec<ChunkSink<'_>>,
image_bytes: Option<u64>,
) -> IoResult<Vec<ChunkDecoded>> {
let deliver = |raw: Vec<u8>, mask: u32, sink: ChunkSink<'_>| -> IoResult<ChunkDecoded> {
match sink {
ChunkSink::Direct { dst, out } => {
let len = out.len();
let bytes = decompress_chunk_into(pipeline, &raw, mask, out)?;
Ok(ChunkDecoded::InPlace { dst, len, bytes })
}
ChunkSink::Staged => Ok(ChunkDecoded::Image(decompress_chunk(
pipeline,
raw,
mask,
image_bytes,
)?)),
}
};
#[cfg(all(feature = "parallel", any(unix, windows)))]
{
use rayon::prelude::*;
let decode =
|(job, sink): (Option<ChunkReadJob>, ChunkSink<'_>)| -> IoResult<ChunkDecoded> {
match job {
Some(j) => deliver(read_chunk_raw(handle, &j)?, j.mask, sink),
None => Ok(ChunkDecoded::Absent),
}
};
let pool = crate::parallel::io_pool().filter(|_| worth_parallel(&jobs));
let work: Vec<_> = jobs.into_iter().zip(sinks).collect();
match pool {
Some(pool) => pool.install(|| {
work.into_par_iter()
.map(&decode)
.collect::<IoResult<Vec<_>>>()
}),
None => work.into_iter().map(decode).collect::<IoResult<Vec<_>>>(),
}
}
#[cfg(all(feature = "parallel", not(any(unix, windows))))]
{
use rayon::prelude::*;
let parallel = worth_parallel(&jobs);
let raws: Vec<Option<(Vec<u8>, u32)>> = jobs
.into_iter()
.map(|job| match job {
Some(j) => Ok(Some((read_chunk_raw(handle, &j)?, j.mask))),
None => Ok(None),
})
.collect::<IoResult<Vec<_>>>()?;
let decode =
|(r, sink): (Option<(Vec<u8>, u32)>, ChunkSink<'_>)| -> IoResult<ChunkDecoded> {
match r {
Some((raw, mask)) => deliver(raw, mask, sink),
None => Ok(ChunkDecoded::Absent),
}
};
let work: Vec<_> = raws.into_iter().zip(sinks).collect();
match crate::parallel::io_pool().filter(|_| parallel) {
Some(pool) => pool.install(|| {
work.into_par_iter()
.map(&decode)
.collect::<IoResult<Vec<_>>>()
}),
None => work.into_iter().map(decode).collect::<IoResult<Vec<_>>>(),
}
}
#[cfg(not(feature = "parallel"))]
{
jobs.into_iter()
.zip(sinks)
.map(|(job, sink)| match job {
Some(j) => deliver(read_chunk_raw(handle, &j)?, j.mask, sink),
None => Ok(ChunkDecoded::Absent),
})
.collect()
}
}
impl Hdf5Reader {
pub fn open_swmr(path: &Path) -> IoResult<Self> {
Self::open(path)
}
pub fn open_swmr_with_locking(
path: &Path,
locking: crate::io::locking::FileLocking,
) -> IoResult<Self> {
Self::open_with_locking(path, locking)
}
pub fn open(path: &Path) -> IoResult<Self> {
Self::open_with_locking(
path,
crate::io::locking::FileLocking::from_env_or(Default::default()),
)
}
pub fn open_with_locking(
path: &Path,
locking: crate::io::locking::FileLocking,
) -> IoResult<Self> {
let mut handle = FileHandle::open_read_with_locking(path, locking)?;
let super_addr = handle
.locate_signature()?
.ok_or(crate::format::FormatError::InvalidSignature)?;
handle.set_base(super_addr);
let sb_buf = handle.read_at_most(0, 1024)?;
let version = detect_superblock_version(&sb_buf)?;
let origin = Origin {
path: path.to_path_buf(),
locking,
};
let mut reader = match version {
0 | 1 => Self::open_v0v1(handle, &sb_buf, origin)?,
2 | 3 => Self::open_v2v3(handle, &sb_buf, origin)?,
v => {
return Err(crate::io::IoError::Format(
crate::format::FormatError::InvalidVersion(v),
))
}
};
let canonical = std::fs::canonicalize(path)?;
reader.source_dir = canonical
.parent()
.map(Path::to_path_buf)
.unwrap_or_default();
VirtualResolveDepth::enter(|| reader.resolve_virtual_extents())?;
Ok(reader)
}
fn open_v2v3(mut handle: FileHandle, sb_buf: &[u8], origin: Origin) -> IoResult<Self> {
let sb = SuperblockV2V3::decode(sb_buf)?;
let ctx = FormatContext {
sizeof_addr: sb.sizeof_offsets,
sizeof_size: sb.sizeof_lengths,
};
let (meta, ext) = Self::read_extension_and_meta(
&mut handle,
ctx,
BTreeV1Config::default(),
sb.superblock_extension_address,
)?;
let root_header =
Self::read_object_header_full(&mut handle, &meta, sb.root_group_object_header_address)?;
let catalog = Self::build_catalog(
&mut handle,
&meta,
Some(&root_header),
sb.root_group_object_header_address,
None,
)?;
let root_attributes = collect_object_attributes(&mut handle, &ctx, &root_header);
let root_link_storage = describe_link_storage(Some(&root_header), &ctx, None);
Ok(Self {
handle,
meta,
ext,
_eof: sb.end_of_file_address,
superblock_version: sb.version,
object_paths: catalog.object_paths(sb.root_group_object_header_address),
datasets: DatasetTable::new(catalog.datasets),
unreadable: catalog.unreadable,
root_attributes,
root_link_storage,
group_attributes: catalog.group_attributes,
group_link_storage: catalog.group_link_storage,
group_paths: catalog.group_paths,
group_aliases: catalog.group_aliases,
links: catalog.links,
datatypes: catalog.datatypes,
path: origin.path,
locking: origin.locking,
elink_prefix: None,
external: Default::default(),
external_resolved: Default::default(),
dataset_access: Default::default(),
vds_resolved: Default::default(),
source_dir: PathBuf::new(),
})
}
fn open_v0v1(mut handle: FileHandle, sb_buf: &[u8], origin: Origin) -> IoResult<Self> {
let sb = SuperblockV0V1::decode(sb_buf)?;
let ctx = FormatContext {
sizeof_addr: sb.sizeof_offsets,
sizeof_size: sb.sizeof_lengths,
};
let sb_btree = BTreeV1Config {
sym_leaf_k: sb.sym_leaf_k,
snode_internal_k: sb.btree_internal_k,
chunk_internal_k: sb
.indexed_storage_k
.unwrap_or(BTreeV1Config::default().chunk_internal_k),
};
let (meta, ext) = Self::read_extension_and_meta(
&mut handle,
ctx,
sb_btree,
sb.superblock_extension_address,
)?;
let ste = &sb.root_symbol_table_entry;
let root_obj_addr = ste.obj_header_addr;
let ste_stab = ste.cached_symbol_table();
let root_hdr = Self::read_object_header_full(&mut handle, &meta, root_obj_addr).ok();
let root_attributes = match root_hdr {
Some(ref h) => collect_object_attributes(&mut handle, &ctx, h),
None => ObjectAttributes::default(),
};
let root_link_storage = describe_link_storage(root_hdr.as_ref(), &ctx, ste_stab);
let catalog = Self::build_catalog(
&mut handle,
&meta,
root_hdr.as_ref(),
root_obj_addr,
ste_stab,
)?;
Ok(Self {
handle,
meta,
ext,
_eof: sb.end_of_file_address,
superblock_version: sb.version,
object_paths: catalog.object_paths(root_obj_addr),
datasets: DatasetTable::new(catalog.datasets),
unreadable: catalog.unreadable,
root_attributes,
root_link_storage,
group_attributes: catalog.group_attributes,
group_link_storage: catalog.group_link_storage,
group_paths: catalog.group_paths,
group_aliases: catalog.group_aliases,
links: catalog.links,
datatypes: catalog.datatypes,
path: origin.path,
locking: origin.locking,
elink_prefix: None,
external: Default::default(),
external_resolved: Default::default(),
dataset_access: Default::default(),
vds_resolved: Default::default(),
source_dir: PathBuf::new(),
})
}
pub(crate) fn read_extension_and_meta(
handle: &mut FileHandle,
ctx: FormatContext,
sb_btree: BTreeV1Config,
ext_addr: u64,
) -> IoResult<(FileMeta, SuperblockExtension)> {
let mut meta = FileMeta {
ctx,
btree: sb_btree,
sohm: None,
};
let ext = Self::superblock_extension_at(handle, ctx, sb_btree, ext_addr)?;
if let Some(k) = ext.btree_k {
meta.btree = BTreeV1Config {
sym_leaf_k: k.sym_leaf_k,
snode_internal_k: k.snode_internal_k,
chunk_internal_k: k.chunk_internal_k,
};
}
let b = &meta.btree;
if b.sym_leaf_k == 0 || b.snode_internal_k == 0 || b.chunk_internal_k == 0 {
return Err(crate::io::IoError::Format(
crate::format::FormatError::InvalidData(format!(
"v1 B-tree K values must be non-zero (sym_leaf={}, snode={}, chunk={})",
b.sym_leaf_k, b.snode_internal_k, b.chunk_internal_k
)),
));
}
if let Some(smt) = &ext.shared_message_table {
meta.sohm = Some(Self::read_sohm_table(handle, &meta.ctx, smt)?);
}
Ok((meta, ext))
}
pub(crate) fn superblock_extension_at(
handle: &mut FileHandle,
ctx: FormatContext,
btree: BTreeV1Config,
ext_addr: u64,
) -> IoResult<SuperblockExtension> {
if ext_addr == UNDEF_ADDR || ext_addr == 0 {
return Ok(SuperblockExtension::default());
}
let meta = FileMeta {
ctx,
btree,
sohm: None,
};
Self::read_superblock_extension(handle, &meta, ext_addr)
}
fn read_superblock_extension(
handle: &mut FileHandle,
meta: &FileMeta,
addr: u64,
) -> IoResult<SuperblockExtension> {
let header = Self::read_object_header_full(handle, meta, addr)?;
let ctx = &meta.ctx;
let mut ext = SuperblockExtension::default();
for msg in &header.messages {
match msg.msg_type {
MSG_SHARED_MESSAGE_TABLE => {
ext.shared_message_table =
Some(SharedMessageTableMessage::decode(&msg.data, ctx)?);
}
MSG_BTREE_K => ext.btree_k = Some(BtreeKMessage::decode(&msg.data)?),
MSG_DRIVER_INFO => ext.driver_info = Some(DriverInfoMessage::decode(&msg.data)?),
MSG_FILE_SPACE_INFO => {
ext.file_space_info = Some(FileSpaceInfoMessage::decode(&msg.data, ctx)?);
}
_ => {}
}
}
Ok(ext)
}
fn read_sohm_table(
handle: &mut FileHandle,
ctx: &FormatContext,
smt: &SharedMessageTableMessage,
) -> IoResult<SohmMasterTable> {
if smt.table_address == UNDEF_ADDR || smt.nindexes == 0 {
return Ok(SohmMasterTable::default());
}
let size = SohmMasterTable::encoded_size(ctx, smt.nindexes);
let buf = handle.read_at(smt.table_address, size)?;
Ok(SohmMasterTable::decode(&buf, ctx, smt.nindexes)?)
}
fn stab_from_header(header: &ObjectHeader, ctx: &FormatContext) -> (u64, u64) {
for msg in &header.messages {
if msg.msg_type == MSG_SYMBOL_TABLE {
let sa = ctx.sizeof_addr as usize;
if msg.data.len() >= 2 * sa {
return (read_uint(&msg.data, sa), read_uint(&msg.data[sa..], sa));
}
}
}
(UNDEF_ADDR, UNDEF_ADDR)
}
fn build_catalog(
handle: &mut FileHandle,
meta: &FileMeta,
root_header: Option<&ObjectHeader>,
root_addr: u64,
root_stab: Option<(u64, u64)>,
) -> IoResult<Catalog> {
let mut walk = CatalogWalk::new(handle, meta, root_addr);
walk.group(root_header, "", 0, root_stab)?;
Ok(walk.finish())
}
pub(crate) fn read_dense_links(
handle: &mut FileHandle,
ctx: &FormatContext,
fractal_heap_addr: u64,
) -> IoResult<Vec<LinkMessage>> {
let hdr_buf = handle.read_at_most(fractal_heap_addr, 512)?;
let fh_header = FractalHeapHeader::decode(&hdr_buf, ctx)?;
let mut br = HandleBlockReader { handle };
let payloads = fractal_heap::collect_managed_objects(&fh_header, ctx, &mut br)?;
let mut links = Vec::new();
for payload in payloads {
let mut pos = 0;
while pos < payload.len() {
if payload[pos] != 1 {
break;
}
match LinkMessage::decode(&payload[pos..], ctx) {
Ok((link, consumed)) if consumed > 0 => {
links.push(link);
pos += consumed;
}
_ => break,
}
}
}
if links.len() < fh_header.man_nobjs as usize {
return Err(crate::io::IoError::InvalidState(format!(
"dense link storage at address {fractal_heap_addr:#x} holds {} managed \
objects but only {} decoded as links",
fh_header.man_nobjs,
links.len()
)));
}
Ok(links)
}
fn collect_snod_addresses(
handle: &mut FileHandle,
meta: &FileMeta,
tree_addr: u64,
depth: usize,
visited: &mut std::collections::HashSet<u64>,
) -> IoResult<Vec<u64>> {
let sizeof_addr = meta.ctx.sizeof_addr as usize;
let sizeof_size = meta.ctx.sizeof_size as usize;
if depth > 256 || !visited.insert(tree_addr) {
return Ok(Vec::new());
}
let node_size = meta.btree.snode_btree_node_size(sizeof_addr, sizeof_size);
let buf = handle.read_at_most(tree_addr, node_size)?;
let node = BTreeV1Node::decode(
&buf,
sizeof_addr,
sizeof_size,
meta.btree.snode_max_entries(),
)?;
if node.level == 0 {
Ok(node.children.clone())
} else {
let mut addrs = Vec::new();
for &child_addr in &node.children {
let child_addrs =
Self::collect_snod_addresses(handle, meta, child_addr, depth + 1, visited)?;
addrs.extend(child_addrs);
}
Ok(addrs)
}
}
fn read_object_header_full(
handle: &mut FileHandle,
meta: &FileMeta,
addr: u64,
) -> IoResult<ObjectHeader> {
crate::io::object_header_io::read_object_header_full(handle, meta, addr)
}
fn committed_datatype(
handle: &mut FileHandle,
header: &ObjectHeader,
meta: &FileMeta,
) -> CommittedDatatypeInfo {
let datatype = header
.messages
.iter()
.find(|m| m.msg_type == MSG_DATATYPE)
.cloned()
.ok_or_else(|| "it holds no datatype message".to_string())
.and_then(|m| {
crate::io::object_header_io::read_datatype_message(handle, meta, &m).map_err(|e| {
match e {
crate::io::IoError::Unsupported(why) => why,
other => format!("its datatype message does not decode: {other}"),
}
})
});
let attributes = header
.messages
.iter()
.filter(|m| m.msg_type == MSG_ATTRIBUTE && m.flags & MSG_FLAG_SHARED == 0)
.filter_map(|m| {
AttributeMessage::decode(&m.data, &meta.ctx)
.ok()
.map(|(a, _)| a)
})
.collect();
CommittedDatatypeInfo {
datatype,
attributes,
}
}
fn classify_object(
handle: &mut FileHandle,
header: &ObjectHeader,
meta: &FileMeta,
name: &str,
addr: u64,
) -> ObjectKind {
let ctx = &meta.ctx;
let present = |t: u8| header.messages.iter().any(|m| m.msg_type == t);
let is_group = present(MSG_LINK)
|| present(MSG_LINK_INFO)
|| present(MSG_SYMBOL_TABLE)
|| present(MSG_GROUP_INFO);
if is_group {
return ObjectKind::Group;
}
if header_is_committed_datatype(header) {
return ObjectKind::CommittedDatatype(Box::new(Self::committed_datatype(
handle, header, meta,
)));
}
let is_dataset =
present(MSG_DATATYPE) && present(MSG_DATASPACE) && present(MSG_DATA_LAYOUT);
if !is_dataset {
return ObjectKind::Group;
}
let mut datatype = None;
let mut dataspace = None;
let mut layout = None;
let mut filter_pipeline = None;
let mut fill_value = None;
let mut fill_defined: u8 = 1;
let mut fill_write_time: u8 = FILL_TIME_IFSET;
let mut alloc_time: u8 = ALLOC_TIME_LATE;
let mut blocked: Option<String> = None;
let mut block = |why: String| {
if blocked.is_none() {
blocked = Some(why);
}
};
let mut external_file_list = None;
for msg in &header.messages {
let shared = msg.flags & MSG_FLAG_SHARED != 0;
if shared && !matches!(msg.msg_type, MSG_DATATYPE | MSG_ATTRIBUTE) {
block(format!(
"its message of type {:#04x} is a shared-message reference, which this \
crate follows only for datatypes",
msg.msg_type
));
continue;
}
match msg.msg_type {
MSG_DATATYPE => {
match crate::io::object_header_io::read_datatype_message(handle, meta, msg) {
Ok(dt) => datatype = Some(dt),
Err(crate::io::IoError::Unsupported(why)) => block(why),
Err(e) => block(format!("its datatype message does not decode: {e}")),
}
}
MSG_DATASPACE => match DataspaceMessage::decode(&msg.data, ctx) {
Ok((ds, _)) => dataspace = Some(ds),
Err(e) => block(format!("its dataspace message does not decode: {e}")),
},
MSG_DATA_LAYOUT => match DataLayoutMessage::decode(&msg.data, ctx) {
Ok((dl, _)) => layout = Some(dl),
Err(e) => block(format!("its data layout message does not decode: {e}")),
},
MSG_FILTER_PIPELINE => match FilterPipeline::decode(&msg.data) {
Ok((fp, _)) => {
if !fp.filters.is_empty() {
filter_pipeline = Some(fp);
}
}
Err(e) => block(format!("its filter pipeline message does not decode: {e}")),
},
MSG_FILL_VALUE => match FillValueMessage::decode(&msg.data) {
Ok((fv, _)) => {
fill_defined = fv.fill_defined;
fill_write_time = fv.fill_write_time;
alloc_time = fv.alloc_time;
if fv.fill_defined == 2 {
fill_value = fv.fill_value;
}
}
Err(e) => block(format!("its fill value message does not decode: {e}")),
},
MSG_EXTERNAL_FILE_LIST => {
match ExternalFileListMessage::decode(&msg.data, ctx) {
Ok((efl, _)) => external_file_list = Some(efl),
Err(e) => block(format!(
"its external file list message does not decode: {e}"
)),
}
}
_ => {}
}
}
if let (Some(ds), Some(dt), Some(dl)) = (&dataspace, &datatype, &layout) {
if let Err(e) = dl.check_against_dataset(ds, dt, ctx) {
block(format!(
"its layout doesn't fit its dataspace and datatype: {e}"
));
}
}
if let Some(why) = blocked {
return ObjectKind::UnreadableDataset(why);
}
let external_files = match external_file_list {
Some(efl) => match Self::resolve_external_file_slots(handle, ctx, &efl) {
Ok(slots) => slots,
Err(e) => {
return ObjectKind::UnreadableDataset(format!(
"its external file list does not resolve: {e}"
))
}
},
None => Vec::new(),
};
let virtual_mappings = match &layout {
Some(DataLayoutMessage::Virtual {
heap_address,
heap_index,
..
}) if *heap_index != 0 => {
match Self::resolve_virtual_mappings(handle, ctx, *heap_address, *heap_index, name)
{
Ok(list) => Some(list),
Err(e) => {
return ObjectKind::UnreadableDataset(format!(
"its virtual dataset mapping list does not resolve: {e}"
))
}
}
}
_ => None,
};
let attributes = collect_object_attributes(handle, ctx, header);
match (datatype, dataspace, layout) {
(Some(dt), Some(ds), Some(dl)) => ObjectKind::Dataset(Box::new(DatasetReadInfo {
name: name.to_string(),
object_header_address: addr,
datatype: dt,
dataspace: ds,
layout: dl,
filter_pipeline,
attributes,
fill_value,
fill_defined,
fill_write_time,
alloc_time,
external_files,
virtual_mappings,
virtual_resolution: None,
virtual_stored_dims: None,
})),
_ => ObjectKind::UnreadableDataset(
"its datatype, dataspace and data layout messages decoded but did not all \
produce a value"
.into(),
),
}
}
fn resolve_virtual_mappings(
handle: &mut FileHandle,
ctx: &FormatContext,
heap_address: u64,
heap_index: u32,
name: &str,
) -> IoResult<VirtualMappingList> {
let coll = read_heap_collection_from(handle, ctx, heap_address)?;
let idx = u16::try_from(heap_index).map_err(|_| {
crate::io::IoError::InvalidState(format!(
"dataset {name:?} virtual mapping heap index {heap_index} does not fit \
the 16-bit on-disk field"
))
})?;
let obj = coll.get_object(idx).ok_or_else(|| {
crate::io::IoError::InvalidState(format!(
"dataset {name:?} virtual mapping list object {idx} not found in the \
global heap collection at address {heap_address:#x}"
))
})?;
VirtualMappingList::decode(obj, ctx).map_err(|e| {
crate::io::IoError::InvalidState(format!(
"dataset {name:?} has a malformed virtual dataset mapping list: {e}"
))
})
}
pub(crate) fn resolve_external_file_slots(
handle: &mut FileHandle,
ctx: &FormatContext,
efl: &ExternalFileListMessage,
) -> IoResult<Vec<ExternalFileSegment>> {
let sa = ctx.sizeof_addr as usize;
let ss = ctx.sizeof_size as usize;
let heap_hdr_buf = handle.read_at_most(efl.heap_addr, 64)?;
let heap_hdr = LocalHeapHeader::decode(&heap_hdr_buf, sa, ss)?;
let heap_data = handle.read_at(heap_hdr.data_addr, heap_hdr.data_size as usize)?;
efl.slots
.iter()
.map(|slot| {
let name = local_heap_get_string(&heap_data, slot.name_offset)?;
Ok(ExternalFileSegment {
name,
offset: slot.offset,
size: slot.size,
})
})
.collect()
}
pub fn dataset_names(&self) -> Vec<&str> {
let mut names: Vec<&str> = self.datasets.iter().map(|d| d.name.as_str()).collect();
names.extend(self.unreadable.keys().map(String::as_str));
names
}
pub fn unreadable_reason(&mut self, path: &str) -> Option<&str> {
if self.external_edge(path).is_some() {
let (owner, local, _) = self.external_owner(path, MAX_EXTERNAL_HOPS).ok()?;
let local = owner.canonical_path(&local);
return owner.unreadable.get(&local).map(String::as_str);
}
let path = self.canonical_path(path);
self.unreadable.get(&path).map(String::as_str)
}
pub fn links(&self) -> &std::collections::BTreeMap<String, LinkClass> {
&self.links
}
pub fn named_datatype_names(&self) -> Vec<&str> {
self.datatypes.keys().map(String::as_str).collect()
}
pub fn named_datatype(&mut self, path: &str) -> IoResult<&DatatypeMessage> {
self.named_datatype_info(path)?
.datatype()
.map_err(|why| crate::io::IoError::Unsupported(why.to_string()))
}
pub fn named_datatype_attr_names(&mut self, path: &str) -> IoResult<Vec<String>> {
let mut names: Vec<String> = self
.named_datatype_info(path)?
.attributes()
.iter()
.map(|a| a.name.clone())
.collect();
names.sort();
Ok(names)
}
pub fn named_datatype_header_attr_count(&mut self, path: &str) -> IoResult<u64> {
Ok(self.named_datatype_info(path)?.attributes().len() as u64)
}
pub fn named_datatype_attr(
&mut self,
path: &str,
attr_name: &str,
) -> IoResult<&AttributeMessage> {
let owned = attr_name.to_string();
self.named_datatype_info(path)?
.attributes()
.iter()
.find(|a| a.name == owned)
.ok_or_else(|| crate::io::IoError::NotFound(format!("{path}:{attr_name}")))
}
pub fn named_datatype_info(&mut self, path: &str) -> IoResult<&CommittedDatatypeInfo> {
if self.external_edge(path).is_some() {
let (owner, local, _) = self.external_owner(path, MAX_EXTERNAL_HOPS)?;
let local = owner.canonical_path(&local);
return owner
.datatypes
.get(&local)
.ok_or(crate::io::IoError::NotFound(local));
}
let local = self.canonical_path(path);
self.datatypes
.get(&local)
.ok_or(crate::io::IoError::NotFound(local))
}
pub fn link_class(&mut self, path: &str) -> Option<&LinkClass> {
let path = path.trim_start_matches('/');
if self.links.contains_key(path) {
return self.links.get(path);
}
if self.external_edge(path).is_some() {
let (owner, local, _) = self.external_owner(path, MAX_EXTERNAL_HOPS).ok()?;
let local = owner.canonical_path(&local);
return owner.links.get(&local);
}
let path = self.canonical_path(path);
self.links.get(&path)
}
fn traverse(&self, name: &str) -> Traversal {
const MAX_TRAVERSALS: usize = 64;
let mut name = name.trim_start_matches('/').to_string();
let mut via = None;
for _ in 0..MAX_TRAVERSALS {
let Some((prefix, rewrite)) = self.longest_rewrite(&name) else {
break;
};
let rest = name[prefix.len()..].to_string();
match rewrite {
Rewrite::Alias(first) => {
name = format!("{first}{rest}").trim_start_matches('/').to_string();
}
Rewrite::Soft(target) => {
let resolved = resolve_link_value(prefix, target);
via = Some(SoftLinkRef {
link: prefix.to_string(),
target: target.to_string(),
});
name = format!("{resolved}{rest}")
.trim_start_matches('/')
.to_string();
}
Rewrite::External { file, path } => {
return Traversal::External {
link: prefix.to_string(),
file: file.to_string(),
path: format!("{path}{rest}"),
};
}
}
}
Traversal::Path { path: name, via }
}
fn longest_rewrite<'a>(&'a self, path: &str) -> Option<(&'a str, Rewrite<'a>)> {
let mut end = path.len();
loop {
let candidate = &path[..end];
if let Some((alias, first)) = self.group_aliases.get_key_value(candidate) {
return Some((alias.as_str(), Rewrite::Alias(first)));
}
match self.links.get_key_value(candidate) {
Some((link, LinkClass::Soft { path })) => {
return Some((link.as_str(), Rewrite::Soft(path)))
}
Some((link, LinkClass::External { file, path })) => {
return Some((link.as_str(), Rewrite::External { file, path }))
}
_ => {}
}
end = candidate.rfind('/')?;
}
}
pub fn superblock_extension(&self) -> &SuperblockExtension {
&self.ext
}
pub fn tracked_free_space(&mut self) -> IoResult<u64> {
let Some(info) = self.ext.file_space_info.clone() else {
return Ok(0);
};
crate::io::free_space_io::tracked_free_space(&mut self.handle, &self.meta.ctx, &info)
}
pub fn userblock_size(&self) -> u64 {
self.handle.base()
}
pub fn superblock_version(&self) -> u8 {
self.superblock_version
}
pub fn canonical_path(&self, name: &str) -> String {
match self.traverse(name) {
Traversal::Path { path, .. } => path,
Traversal::External { .. } => name.trim_start_matches('/').to_string(),
}
}
pub(crate) fn external_edge(&self, name: &str) -> Option<ExternalEdge> {
match self.traverse(name) {
Traversal::Path { .. } => None,
Traversal::External { link, file, path } => Some(ExternalEdge { link, file, path }),
}
}
fn prefix_open_candidates(
&self,
env_var: &str,
prop_prefix: Option<&Path>,
file: &str,
) -> Vec<PathBuf> {
let raw = Path::new(file);
let mut candidates = Vec::new();
if raw.is_absolute() {
candidates.push(raw.to_path_buf());
}
let base: &Path = if raw.is_absolute() {
Path::new(raw.file_name().unwrap_or(raw.as_os_str()))
} else {
raw
};
if let Ok(prefixes) = std::env::var(env_var) {
candidates.extend(
prefixes
.split(':')
.filter(|p| !p.is_empty())
.map(|p| Path::new(p).join(base)),
);
}
if let Some(prefix) = prop_prefix {
candidates.push(prefix.join(base));
}
if let Some(dir) = self.path.parent().filter(|d| !d.as_os_str().is_empty()) {
candidates.push(dir.join(base));
}
candidates.push(base.to_path_buf());
if !self.source_dir.as_os_str().is_empty() {
candidates.push(self.source_dir.join(base));
}
candidates
}
pub(crate) fn set_elink_prefix(&mut self, prefix: Option<String>) {
self.elink_prefix = prefix;
}
fn external_candidates(&self, file: &str) -> Vec<PathBuf> {
let prop = self.elink_prefix.as_deref().map(Path::new);
self.prefix_open_candidates("HDF5_EXT_PREFIX", prop, file)
}
fn vds_candidates(&self, access: &DatasetAccess, file: &str) -> Vec<PathBuf> {
let prop = resolve_vdsfile_prefix(access.virtual_prefix_value(), &self.source_dir);
self.prefix_open_candidates("HDF5_VDS_PREFIX", prop.as_deref(), file)
}
fn external_target(&mut self, link: &str, file: &str) -> IoResult<&mut Hdf5Reader> {
let resolved = match self.external_resolved.get(file) {
Some(resolved) => resolved.clone(),
None => {
let candidates = self.external_candidates(file);
let resolved = candidates
.iter()
.find(|p| p.is_file())
.cloned()
.ok_or_else(|| crate::io::IoError::ExternalFileNotFound {
link: link.to_string(),
file: file.to_string(),
searched: candidates.iter().map(|p| p.display().to_string()).collect(),
})?;
self.external_resolved
.insert(file.to_string(), resolved.clone());
resolved
}
};
self.cross_file(resolved, CrossFileOwner::Reader)
}
fn cross_file(
&mut self,
resolved: PathBuf,
owner: CrossFileOwner,
) -> IoResult<&mut Hdf5Reader> {
let locking = self.locking;
let elink_prefix = self.elink_prefix.clone();
match self.external.entry(resolved) {
std::collections::btree_map::Entry::Occupied(e) => {
let e = e.into_mut();
e.owner.widen(owner);
Ok(&mut *e.reader)
}
std::collections::btree_map::Entry::Vacant(e) => {
let mut reader = Hdf5Reader::open_with_locking(e.key(), locking)?;
reader.elink_prefix.clone_from(&elink_prefix);
Ok(&mut *e
.insert(CrossFileEntry {
reader: Box::new(reader),
owner,
})
.reader)
}
}
}
pub(crate) fn release_closed_virtual_sources(&mut self) {
let dead: Vec<PathBuf> = self
.external
.iter()
.filter(|(_, e)| match &e.owner {
CrossFileOwner::Reader => false,
CrossFileOwner::VirtualOpens(vds) => !vds.iter().any(|v| self.is_open_dataset(v)),
})
.map(|(p, _)| p.clone())
.collect();
for path in dead {
self.external.remove(&path);
}
}
fn is_open_dataset(&self, name: &str) -> bool {
self.dataset_access
.get(name)
.is_some_and(|e| e.open.strong_count() > 0)
}
fn vds_source_file(&mut self, vds: &str, file_name: &str) -> Option<&mut Hdf5Reader> {
let key = (vds.to_string(), file_name.to_string());
let resolved = match self.vds_resolved.get(&key) {
Some(resolved) => resolved.clone(),
None => {
let access = self.access_in_force(vds);
let resolved = self
.vds_candidates(&access, file_name)
.into_iter()
.find(|p| p.is_file())?;
self.vds_resolved.insert(key, resolved.clone());
resolved
}
};
self.cross_file(resolved, CrossFileOwner::virtual_open(vds))
.ok()
}
fn external_owner(
&mut self,
name: &str,
hops: usize,
) -> IoResult<(&mut Self, String, Option<ExternalEdge>)> {
let path = name.trim_start_matches('/').to_string();
let Some(edge) = self.external_edge(&path) else {
return Ok((self, path, None));
};
if hops == 0 {
return Err(crate::io::IoError::InvalidState(format!(
"resolving '{name}' crossed more than {MAX_EXTERNAL_HOPS} external links \
(libhdf5 stops at the same H5L_NUM_LINKS); the links may form a cycle"
)));
}
let target = self.external_target(&edge.link, &edge.file)?;
let (owner, path, deeper) = target.external_owner(&edge.path, hops - 1)?;
Ok((owner, path, deeper.or(Some(edge))))
}
pub fn dataset_info(&mut self, name: &str) -> Option<&DatasetReadInfo> {
if self.external_edge(name).is_some() {
let (owner, path, _) = self.external_owner(name, MAX_EXTERNAL_HOPS).ok()?;
return owner.dataset_info_local(&path);
}
self.dataset_info_local(name)
}
fn dataset_info_local(&self, name: &str) -> Option<&DatasetReadInfo> {
let name = self.canonical_path(name);
self.datasets.get(&name)
}
fn dataset_position(&self, name: &str) -> IoResult<usize> {
let canonical = self.canonical_path(name);
self.datasets
.position(&canonical)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))
}
pub fn open_dataset_with(
&mut self,
name: &str,
access: &DatasetAccess,
) -> IoResult<(Option<DatasetOpenToken>, &DatasetReadInfo)> {
if self.external_edge(name).is_some() {
let (owner, path, edge) = self.external_owner(name, MAX_EXTERNAL_HOPS)?;
let open = owner.apply_dataset_access(&path, access)?;
return match owner.open_dataset_local(&path) {
Err(crate::io::IoError::NotFound(_)) => {
Err(edge.map_or_else(|| crate::io::IoError::NotFound(path), |e| e.dangling()))
}
other => other.map(|info| (open, info)),
};
}
let open = self.apply_dataset_access(name, access)?;
self.open_dataset_local(name).map(|info| (open, info))
}
fn open_dataset_local(&self, name: &str) -> IoResult<&DatasetReadInfo> {
let Traversal::Path { path, via } = self.traverse(name) else {
return Err(crate::io::IoError::NotFound(name.to_string()));
};
if let Some(info) = self.datasets.get(&path) {
return Ok(info);
}
if let Some(why) = self.unreadable.get(&path) {
return Err(crate::io::IoError::Unsupported(format!(
"'{name}' is a dataset this crate cannot read: {why}"
)));
}
if let Some(SoftLinkRef { link, target }) = via {
return Err(crate::io::IoError::DanglingLink { link, target });
}
Err(crate::io::IoError::NotFound(name.to_string()))
}
fn resolve_attr<'a>(
attrs: &'a ObjectAttributes,
owner: &str,
name: &str,
) -> IoResult<&'a AttributeMessage> {
match attrs.entries.iter().find(|a| a.name() == name) {
Some(entry) => entry.decoded().map_err(|reason| {
crate::io::IoError::Unsupported(format!(
"attribute '{name}' on '{owner}' cannot be decoded: {reason}"
))
}),
None => match attrs.unreadable_reason() {
Some(reason) => Err(incomplete_error(owner, reason)),
None => Err(crate::io::IoError::NotFound(format!("{owner}:{name}"))),
},
}
}
fn attr_reason<'a>(attrs: &'a ObjectAttributes, name: &str) -> Option<&'a str> {
attrs
.entries
.iter()
.find(|a| a.name() == name)?
.unreadable_reason()
}
pub fn dataset_attr_unreadable_reason(
&mut self,
ds_name: &str,
attr_name: &str,
) -> Option<&str> {
Self::attr_reason(&self.dataset_info(ds_name)?.attributes, attr_name)
}
pub fn root_attr_unreadable_reason(&self, name: &str) -> Option<&str> {
Self::attr_reason(&self.root_attributes, name)
}
pub fn group_attr_unreadable_reason(&self, group_path: &str, name: &str) -> Option<&str> {
Self::attr_reason(
self.group_attributes
.get(&self.canonical_path(group_path))?,
name,
)
}
pub fn dataset_attrs_unreadable_reason(&mut self, ds_name: &str) -> Option<&str> {
self.dataset_info(ds_name)?.attributes.unreadable_reason()
}
pub fn dataset_attr_storage(&mut self, ds_name: &str) -> IoResult<AttributeStorage> {
Ok(self
.dataset_info(ds_name)
.ok_or_else(|| crate::io::IoError::NotFound(ds_name.to_string()))?
.attributes
.storage())
}
pub fn dataset_header_attr_count(&mut self, ds_name: &str) -> IoResult<u64> {
let info = self
.dataset_info(ds_name)
.ok_or_else(|| crate::io::IoError::NotFound(ds_name.to_string()))?;
info.attributes.header_count(ds_name)
}
pub fn root_attrs_unreadable_reason(&self) -> Option<&str> {
self.root_attributes.unreadable_reason()
}
pub fn group_attrs_unreadable_reason(&self, group_path: &str) -> Option<&str> {
self.group_attributes
.get(&self.canonical_path(group_path))?
.unreadable_reason()
}
pub fn dataset_attr_names(&mut self, name: &str) -> IoResult<Vec<String>> {
let info = self
.dataset_info(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?;
info.attributes.ordered_names(name)
}
pub fn dataset_attr(&mut self, ds_name: &str, attr_name: &str) -> IoResult<&AttributeMessage> {
let info = self
.dataset_info(ds_name)
.ok_or_else(|| crate::io::IoError::NotFound(ds_name.to_string()))?;
Self::resolve_attr(&info.attributes, ds_name, attr_name)
}
pub fn root_attr_names(&self) -> IoResult<Vec<String>> {
self.root_attributes.ordered_names("/")
}
pub fn root_attr(&self, name: &str) -> IoResult<&AttributeMessage> {
Self::resolve_attr(&self.root_attributes, "/", name)
}
pub fn root_attr_creation_order(&self) -> CreationOrder {
self.root_attributes.creation_order()
}
pub fn root_attr_storage(&self) -> AttributeStorage {
self.root_attributes.storage()
}
pub fn root_header_attr_count(&self) -> IoResult<u64> {
self.root_attributes.header_count("/")
}
pub fn root_link_creation_order(&self) -> CreationOrder {
self.root_link_storage.1
}
pub fn root_link_storage(&self) -> LinkStorage {
self.root_link_storage.0
}
pub fn group_attr_names(&mut self, group_path: &str) -> IoResult<Vec<String>> {
if self.external_edge(group_path).is_some() {
let (owner, local, _) = self.external_owner(group_path, MAX_EXTERNAL_HOPS)?;
if local.is_empty() {
return owner.root_attr_names();
}
return owner.group_attr_names_local(&local);
}
self.group_attr_names_local(group_path)
}
fn group_attr_names_local(&self, group_path: &str) -> IoResult<Vec<String>> {
let Some(attrs) = self.group_attributes.get(&self.canonical_path(group_path)) else {
return Ok(Vec::new());
};
attrs.ordered_names(group_path)
}
pub fn group_attr_creation_order(&self, group_path: &str) -> CreationOrder {
self.group_attributes
.get(&self.canonical_path(group_path))
.map(ObjectAttributes::creation_order)
.unwrap_or_default()
}
pub fn group_attr_storage(&self, group_path: &str) -> AttributeStorage {
self.group_attributes
.get(&self.canonical_path(group_path))
.map(ObjectAttributes::storage)
.unwrap_or_default()
}
pub fn group_header_attr_count(&self, group_path: &str) -> IoResult<u64> {
let Some(attrs) = self.group_attributes.get(&self.canonical_path(group_path)) else {
return Ok(0);
};
attrs.header_count(group_path)
}
pub fn group_link_creation_order(&self, group_path: &str) -> CreationOrder {
self.group_link_storage
.get(&self.canonical_path(group_path))
.map_or(CreationOrder::Untracked, |(_, order)| *order)
}
pub fn group_link_storage(&self, group_path: &str) -> LinkStorage {
self.group_link_storage
.get(&self.canonical_path(group_path))
.map_or(LinkStorage::Compact, |(storage, _)| *storage)
}
pub fn group_attr(&mut self, group_path: &str, name: &str) -> IoResult<&AttributeMessage> {
if self.external_edge(group_path).is_some() {
let (owner, local, _) = self.external_owner(group_path, MAX_EXTERNAL_HOPS)?;
if local.is_empty() {
return owner.root_attr(name);
}
return owner.group_attr_local(&local, name);
}
self.group_attr_local(group_path, name)
}
fn group_attr_local(&self, group_path: &str, name: &str) -> IoResult<&AttributeMessage> {
match self.group_attributes.get(&self.canonical_path(group_path)) {
Some(attrs) => Self::resolve_attr(attrs, group_path, name),
None => Err(crate::io::IoError::NotFound(format!("{group_path}:{name}"))),
}
}
pub fn group_paths(&self) -> &std::collections::BTreeSet<String> {
&self.group_paths
}
pub fn has_group(&self, group_path: &str) -> bool {
if group_path.is_empty() || self.group_paths.contains(group_path) {
return true;
}
let canon = self.canonical_path(group_path);
canon.is_empty() || self.group_paths.contains(&canon)
}
fn read_heap_collection(&mut self, addr: u64) -> IoResult<GlobalHeapCollection> {
read_heap_collection_from(&mut self.handle, &self.meta.ctx, addr)
}
pub fn attr_string_value(&mut self, attr: &AttributeMessage) -> IoResult<String> {
use crate::format::messages::datatype::DatatypeMessage;
if !matches!(attr.datatype, DatatypeMessage::VarLenString { .. }) {
return fixed_string_attr_value(attr);
}
if attr.data.len() < vlen_reference_size(&self.meta.ctx) {
return Ok(String::new());
}
let (_seq, coll_addr, obj_index) = decode_vlen_reference(&attr.data, &self.meta.ctx)?;
if coll_addr == UNDEF_ADDR || coll_addr == 0 {
return Ok(String::new());
}
let coll = self.read_heap_collection(coll_addr)?;
let idx = u16::try_from(obj_index).map_err(|_| {
crate::io::IoError::InvalidState(format!(
"global heap object index {obj_index} does not fit the 16-bit on-disk field"
))
})?;
let obj = coll.get_object(idx).ok_or_else(|| {
crate::io::IoError::InvalidState(format!(
"global heap object {idx} not found in the collection at address {coll_addr:#x}"
))
})?;
Ok(String::from_utf8_lossy(obj).to_string())
}
pub fn path_for_object(&self, addr: u64) -> Option<&str> {
self.object_paths.get(&addr).map(String::as_str)
}
pub fn object_message_storage(&mut self, path: &str) -> IoResult<Vec<(u8, MessageStorage)>> {
let addr = self.object_header_address(path)?;
crate::io::object_header_io::read_header_message_storage(&mut self.handle, &self.meta, addr)
}
pub fn object_message_flags(&mut self, path: &str) -> IoResult<Vec<(u8, u8)>> {
let addr = self.object_header_address(path)?;
crate::io::object_header_io::read_header_message_flags(&mut self.handle, &self.meta, addr)
}
pub fn object_datatype_versions(
&mut self,
path: &str,
) -> IoResult<Vec<crate::format::messages::datatype::DatatypeNodeVersion>> {
let addr = self.object_header_address(path)?;
crate::io::object_header_io::read_header_datatype_versions(
&mut self.handle,
&self.meta,
addr,
)
}
pub fn object_records_times(&mut self, path: &str) -> IoResult<bool> {
let addr = self.object_header_address(path)?;
Ok(crate::io::object_header_io::read_header_recorded_times(
&mut self.handle,
&self.meta,
addr,
)?
.is_some())
}
fn object_header_address(&mut self, path: &str) -> IoResult<u64> {
if self.external_edge(path).is_some() {
return Err(crate::io::IoError::NotFound(format!(
"{path} is in another file; its header is not this file's to read"
)));
}
match self.dataset_info(path) {
Some(info) => Ok(info.object_header_address),
None => {
let want = absolute_path(&self.canonical_path(path));
self.object_paths
.iter()
.find(|(_, p)| **p == want)
.map(|(addr, _)| *addr)
.ok_or_else(|| crate::io::IoError::NotFound(path.to_string()))
}
}
}
pub fn read_references(&mut self, name: &str) -> IoResult<Vec<Reference>> {
let datatype = self
.dataset_info(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?
.datatype
.clone();
let raw = self.read_dataset_raw(name)?;
self.decode_references(&datatype, &raw)
}
pub fn attr_references(&mut self, attr: &AttributeMessage) -> IoResult<Vec<Reference>> {
self.decode_references(&attr.datatype, &attr.data)
}
fn decode_references(
&mut self,
datatype: &DatatypeMessage,
bytes: &[u8],
) -> IoResult<Vec<Reference>> {
let DatatypeMessage::Reference { size, kind } = datatype else {
return Err(crate::io::IoError::InvalidState(format!(
"datatype {datatype} is not a reference"
)));
};
let (size, kind) = (*size as usize, *kind);
if size == 0 {
return Err(crate::io::IoError::InvalidState(
"reference datatype has zero width".into(),
));
}
let encoding = kind.encoding();
let mut heaps = std::collections::HashMap::new();
let mut out = Vec::with_capacity(bytes.len() / size);
for elem in bytes.chunks_exact(size) {
out.push(self.decode_reference_element(elem, encoding, &mut heaps)?);
}
Ok(out)
}
fn decode_reference_element(
&mut self,
elem: &[u8],
encoding: ReferenceEncoding,
heaps: &mut std::collections::HashMap<u64, GlobalHeapCollection>,
) -> IoResult<Reference> {
match encoding {
ReferenceEncoding::Old(OldReferenceKind::Object) => {
match decode_object_element(elem, &self.meta.ctx)? {
None => Ok(Reference::Null),
Some(address) => Ok(self.resolve_reference(DecodedReference {
address,
file: None,
target: ReferenceTarget::Object,
})),
}
}
ReferenceEncoding::Old(OldReferenceKind::DatasetRegion) => {
let Some((coll_addr, obj_index)) = decode_region_element(elem, &self.meta.ctx)?
else {
return Ok(Reference::Null);
};
let obj = self.heap_object(coll_addr, obj_index, heaps)?;
let (address, selection) = decode_region_heap_object(obj, &self.meta.ctx)?;
Ok(self.resolve_reference(DecodedReference {
address,
file: None,
target: ReferenceTarget::Region(selection),
}))
}
ReferenceEncoding::Revised => {
let (kind, external, body) = match decode_revised_element(elem, &self.meta.ctx)? {
RevisedElement::Null => return Ok(Reference::Null),
RevisedElement::Inline { kind, body } => (kind, false, body.to_vec()),
RevisedElement::Heap {
kind,
external,
collection,
index,
} => (
kind,
external,
self.heap_object(collection, index, heaps)?.to_vec(),
),
};
match decode_revised_body(kind, external, &body, &self.meta.ctx)? {
None => Ok(Reference::Null),
Some(decoded) => Ok(self.resolve_reference(decoded)),
}
}
}
}
fn resolve_reference(&mut self, decoded: DecodedReference) -> Reference {
let DecodedReference {
address,
file,
target,
} = decoded;
let path = match &file {
None => self.path_for_object(address).map(str::to_string),
Some(name) => self
.cross_file(PathBuf::from(name), CrossFileOwner::Reader)
.ok()
.and_then(|target| target.path_for_object(address).map(str::to_string)),
};
match target {
ReferenceTarget::Object => Reference::Object {
address,
file,
path,
},
ReferenceTarget::Region(selection) => Reference::Region {
address,
file,
path,
selection,
},
ReferenceTarget::Attribute(name) => Reference::Attr {
address,
file,
path,
name,
},
}
}
fn heap_object<'h>(
&mut self,
collection: u64,
index: u32,
heaps: &'h mut std::collections::HashMap<u64, GlobalHeapCollection>,
) -> IoResult<&'h [u8]> {
if let std::collections::hash_map::Entry::Vacant(slot) = heaps.entry(collection) {
slot.insert(self.read_heap_collection(collection)?);
}
let idx = u16::try_from(index).map_err(|_| {
crate::io::IoError::InvalidState(format!(
"global heap object index {index} does not fit the 16-bit on-disk field"
))
})?;
heaps[&collection].get_object(idx).ok_or_else(|| {
crate::io::IoError::InvalidState(format!(
"global heap object {idx} not found in the collection at address {collection:#x}"
))
})
}
pub fn dataset_shape(&mut self, name: &str) -> IoResult<Vec<u64>> {
let info = self
.dataset_info(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?;
Ok(info.dataspace.dims.clone())
}
fn raw_size_and_datatype(&self, name: &str) -> IoResult<(DatatypeMessage, u64)> {
let info = self
.dataset_info_local(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?;
Ok((info.datatype.clone(), Self::raw_size_of(info)))
}
fn raw_size_of(info: &DatasetReadInfo) -> u64 {
if info.dataspace.is_null() {
0
} else {
saturating_byte_len(&info.dataspace.dims, info.datatype.element_size() as u64)
}
}
pub fn dataset_raw_size(&mut self, name: &str) -> IoResult<u64> {
if self.external_edge(name).is_some() {
let (owner, path, _) = self.external_owner(name, MAX_EXTERNAL_HOPS)?;
return owner.dataset_raw_size(&path);
}
let info = self
.dataset_info_local(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?;
Ok(Self::raw_size_of(info))
}
#[cfg(feature = "mmap")]
pub(crate) fn dataset_view_source(&mut self, name: &str) -> IoResult<DatasetViewSource> {
if self.external_edge(name).is_some() {
let (owner, path, _) = self.external_owner(name, MAX_EXTERNAL_HOPS)?;
return owner.dataset_view_source(&path);
}
let base = self.handle.base();
let map = self.handle.map_snapshot();
let info = self
.dataset_info_local(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?;
let len = Self::raw_size_of(info);
let storage = match &info.layout {
DataLayoutMessage::Contiguous { .. } if !info.external_files.is_empty() => {
ViewStorage::Elsewhere("its raw data is in external data files")
}
DataLayoutMessage::Contiguous { address, .. } if *address == UNDEF_ADDR => {
ViewStorage::Unallocated
}
DataLayoutMessage::Contiguous { address, .. } => {
let offset = address.checked_add(base).ok_or_else(|| {
crate::io::IoError::InvalidState(format!(
"dataset '{name}' claims raw data at {address}, which overflows \
past the userblock at {base}"
))
})?;
ViewStorage::Contiguous { offset, len }
}
DataLayoutMessage::Compact { .. } => {
ViewStorage::Elsewhere("its raw data is compact, stored inside the object header")
}
DataLayoutMessage::ChunkedV3 { .. } | DataLayoutMessage::ChunkedV4 { .. } => {
ViewStorage::Elsewhere("it is chunked")
}
DataLayoutMessage::Virtual { .. } => {
ViewStorage::Elsewhere("it is virtual, mapped from other datasets")
}
};
Ok(DatasetViewSource {
map,
storage,
datatype: info.datatype.clone(),
dims: info.dataspace.dims.clone(),
})
}
pub fn read_dataset_raw(&mut self, name: &str) -> IoResult<Vec<u8>> {
if self.external_edge(name).is_some() {
let (owner, path, _) = self.external_owner(name, MAX_EXTERNAL_HOPS)?;
return owner.read_dataset_raw(&path);
}
let (datatype, total) = self.raw_size_and_datatype(name)?;
read_image_into_new(total as usize, |data| {
self.read_dataset_raw_into_unconverted(name, data, ReadDst::Fresh)?;
Self::apply_post_filter_conversion(data, &datatype)
})
}
pub fn read_dataset_raw_into(&mut self, name: &str, out: &mut [u8]) -> IoResult<()> {
self.read_dataset_raw_into_dst(name, out, ReadDst::Reused)
}
pub(crate) fn read_dataset_raw_into_dst(
&mut self,
name: &str,
out: &mut [u8],
dst: ReadDst,
) -> IoResult<()> {
if self.external_edge(name).is_some() {
let (owner, path, _) = self.external_owner(name, MAX_EXTERNAL_HOPS)?;
return owner.read_dataset_raw_into_dst(&path, out, dst);
}
let (datatype, total) = self.raw_size_and_datatype(name)?;
if out.len() as u64 != total {
return Err(crate::io::IoError::InvalidState(format!(
"read_dataset_raw_into: buffer is {} bytes but dataset needs {}",
out.len(),
total
)));
}
self.read_dataset_raw_into_unconverted(name, out, dst)?;
Self::apply_post_filter_conversion(out, &datatype)?;
Ok(())
}
fn read_dataset_raw_into_unconverted(
&mut self,
name: &str,
out: &mut [u8],
dst: ReadDst,
) -> IoResult<()> {
let info = self
.dataset_info_local(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?;
let layout = info.layout.clone();
let pipeline = info.filter_pipeline.clone();
let fill_value = info.fill_value.clone();
let external_files = info.external_files.clone();
match &layout {
DataLayoutMessage::Contiguous { .. } if !external_files.is_empty() => {
let prefix = self.extfile_prefix_in_force(name);
read_external_file_bytes(&external_files, prefix.as_deref(), 0, out)?;
}
DataLayoutMessage::Contiguous { address, .. } => {
if *address == UNDEF_ADDR {
fill_tiled_into(out, fill_value.as_deref());
} else {
self.handle.read_exact_at_into(*address, out, dst)?;
}
}
DataLayoutMessage::Compact { data } => {
let n = out.len().min(data.len());
out[..n].copy_from_slice(&data[..n]);
if n < out.len() {
fill_tiled_into(&mut out[n..], fill_value.as_deref());
}
}
DataLayoutMessage::ChunkedV3 {
chunk_dims,
b_tree_address,
} => {
let real_chunk_dims = &chunk_dims[..chunk_dims.len() - 1];
self.read_chunked_btree_v1(
name,
real_chunk_dims,
*b_tree_address,
ChunkReadRequest {
pipeline: pipeline.as_ref(),
target: ChunkTarget::Full,
fill_value: fill_value.as_deref(),
dst,
},
out,
)?;
}
DataLayoutMessage::ChunkedV4 {
chunk_dims,
index_address,
index_type,
earray_params,
single_chunk_filter,
..
} => {
let real_chunk_dims = &chunk_dims[..chunk_dims.len() - 1];
self.read_chunked_v4(
name,
real_chunk_dims,
ChunkIndexDesc {
index_type: *index_type,
index_address: *index_address,
earray_params: earray_params.as_ref(),
single_chunk_filter: *single_chunk_filter,
},
ChunkReadRequest {
pipeline: pipeline.as_ref(),
target: ChunkTarget::Full,
fill_value: fill_value.as_deref(),
dst,
},
out,
)?;
}
DataLayoutMessage::Virtual { .. } => {
fill_tiled_into(out, fill_value.as_deref());
self.read_virtual_into(name, out, 0)?;
}
}
Ok(())
}
fn resolve_virtual_extents(&mut self) -> IoResult<()> {
let targets: Vec<usize> = self
.datasets
.iter()
.enumerate()
.filter(|(_, d)| d.virtual_mappings.is_some())
.map(|(i, _)| i)
.collect();
for i in targets {
let stored = self.datasets[i].dataspace.dims.clone();
self.datasets.entry_mut(i).virtual_stored_dims = Some(stored.clone());
self.resolve_virtual_extent_of(i, &stored)?;
}
self.release_closed_virtual_sources();
Ok(())
}
fn resolve_virtual_extent_of(&mut self, i: usize, stored: &[u64]) -> IoResult<()> {
let Some(mappings) = self.datasets[i].virtual_mappings.clone() else {
return Ok(());
};
let vds = self.datasets[i].name.clone();
let access = self.access_in_force(&vds);
let (resolution, dims) =
self.resolve_one_virtual_extent(&vds, &mappings, stored, &access)?;
let entry = self.datasets.entry_mut(i);
entry.dataspace.dims = dims;
entry.virtual_resolution = Some(resolution);
Ok(())
}
fn access_in_force(&self, canonical: &str) -> DatasetAccess {
self.dataset_access
.get(canonical)
.map(|e| e.access.clone())
.unwrap_or_default()
}
fn apply_dataset_access(
&mut self,
name: &str,
access: &DatasetAccess,
) -> IoResult<Option<DatasetOpenToken>> {
let canonical = self.canonical_path(name);
let Some(i) = self.datasets.position(&canonical) else {
return Ok(None);
};
if let Some(open) = self
.dataset_access
.get(&canonical)
.and_then(|e| e.open.upgrade())
{
let in_force = self.extfile_prefix_of(&self.access_in_force(&canonical));
if in_force != self.extfile_prefix_of(access) {
return Err(crate::io::IoError::InvalidState(format!(
"dataset {canonical:?} is already open under a different external file \
prefix, and libhdf5 refuses to join an open that disagrees about one"
)));
}
return Ok(Some(open));
}
let token: DatasetOpenToken = std::sync::Arc::new(());
let unchanged = &self.access_in_force(&canonical) == access;
let access = access.clone();
self.dataset_access.insert(
canonical,
AccessInForce {
access,
open: std::sync::Arc::downgrade(&token),
},
);
if !unchanged {
if let Some(stored) = self.datasets[i].virtual_stored_dims.clone() {
VirtualResolveDepth::enter(|| self.resolve_virtual_extent_of(i, &stored))?;
}
}
Ok(Some(token))
}
fn extfile_prefix_of(&self, access: &DatasetAccess) -> Option<PathBuf> {
resolve_extfile_prefix(access.efile_prefix_value(), &self.source_dir)
}
fn extfile_prefix_in_force(&self, name: &str) -> Option<PathBuf> {
let canonical = self.canonical_path(name);
self.extfile_prefix_of(&self.access_in_force(&canonical))
}
fn resolve_one_virtual_extent(
&mut self,
vds: &str,
mappings: &VirtualMappingList,
curr_dims: &[u64],
access: &DatasetAccess,
) -> IoResult<(Vec<MappingResolution>, Vec<u64>)> {
let rank = curr_dims.len();
let mut resolution = Vec::with_capacity(mappings.mappings.len());
let mut new_dims: Vec<Option<u64>> = vec![None; rank];
let mut min_dims = vec![0u64; rank];
let incl_trail = access.view() == VirtualView::FirstMissing;
let take_clip = |slot: &mut Option<u64>, clip: u64| {
if slot.is_none_or(|d| if incl_trail { clip < d } else { clip > d }) {
*slot = Some(clip);
}
};
for m in &mappings.mappings {
let unlim_virtual = m.virtual_selection.unlim_dim();
let res = match (unlim_virtual, m.source_selection.unlim_dim()) {
(Some(vd), Some(sd)) => {
let source_clip = self
.virtual_source_dims(vds, m, access)
.ok()
.flatten()
.and_then(|d| d.get(sd).copied())
.unwrap_or(0);
let virtual_clip = match (
regular_hyperslab(&m.virtual_selection),
regular_hyperslab(&m.source_selection),
) {
(Some(v), Some(sr)) => {
v.clip_extent(sr.num_slices(source_clip), incl_trail)
}
_ => 0,
};
take_clip(&mut new_dims[vd], virtual_clip);
MappingResolution::Unlimited {
virtual_clip,
source_clip,
}
}
(Some(vd), None) => {
let (blocks, present) = self.printf_blocks_present(vds, m, access);
let virtual_clip = match (blocks, regular_hyperslab(&m.virtual_selection)) {
(0, _) | (_, None) => 0,
(n, Some(r)) => match access.view() {
VirtualView::LastAvailable => {
let last = r.unlim_block(n - 1);
last.start[vd] + last.block[vd]
}
VirtualView::FirstMissing => r.unlim_block(n).start[vd],
},
};
take_clip(&mut new_dims[vd], virtual_clip);
MappingResolution::Printf { blocks, present }
}
_ => MappingResolution::Bounded,
};
if let Some((_, hi)) = m.virtual_selection.bounds() {
for (d, &e) in hi.iter().enumerate().take(rank) {
if Some(d) != unlim_virtual && e + 1 > min_dims[d] {
min_dims[d] = e + 1;
}
}
}
resolution.push(res);
}
let dims = (0..rank)
.map(|d| match new_dims[d] {
Some(v) => v.max(min_dims[d]),
None => curr_dims[d],
})
.collect();
Ok((resolution, dims))
}
fn virtual_source_dims(
&mut self,
vds: &str,
m: &VirtualMapping,
access: &DatasetAccess,
) -> IoResult<Option<Vec<u64>>> {
let m = built_names(m, 0)?;
Ok(self.source_dims(vds, &m.source_file_name, &m.source_dset_name, access))
}
fn printf_blocks_present(
&mut self,
vds: &str,
m: &VirtualMapping,
access: &DatasetAccess,
) -> (u64, Vec<u64>) {
let gap = access.effective_printf_gap();
let mut first_missing = 0u64;
let mut present = Vec::new();
let mut j = 0u64;
while j - first_missing <= gap {
let Ok(built) = built_names(m, j) else {
break;
};
if self
.source_dims(
vds,
&built.source_file_name,
&built.source_dset_name,
access,
)
.is_some()
{
first_missing = j + 1;
present.push(j);
}
j += 1;
}
(first_missing, present)
}
fn source_dims(
&mut self,
vds: &str,
file_name: &str,
dset_name: &str,
access: &DatasetAccess,
) -> Option<Vec<u64>> {
let dset_name = dset_name.trim_start_matches('/');
if file_name == "." {
self.apply_dataset_access(dset_name, access).ok()?;
return self
.dataset_info_local(dset_name)
.map(|i| i.dataspace.dims.clone());
}
let reader = self.vds_source_file(vds, file_name)?;
reader.apply_dataset_access(dset_name, access).ok()?;
reader
.dataset_info(dset_name)
.map(|i| i.dataspace.dims.clone())
}
fn read_virtual_into(&mut self, name: &str, out: &mut [u8], depth: usize) -> IoResult<()> {
if depth >= MAX_VIRTUAL_DEPTH {
return Err(crate::io::IoError::InvalidState(format!(
"dataset {name:?}: virtual dataset mapping nests {MAX_VIRTUAL_DEPTH} levels \
deep, aborting (possible cyclic mapping)"
)));
}
let info = self
.dataset_info(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?;
let dims = info.dataspace.dims.clone();
let element_size = info.datatype.element_size() as u64;
let Some(mappings) = info.virtual_mappings.clone() else {
return Ok(());
};
let resolution = info.virtual_resolution.clone().unwrap_or_default();
let mappings = concrete_virtual_mappings(&mappings, &resolution)?;
let canonical = self.canonical_path(name);
let access = self.access_in_force(&canonical);
for mapping in &mappings {
let virtual_sel = mapping.virtual_selection.resolve(&dims).map_err(|e| {
crate::io::IoError::InvalidState(format!(
"dataset {name:?}: virtual mapping's virtual selection is not \
supported: {e}"
))
})?;
if virtual_sel.runs.is_empty() {
continue;
}
let source_name = mapping.source_dset_name.trim_start_matches('/');
if mapping.source_file_name == "." {
self.apply_dataset_access(source_name, &access)?;
let Some(src_dims) = self
.dataset_info(source_name)
.map(|i| i.dataspace.dims.clone())
else {
continue;
};
let source_sel = mapping.source_selection.resolve(&src_dims).map_err(|e| {
crate::io::IoError::InvalidState(format!(
"dataset {name:?}: virtual mapping's source selection is not \
supported: {e}"
))
})?;
copy_matched_selections(
|s, c, buf| {
self.read_slice_into_unconverted(
source_name,
s,
c,
buf,
depth + 1,
ReadDst::Fresh,
)
},
&source_sel,
&virtual_sel,
element_size,
out,
)?;
} else {
let Some(src_reader) = self.vds_source_file(&canonical, &mapping.source_file_name)
else {
continue;
};
src_reader.apply_dataset_access(source_name, &access)?;
let Some(src_dims) = src_reader
.dataset_info(source_name)
.map(|i| i.dataspace.dims.clone())
else {
continue;
};
let source_sel = mapping.source_selection.resolve(&src_dims).map_err(|e| {
crate::io::IoError::InvalidState(format!(
"dataset {name:?}: virtual mapping's source selection is not \
supported: {e}"
))
})?;
copy_matched_selections(
|s, c, buf| {
src_reader.read_slice_into_unconverted(
source_name,
s,
c,
buf,
depth + 1,
ReadDst::Fresh,
)
},
&source_sel,
&virtual_sel,
element_size,
out,
)?;
}
}
Ok(())
}
fn apply_post_filter_conversion(buffer: &mut [u8], datatype: &DatatypeMessage) -> IoResult<()> {
use crate::format::nbit_scaleoffset::{
apply_datatype_conversion, datatype_needs_bit_conversion,
};
if datatype_needs_bit_conversion(datatype) {
apply_datatype_conversion(buffer, datatype)?;
}
Ok(())
}
pub fn refresh(&mut self) -> IoResult<()> {
self.handle.refresh_read_source();
let sb_buf = self.handle.read_at_most(0, 256)?;
let sb = SuperblockV2V3::decode(&sb_buf)?;
let ctx = FormatContext {
sizeof_addr: sb.sizeof_offsets,
sizeof_size: sb.sizeof_lengths,
};
let (meta, ext) = Self::read_extension_and_meta(
&mut self.handle,
ctx,
self.meta.btree,
sb.superblock_extension_address,
)?;
let root_header = Self::read_object_header_full(
&mut self.handle,
&meta,
sb.root_group_object_header_address,
)?;
let catalog = Self::build_catalog(
&mut self.handle,
&meta,
Some(&root_header),
sb.root_group_object_header_address,
None,
)?;
let root_link_storage = describe_link_storage(Some(&root_header), &meta.ctx, None);
self._eof = sb.end_of_file_address;
self.meta = meta;
self.ext = ext;
self.object_paths = catalog.object_paths(sb.root_group_object_header_address);
self.datasets = DatasetTable::new(catalog.datasets);
self.unreadable = catalog.unreadable;
self.root_link_storage = root_link_storage;
self.group_attributes = catalog.group_attributes;
self.group_link_storage = catalog.group_link_storage;
self.group_paths = catalog.group_paths;
self.group_aliases = catalog.group_aliases;
self.links = catalog.links;
self.datatypes = catalog.datatypes;
self.resolve_virtual_extents()?;
Ok(())
}
fn decoded_chunk_index<F>(
&mut self,
pos: usize,
index_address: u64,
decode: F,
) -> IoResult<std::sync::Arc<DecodedChunkIndex>>
where
F: FnOnce(&mut Self) -> IoResult<DecodedChunkIndex>,
{
if let Some(hit) = self.datasets.chunk_index(pos, index_address) {
return Ok(std::sync::Arc::clone(hit));
}
let decoded = decode(self)?;
Ok(self.datasets.cache_chunk_index(pos, index_address, decoded))
}
fn read_chunked_v4(
&mut self,
name: &str,
chunk_dims: &[u64],
desc: ChunkIndexDesc<'_>,
req: ChunkReadRequest,
output: &mut [u8],
) -> IoResult<()> {
let ChunkReadRequest {
pipeline, target, ..
} = req;
let ChunkIndexDesc {
index_type,
index_address,
earray_params,
single_chunk_filter,
} = desc;
let pos = self.dataset_position(name)?;
let info = &self.datasets[pos];
let dims = info.dataspace.dims.clone();
let element_size = info.datatype.element_size() as u64;
match index_type {
data_layout::ChunkIndexType::SingleChunk => {
let total_size: u64 = saturating_byte_len(&dims, element_size);
let geo = ChunkOutputGeometry {
dims: &dims,
chunk_dims,
element_size,
};
if index_address == UNDEF_ADDR || total_size == 0 {
return place_chunk_jobs(
&self.handle,
Vec::new(),
&[],
req,
&geo,
None,
output,
);
}
let job = match (pipeline, single_chunk_filter) {
(Some(_), Some(scf)) => ChunkReadJob {
addr: index_address,
len: scf.nbytes as usize,
at_most: false,
mask: scf.filter_mask,
},
(Some(_), None) => ChunkReadJob {
addr: index_address,
len: total_size.saturating_mul(2) as usize,
at_most: true,
mask: 0,
},
(None, _) => ChunkReadJob {
addr: index_address,
len: total_size as usize,
at_most: false,
mask: 0,
},
};
let index = self.decoded_chunk_index(pos, index_address, |_| {
Ok(DecodedChunkIndex::new(
vec![(job.addr, job.len as u64, job.mask)],
vec![0u64; dims.len()],
))
})?;
place_chunk_jobs(
&self.handle,
vec![Some(job)],
&index.coords,
req,
&geo,
Some(&index.images),
output,
)
}
data_layout::ChunkIndexType::Implicit => {
self.read_chunked_implicit(name, chunk_dims, index_address, req, output)
}
data_layout::ChunkIndexType::FixedArray => {
self.read_chunked_fixed_array(name, chunk_dims, index_address, req, output)
}
data_layout::ChunkIndexType::BTreeV2 => {
self.read_chunked_btree_v2(name, chunk_dims, index_address, req, output)
}
data_layout::ChunkIndexType::ExtensibleArray => {
let params = earray_params.ok_or_else(|| {
crate::io::IoError::InvalidState("missing earray params".into())
})?;
if index_address == UNDEF_ADDR {
let geo = ChunkOutputGeometry {
dims: &dims,
chunk_dims,
element_size,
};
return place_chunk_jobs(
&self.handle,
Vec::new(),
&[],
req,
&geo,
None,
output,
);
}
let max_dims = self.datasets[pos].dataspace.max_dims.clone();
let rank = dims.len();
let index = self.decoded_chunk_index(pos, index_address, |reader| {
let grid =
crate::io::chunk_grid::index_grid(&dims, max_dims.as_deref(), chunk_dims)?;
let chunks_total: u64 = grid.iter().fold(1u64, |acc, &n| acc.saturating_mul(n));
let mut entries = reader.collect_ea_chunk_entries(
index_address,
params,
&dims,
max_dims.as_deref(),
chunk_dims,
element_size,
)?;
entries.truncate(std::cmp::min(chunks_total as usize, entries.len()));
let coords = crate::io::chunk_grid::coords_table(
&dims,
max_dims.as_deref(),
chunk_dims,
entries.len(),
)?;
Ok(DecodedChunkIndex::new(entries, coords))
})?;
let chunk_entries = &index.entries;
let slot_coords = &index.coords;
let chunk_coords = |i: usize| -> &[u64] { &slot_coords[i * rank..(i + 1) * rank] };
let jobs: Vec<Option<ChunkReadJob>> = if pipeline.is_some() {
let file_size = self.handle.file_size()?;
chunk_entries
.iter()
.enumerate()
.map(|(i, &(addr, nbytes, mask))| {
if addr == UNDEF_ADDR
|| nbytes == 0
|| addr >= file_size
|| nbytes > file_size
|| !target.overlaps(chunk_coords(i), chunk_dims)
{
None
} else {
Some(ChunkReadJob {
addr,
len: nbytes as usize,
at_most: false,
mask,
})
}
})
.collect()
} else {
chunk_entries
.iter()
.enumerate()
.map(|(i, &(addr, nbytes, _))| {
if addr == UNDEF_ADDR || !target.overlaps(chunk_coords(i), chunk_dims) {
None
} else {
Some(ChunkReadJob {
addr,
len: nbytes as usize,
at_most: true,
mask: 0,
})
}
})
.collect()
};
let geo = ChunkOutputGeometry {
dims: &dims,
chunk_dims,
element_size,
};
place_chunk_jobs(
&self.handle,
jobs,
slot_coords,
req,
&geo,
Some(&index.images),
output,
)
}
}
}
fn collect_fa_chunk_entries(
&mut self,
chunk_dims: &[u64],
ndims: usize,
element_size: u64,
index_address: u64,
) -> IoResult<Vec<(u64, u64, u32)>> {
use crate::format::chunk_index::fixed_array::*;
if index_address == UNDEF_ADDR {
return Ok(Vec::new());
}
let hdr_buf = self.handle.read_at_most(index_address, 256)?;
let fa_hdr = FixedArrayHeader::decode(&hdr_buf, &self.meta.ctx)?;
if fa_hdr.data_blk_addr == UNDEF_ADDR {
return Ok(Vec::new());
}
if chunk_dims.len() != ndims {
return Err(crate::io::IoError::InvalidState(format!(
"fixed-array dataset rank {} does not match chunk rank {}",
ndims,
chunk_dims.len()
)));
}
let is_filtered = fa_hdr.client_id == FA_CLIENT_FILT_CHUNK;
let sizeof_addr = self.meta.ctx.sizeof_addr as usize;
let chunk_size_len = if is_filtered {
(fa_hdr.element_size as usize)
.checked_sub(sizeof_addr + 4)
.ok_or_else(|| {
crate::io::IoError::InvalidState(
"fixed array filtered element_size too small".into(),
)
})?
} else {
0
};
if chunk_size_len > 8 {
return Err(crate::io::IoError::InvalidState(format!(
"fixed array filtered chunk-size width {chunk_size_len} exceeds 8 bytes"
)));
}
let chunk_bytes: u64 = saturating_byte_len(chunk_dims, element_size);
let elem_size = if is_filtered {
sizeof_addr + chunk_size_len + 4
} else {
sizeof_addr
};
let file_size = self.handle.file_size()?;
let num_elmts = usize::try_from(fa_hdr.num_elmts)
.ok()
.filter(|&n| {
(n as u64)
.checked_mul(elem_size as u64)
.is_some_and(|bytes| bytes <= file_size)
})
.ok_or_else(|| {
crate::io::IoError::InvalidState(format!(
"fixed array declares {} elements of {elem_size} bytes, more than the \
{file_size}-byte file holds",
fa_hdr.num_elmts
))
})?;
let mut chunk_entries: Vec<(u64, u64, u32)> = Vec::with_capacity(num_elmts);
if fa_hdr.is_paged() {
let npages = fa_hdr.npages();
let dblk_page_nelmts = fa_hdr.dblk_page_nelmts();
let prefix_len = 4 + 1 + 1 + sizeof_addr + (npages as usize).div_ceil(8) + 4;
let prefix_buf = self.handle.read_at_most(fa_hdr.data_blk_addr, prefix_len)?;
let prefix = FixedArrayPagedPrefix::decode(&prefix_buf, &self.meta.ctx, npages)?;
let page_stride = dblk_page_nelmts as usize * elem_size + 4;
let pages_base = fa_hdr.data_blk_addr + prefix.prefix_size as u64;
for p in 0..npages as usize {
let page_nelmts = if p + 1 == npages as usize {
let rem = fa_hdr.num_elmts % dblk_page_nelmts;
if rem == 0 {
dblk_page_nelmts
} else {
rem
}
} else {
dblk_page_nelmts
} as usize;
if !prefix.page_initialized(p) {
chunk_entries
.extend(std::iter::repeat_n((UNDEF_ADDR, 0u64, 0u32), page_nelmts));
continue;
}
let page_addr = pages_base + (p as u64) * page_stride as u64;
let page_size = page_nelmts * elem_size + 4;
let page_buf = self.handle.read_at_most(page_addr, page_size)?;
if is_filtered {
let elems = decode_filtered_page(
&page_buf,
&self.meta.ctx,
page_nelmts,
chunk_size_len,
)?;
for e in elems {
chunk_entries.push((e.address, e.chunk_size, e.filter_mask));
}
} else {
let addrs = decode_unfiltered_page(&page_buf, &self.meta.ctx, page_nelmts)?;
for addr in addrs {
chunk_entries.push((addr, chunk_bytes, 0));
}
}
}
} else {
let dblk_size = 4 + 1 + 1 + sizeof_addr + num_elmts * elem_size + 4;
let dblk_buf = self.handle.read_at_most(fa_hdr.data_blk_addr, dblk_size)?;
if is_filtered {
let fa_dblk = FixedArrayDataBlock::decode_filtered(
&dblk_buf,
&self.meta.ctx,
num_elmts,
chunk_size_len,
)?;
for e in &fa_dblk.filtered_elements {
chunk_entries.push((e.address, e.chunk_size, e.filter_mask));
}
} else {
let fa_dblk =
FixedArrayDataBlock::decode_unfiltered(&dblk_buf, &self.meta.ctx, num_elmts)?;
for &addr in &fa_dblk.elements {
chunk_entries.push((addr, chunk_bytes, 0));
}
}
}
Ok(chunk_entries)
}
fn read_chunked_fixed_array(
&mut self,
name: &str,
chunk_dims: &[u64],
index_address: u64,
req: ChunkReadRequest,
output: &mut [u8],
) -> IoResult<()> {
let ChunkReadRequest {
pipeline, target, ..
} = req;
let pos = self.dataset_position(name)?;
let info = &self.datasets[pos];
let dims = info.dataspace.dims.clone();
let element_size = info.datatype.element_size() as u64;
let max_dims = info.dataspace.max_dims.clone();
let ndims = dims.len();
let geo = ChunkOutputGeometry {
dims: &dims,
chunk_dims,
element_size,
};
let index = self.decoded_chunk_index(pos, index_address, |reader| {
let entries =
reader.collect_fa_chunk_entries(chunk_dims, ndims, element_size, index_address)?;
let coords = crate::io::chunk_grid::coords_table(
&dims,
max_dims.as_deref(),
chunk_dims,
entries.len(),
)?;
Ok(DecodedChunkIndex::new(entries, coords))
})?;
if index.entries.is_empty() {
return place_chunk_jobs(&self.handle, Vec::new(), &[], req, &geo, None, output);
}
let chunk_bytes: u64 = saturating_byte_len(chunk_dims, element_size);
let chunk_coords = |i: usize| -> &[u64] { &index.coords[i * ndims..(i + 1) * ndims] };
let jobs: Vec<Option<ChunkReadJob>> = index
.entries
.iter()
.enumerate()
.map(|(linear_idx, &(addr, comp_size, mask))| {
if addr == UNDEF_ADDR || !target.overlaps(chunk_coords(linear_idx), chunk_dims) {
None
} else if pipeline.is_some() {
let read_len = if comp_size > 0 {
comp_size as usize
} else {
chunk_bytes as usize * 2
};
Some(ChunkReadJob {
addr,
len: read_len,
at_most: true,
mask,
})
} else {
Some(ChunkReadJob {
addr,
len: chunk_bytes as usize,
at_most: false,
mask,
})
}
})
.collect();
place_chunk_jobs(
&self.handle,
jobs,
&index.coords,
req,
&geo,
Some(&index.images),
output,
)
}
fn read_chunked_implicit(
&mut self,
name: &str,
chunk_dims: &[u64],
index_address: u64,
req: ChunkReadRequest,
output: &mut [u8],
) -> IoResult<()> {
let target = req.target;
let pos = self.dataset_position(name)?;
let info = &self.datasets[pos];
let dims = info.dataspace.dims.clone();
let element_size = info.datatype.element_size() as u64;
let ndims = dims.len();
if index_address == UNDEF_ADDR {
let geo = ChunkOutputGeometry {
dims: &dims,
chunk_dims,
element_size,
};
return place_chunk_jobs(
&self.handle,
Vec::new(),
&[],
ChunkReadRequest {
pipeline: None,
..req
},
&geo,
None,
output,
);
}
if chunk_dims.len() != ndims {
return Err(crate::io::IoError::InvalidState(format!(
"implicit-index dataset rank {} does not match chunk rank {}",
ndims,
chunk_dims.len()
)));
}
let max_dims = self.datasets[pos].dataspace.max_dims.clone();
let chunk_bytes: u64 = saturating_byte_len(chunk_dims, element_size);
let index = self.decoded_chunk_index(pos, index_address, |_| {
let grid = crate::io::chunk_grid::index_grid(&dims, max_dims.as_deref(), chunk_dims)?;
let chunks_total: u64 = grid.iter().fold(1u64, |acc, &n| acc.saturating_mul(n));
let coords = crate::io::chunk_grid::coords_table(
&dims,
max_dims.as_deref(),
chunk_dims,
chunks_total as usize,
)?;
let entries = (0..chunks_total)
.map(|i| (index_address + i * chunk_bytes, chunk_bytes, 0))
.collect();
Ok(DecodedChunkIndex::new(entries, coords))
})?;
let slot_coords = &index.coords;
let jobs: Vec<Option<ChunkReadJob>> = index
.entries
.iter()
.enumerate()
.map(|(i, &(addr, len, _))| {
let coords = &slot_coords[i * ndims..(i + 1) * ndims];
if !target.overlaps(coords, chunk_dims) {
None
} else {
Some(ChunkReadJob {
addr,
len: len as usize,
at_most: false,
mask: 0,
})
}
})
.collect();
let geo = ChunkOutputGeometry {
dims: &dims,
chunk_dims,
element_size,
};
place_chunk_jobs(
&self.handle,
jobs,
slot_coords,
ChunkReadRequest {
pipeline: None,
..req
},
&geo,
Some(&index.images),
output,
)
}
fn collect_bt2_chunk_entries(
&mut self,
chunk_dims: &[u64],
ndims: usize,
element_size: u64,
index_address: u64,
) -> IoResult<Vec<Bt2ChunkEntry>> {
use crate::format::chunk_index::btree_v2::*;
if index_address == UNDEF_ADDR {
return Ok(Vec::new());
}
let hdr_buf = self.handle.read_at_most(index_address, 256)?;
let bt2_hdr = Bt2Header::decode(&hdr_buf, &self.meta.ctx)?;
if bt2_hdr.root_node_addr == UNDEF_ADDR || bt2_hdr.total_num_records == 0 {
return Ok(Vec::new());
}
let ctx = self.meta.ctx;
let record_bytes = collect_btree_v2_records(
&bt2_hdr,
&ctx,
&mut HandleBlockReader {
handle: &mut self.handle,
},
)?;
let total_records = if bt2_hdr.record_size > 0 {
record_bytes.len() / bt2_hdr.record_size as usize
} else {
0
};
let chunk_bytes: u64 = saturating_byte_len(chunk_dims, element_size);
let entries: Vec<Bt2ChunkEntry> = if bt2_hdr.record_type == BT2_TYPE_CHUNK_UNFILT {
Bt2ChunkIndex::decode_unfiltered_records(
&record_bytes,
total_records,
ndims,
&self.meta.ctx,
)?
.into_iter()
.map(|r| (r.chunk_address, chunk_bytes as usize, r.scaled_offsets, 0))
.collect()
} else {
Bt2ChunkIndex::decode_filtered_records(
&record_bytes,
total_records,
ndims,
bt2_hdr.record_size,
&self.meta.ctx,
)?
.into_iter()
.map(|r| {
(
r.chunk_address,
r.chunk_size as usize,
r.scaled_offsets,
r.filter_mask,
)
})
.collect()
};
Ok(entries)
}
fn read_chunked_btree_v2(
&mut self,
name: &str,
chunk_dims: &[u64],
index_address: u64,
req: ChunkReadRequest,
output: &mut [u8],
) -> IoResult<()> {
let target = req.target;
let pos = self.dataset_position(name)?;
let info = &self.datasets[pos];
let dims = info.dataspace.dims.clone();
let element_size = info.datatype.element_size() as u64;
let ndims = dims.len();
let geo = ChunkOutputGeometry {
dims: &dims,
chunk_dims,
element_size,
};
let index = self.decoded_chunk_index(pos, index_address, |reader| {
let records =
reader.collect_bt2_chunk_entries(chunk_dims, ndims, element_size, index_address)?;
let mut entries = Vec::with_capacity(records.len());
let mut coords = Vec::with_capacity(records.len() * ndims);
for (addr, read_size, scaled, mask) in records {
entries.push((addr, read_size as u64, mask));
coords.extend_from_slice(&scaled);
}
Ok(DecodedChunkIndex::new(entries, coords))
})?;
if index.entries.is_empty() {
return place_chunk_jobs(&self.handle, Vec::new(), &[], req, &geo, None, output);
}
let mut jobs: Vec<Option<ChunkReadJob>> = Vec::with_capacity(index.entries.len());
for (i, &(addr, read_size, mask)) in index.entries.iter().enumerate() {
let scaled = &index.coords[i * ndims..(i + 1) * ndims];
if addr == UNDEF_ADDR || read_size == 0 || !target.overlaps(scaled, chunk_dims) {
jobs.push(None);
} else {
jobs.push(Some(ChunkReadJob {
addr,
len: read_size as usize,
at_most: false,
mask,
}));
}
}
place_chunk_jobs(
&self.handle,
jobs,
&index.coords,
req,
&geo,
Some(&index.images),
output,
)
}
fn read_chunked_btree_v1(
&mut self,
name: &str,
chunk_dims: &[u64],
b_tree_address: u64,
req: ChunkReadRequest,
output: &mut [u8],
) -> IoResult<()> {
let ChunkReadRequest {
pipeline, target, ..
} = req;
let pos = self.dataset_position(name)?;
let info = &self.datasets[pos];
let dims = info.dataspace.dims.clone();
let element_size = info.datatype.element_size() as u64;
let ndims = dims.len();
if chunk_dims.len() != ndims {
return Err(crate::io::IoError::InvalidState(format!(
"B-tree-v1 dataset rank {} does not match chunk rank {}",
ndims,
chunk_dims.len()
)));
}
let total_size: u64 = saturating_byte_len(&dims, element_size);
if b_tree_address == UNDEF_ADDR || total_size == 0 {
let geo = ChunkOutputGeometry {
dims: &dims,
chunk_dims,
element_size,
};
return place_chunk_jobs(&self.handle, Vec::new(), &[], req, &geo, None, output);
}
let file_size = self.handle.file_size()?;
let index = self.decoded_chunk_index(pos, b_tree_address, |reader| {
let mut leaves: Vec<(Vec<u64>, u64, u32, u32)> = Vec::new();
reader.collect_btree_v1_chunks(b_tree_address, ndims, file_size, 0, &mut leaves)?;
let mut entries = Vec::with_capacity(leaves.len());
let mut coords = Vec::with_capacity(leaves.len() * ndims);
for (offsets, addr, chunk_size, mask) in leaves {
for (d, &cd) in chunk_dims.iter().enumerate().take(ndims) {
coords.push(offsets[d].checked_div(cd).unwrap_or(0));
}
entries.push((addr, chunk_size as u64, mask));
}
Ok(DecodedChunkIndex::new(entries, coords))
})?;
let chunk_bytes: u64 = saturating_byte_len(chunk_dims, element_size);
let mut jobs: Vec<Option<ChunkReadJob>> = Vec::with_capacity(index.entries.len());
for (i, &(addr, chunk_size, mask)) in index.entries.iter().enumerate() {
let skip = addr == UNDEF_ADDR
|| chunk_size == 0
|| addr >= file_size
|| chunk_size > file_size
|| !target.overlaps(&index.coords[i * ndims..(i + 1) * ndims], chunk_dims);
jobs.push(if skip {
None
} else {
Some(ChunkReadJob {
addr,
len: chunk_size as usize,
at_most: false,
mask,
})
});
}
let geo = ChunkOutputGeometry {
dims: &dims,
chunk_dims,
element_size,
};
place_chunk_jobs(
&self.handle,
jobs,
&index.coords,
req,
&geo,
Some(&index.images),
output,
)?;
if pipeline.is_none() {
for &(addr, chunk_size, _) in &index.entries {
if addr != UNDEF_ADDR && chunk_size != chunk_bytes && chunk_size != 0 {
return Err(crate::io::IoError::InvalidState(format!(
"chunk B-tree v1: unfiltered chunk size {} != expected {}",
chunk_size, chunk_bytes
)));
}
}
}
Ok(())
}
fn collect_btree_v1_chunks(
&mut self,
addr: u64,
rank: usize,
file_size: u64,
depth: u32,
out: &mut Vec<(Vec<u64>, u64, u32, u32)>,
) -> IoResult<()> {
if depth > 256 {
return Err(crate::io::IoError::InvalidState(
"chunk B-tree v1 exceeds maximum depth".into(),
));
}
if addr == UNDEF_ADDR || addr >= file_size {
return Ok(());
}
let sa = self.meta.ctx.sizeof_addr as usize;
let node_size = self.meta.btree.chunk_btree_node_size(sa, rank);
let buf = self.handle.read_at_most(addr, node_size)?;
let node = ChunkBTreeV1Node::decode(&buf, sa, rank, self.meta.btree.chunk_max_entries())?;
if node.level == 0 {
for (i, &child_addr) in node.children.iter().enumerate() {
let key = &node.keys[i];
out.push((
key.offsets[..rank].to_vec(),
child_addr,
key.chunk_size,
key.filter_mask,
));
}
} else {
let children: Vec<u64> = node.children.clone();
for child_addr in children {
if child_addr == UNDEF_ADDR || child_addr >= file_size {
continue;
}
self.collect_btree_v1_chunks(child_addr, rank, file_size, depth + 1, out)?;
}
}
Ok(())
}
pub fn read_vlen_strings(&mut self, name: &str) -> IoResult<Vec<String>> {
Ok(self
.read_vlen_objects(name)?
.into_iter()
.map(|bytes| String::from_utf8_lossy(&bytes).to_string())
.collect())
}
pub fn read_vlen_bytes(&mut self, name: &str) -> IoResult<Vec<Vec<u8>>> {
self.read_vlen_objects(name)
}
fn read_vlen_objects(&mut self, name: &str) -> IoResult<Vec<Vec<u8>>> {
if self.external_edge(name).is_some() {
let (owner, path, _) = self.external_owner(name, MAX_EXTERNAL_HOPS)?;
return owner.read_vlen_objects(&path);
}
let info = self
.dataset_info_local(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?;
let dims = info.dataspace.dims.clone();
let layout = info.layout.clone();
let external_files = info.external_files.clone();
let total_elements: u64 = dims.iter().fold(1u64, |acc, &d| acc.saturating_mul(d));
let raw = match &layout {
DataLayoutMessage::Contiguous { size, .. } if !external_files.is_empty() => {
let prefix = self.extfile_prefix_in_force(name);
let mut buf = vec![0u8; *size as usize];
read_external_file_bytes(&external_files, prefix.as_deref(), 0, &mut buf)?;
buf
}
DataLayoutMessage::Contiguous { address, size } => {
if *address == UNDEF_ADDR {
return Ok(vec![]);
}
self.handle.read_at(*address, *size as usize)?
}
DataLayoutMessage::Compact { data } => data.clone(),
_ => {
self.read_dataset_raw(name)?
}
};
let ref_size = vlen_reference_size(&self.meta.ctx);
let mut items: Vec<Vec<u8>> =
Vec::with_capacity((total_elements as usize).min(raw.len() / ref_size));
let mut heap_cache: std::collections::HashMap<
u64,
(GlobalHeapCollection, std::collections::HashMap<u16, usize>),
> = std::collections::HashMap::new();
for i in 0..total_elements as usize {
let offset = i * ref_size;
if offset + ref_size > raw.len() {
break;
}
let (_seq_len, collection_addr, obj_index) =
decode_vlen_reference(&raw[offset..], &self.meta.ctx)?;
if collection_addr == UNDEF_ADDR || collection_addr == 0 {
items.push(Vec::new());
continue;
}
if let std::collections::hash_map::Entry::Vacant(e) = heap_cache.entry(collection_addr)
{
let coll = self.read_heap_collection(collection_addr)?;
let lookup: std::collections::HashMap<u16, usize> = coll
.objects
.iter()
.enumerate()
.map(|(i, o)| (o.index, i))
.collect();
e.insert((coll, lookup));
}
let idx = u16::try_from(obj_index).map_err(|_| {
crate::io::IoError::InvalidState(format!(
"global heap object index {obj_index} does not fit the 16-bit on-disk field \
(element {i} of \"{name}\")"
))
})?;
let (coll, lookup) = &heap_cache[&collection_addr];
let &oi = lookup.get(&idx).ok_or_else(|| {
crate::io::IoError::InvalidState(format!(
"global heap object {idx} not found in the collection at address \
{collection_addr:#x} (element {i} of \"{name}\")"
))
})?;
items.push(coll.objects[oi].data.clone());
}
Ok(items)
}
fn collect_ea_chunk_entries(
&mut self,
index_address: u64,
params: &data_layout::EarrayParams,
dims: &[u64],
max_dims: Option<&[u64]>,
chunk_dims: &[u64],
element_size: u64,
) -> IoResult<Vec<(u64, u64, u32)>> {
use crate::format::chunk_index::extensible_array::{self as ea, *};
if index_address == UNDEF_ADDR {
return Ok(vec![]);
}
let hdr_buf = self.handle.read_at_most(index_address, 256)?;
let ea_hdr = ExtensibleArrayHeader::decode(&hdr_buf, &self.meta.ctx)?;
if ea_hdr.idx_blk_addr == UNDEF_ADDR {
return Ok(vec![]);
}
let chunks_dim0: usize = crate::io::chunk_grid::index_grid(dims, max_dims, chunk_dims)?
.iter()
.fold(1usize, |acc, &n| acc.saturating_mul(n as usize));
let geo = EaGeometry::new(
params.idx_blk_elmts,
params.data_blk_min_elmts,
params.sup_blk_min_data_ptrs,
params.max_nelmts_bits,
params.max_dblk_page_nelmts_bits,
)?;
let chunk_bytes = saturating_byte_len(chunk_dims, element_size);
let is_filtered = ea_hdr.class_id == ea::EA_CLS_FILT_CHUNK;
let chunk_size_len = if is_filtered {
ea_hdr.raw_elmt_size - self.meta.ctx.sizeof_addr - 4
} else {
0
};
let max_nelmts_bits = params.max_nelmts_bits;
let mut entries: Vec<(u64, u64, u32)> = Vec::new();
let (dblk_addrs, sblk_addrs): (Vec<u64>, Vec<u64>) = if is_filtered {
let buf = self.handle.read_at_most(ea_hdr.idx_blk_addr, 65536)?;
let fiblk = ea::FilteredIndexBlock::decode(
&buf,
&self.meta.ctx,
params.idx_blk_elmts as usize,
geo.ndblk_addrs,
geo.nsblk_addrs,
chunk_size_len,
)?;
for e in &fiblk.elements {
entries.push((e.addr, e.nbytes, e.filter_mask));
}
(fiblk.dblk_addrs, fiblk.sblk_addrs)
} else {
let buf = self.handle.read_at_most(ea_hdr.idx_blk_addr, 65536)?;
let iblk = ExtensibleArrayIndexBlock::decode(
&buf,
&self.meta.ctx,
params.idx_blk_elmts as usize,
geo.ndblk_addrs,
geo.nsblk_addrs,
)?;
for &addr in &iblk.elements {
entries.push((addr, chunk_bytes, 0));
}
(iblk.dblk_addrs, iblk.sblk_addrs)
};
let sa = self.meta.ctx.sizeof_addr as usize;
let raw_elmt_size = if is_filtered {
ea::FilteredChunkEntry::raw_size(self.meta.ctx.sizeof_addr, chunk_size_len) as usize
} else {
sa
};
'outer: for (u, s) in geo.sblk.iter().enumerate() {
if entries.len() >= chunks_dim0 {
break;
}
let dblk_nelmts = s.dblk_nelmts as usize;
let paged = geo.is_sblk_paged(u);
let (this_dblk_addrs, page_init): (Vec<u64>, Vec<u8>) = if u < geo.iblock_nsblks {
let start = s.start_dblk as usize;
(
(0..s.ndblks as usize)
.map(|d| dblk_addrs.get(start + d).copied().unwrap_or(UNDEF_ADDR))
.collect(),
Vec::new(),
)
} else {
let sblk_addr = sblk_addrs
.get(u - geo.iblock_nsblks)
.copied()
.unwrap_or(UNDEF_ADDR);
if sblk_addr == UNDEF_ADDR {
(vec![UNDEF_ADDR; s.ndblks as usize], Vec::new())
} else {
let page_init_total = if paged {
s.ndblks as usize * geo.dblk_page_init_size(u)
} else {
0
};
let sblk_size =
4 + 1 + 1 + sa + 8 + page_init_total + s.ndblks as usize * sa + 4;
let buf = self.handle.read_at_most(sblk_addr, sblk_size)?;
let sb = ExtensibleArraySuperBlock::decode(
&buf,
&self.meta.ctx,
max_nelmts_bits,
s.ndblks as usize,
page_init_total,
)?;
(sb.dblk_addrs, sb.page_init)
}
};
let npages = geo.npages(u) as usize;
let page_size = geo.dblk_page_size(raw_elmt_size);
let prefix = geo.dblk_prefix_size(self.meta.ctx.sizeof_addr, max_nelmts_bits);
for (d, &dblk_addr) in this_dblk_addrs.iter().enumerate() {
if dblk_addr == UNDEF_ADDR {
entries.extend(std::iter::repeat_n((UNDEF_ADDR, 0, 0), dblk_nelmts));
} else if paged {
for p in 0..npages {
let bit = d * npages + p;
let initialized = page_init[bit / 8] & (0x80u8 >> (bit % 8)) != 0;
if !initialized {
entries.extend(std::iter::repeat_n(
(UNDEF_ADDR, 0, 0),
geo.dblk_page_nelmts as usize,
));
continue;
}
let page_addr = dblk_addr + prefix as u64 + (p as u64) * page_size as u64;
let page = self.handle.read_at(page_addr, page_size)?;
for k in 0..geo.dblk_page_nelmts as usize {
let off = k * raw_elmt_size;
if is_filtered {
let e = ea::FilteredChunkEntry::decode(
&page[off..],
sa,
chunk_size_len as usize,
);
entries.push((e.addr, e.nbytes, e.filter_mask));
} else {
entries.push((read_addr(&page[off..], sa), chunk_bytes, 0));
}
}
}
} else if is_filtered {
let dblk_size = prefix + dblk_nelmts * raw_elmt_size;
let buf = self.handle.read_at_most(dblk_addr, dblk_size)?;
let dblk = ea::FilteredDataBlock::decode(
&buf,
&self.meta.ctx,
max_nelmts_bits,
dblk_nelmts,
chunk_size_len,
)?;
for e in &dblk.elements {
entries.push((e.addr, e.nbytes, e.filter_mask));
}
} else {
let dblk_size = prefix + dblk_nelmts * raw_elmt_size;
let buf = self.handle.read_at_most(dblk_addr, dblk_size)?;
let dblk = ExtensibleArrayDataBlock::decode(
&buf,
&self.meta.ctx,
max_nelmts_bits,
dblk_nelmts,
)?;
for &addr in &dblk.elements {
entries.push((addr, chunk_bytes, 0));
}
}
if entries.len() >= chunks_dim0 {
break 'outer;
}
}
}
Ok(entries)
}
pub fn read_slice(&mut self, name: &str, starts: &[u64], counts: &[u64]) -> IoResult<Vec<u8>> {
if self.external_edge(name).is_some() {
let (owner, path, _) = self.external_owner(name, MAX_EXTERNAL_HOPS)?;
return owner.read_slice(&path, starts, counts);
}
let (datatype, out_bytes) = self.slice_size_and_datatype(name, starts, counts)?;
read_image_into_new::<u8, _, _>(out_bytes as usize, |image| {
self.read_slice_into_unconverted(name, starts, counts, image, 0, ReadDst::Fresh)?;
Self::apply_post_filter_conversion(image, &datatype)
})
}
pub fn read_slice_into(
&mut self,
name: &str,
starts: &[u64],
counts: &[u64],
out: &mut [u8],
) -> IoResult<()> {
self.read_slice_into_dst(name, starts, counts, out, ReadDst::Reused)
}
pub(crate) fn read_slice_into_dst(
&mut self,
name: &str,
starts: &[u64],
counts: &[u64],
out: &mut [u8],
dst: ReadDst,
) -> IoResult<()> {
if self.external_edge(name).is_some() {
let (owner, path, _) = self.external_owner(name, MAX_EXTERNAL_HOPS)?;
return owner.read_slice_into_dst(&path, starts, counts, out, dst);
}
let (datatype, out_bytes) = self.slice_size_and_datatype(name, starts, counts)?;
if out.len() as u64 != out_bytes {
return Err(crate::io::IoError::InvalidState(format!(
"read_slice_into: buffer is {} bytes but selection needs {}",
out.len(),
out_bytes
)));
}
self.read_slice_into_unconverted(name, starts, counts, out, 0, dst)?;
Self::apply_post_filter_conversion(out, &datatype)?;
Ok(())
}
fn slice_size_and_datatype(
&self,
name: &str,
starts: &[u64],
counts: &[u64],
) -> IoResult<(DatatypeMessage, u64)> {
let info = self
.dataset_info_local(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?;
check_hyperslab(&info.dataspace.dims, starts, counts)?;
let out_bytes = saturating_byte_len(counts, info.datatype.element_size() as u64);
Ok((info.datatype.clone(), out_bytes))
}
fn read_slice_into_unconverted(
&mut self,
name: &str,
starts: &[u64],
counts: &[u64],
out: &mut [u8],
depth: usize,
dst: ReadDst,
) -> IoResult<()> {
let info = self
.dataset_info_local(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?;
let dims = info.dataspace.dims.clone();
let element_size = info.datatype.element_size() as u64;
let layout = info.layout.clone();
let pipeline = info.filter_pipeline.clone();
let fill_value = info.fill_value.clone();
let external_files = info.external_files.clone();
let ndims = dims.len();
check_hyperslab(&dims, starts, counts)?;
if ndims == 0 {
return Err(crate::io::IoError::InvalidState(
"read_slice does not support scalar datasets; use read_dataset_raw".into(),
));
}
match &layout {
DataLayoutMessage::Contiguous { .. } if !external_files.is_empty() => {
let prefix = self.extfile_prefix_in_force(name);
for_each_contiguous_run(
&dims,
starts,
counts,
element_size,
|src_off, out_off, len| {
read_external_file_bytes(
&external_files,
prefix.as_deref(),
src_off,
&mut out[out_off..out_off + len],
)
},
)?;
}
DataLayoutMessage::Contiguous { address, .. } => {
if *address == UNDEF_ADDR {
fill_tiled_into(out, fill_value.as_deref());
} else {
let base = *address;
for_each_contiguous_run(
&dims,
starts,
counts,
element_size,
|src_off, out_off, len| {
let at = base.checked_add(src_off).ok_or_else(|| {
crate::io::IoError::InvalidState(format!(
"dataset '{name}' claims raw data at {base}, which \
overflows {src_off} bytes into the selection"
))
})?;
self.handle
.read_exact_at_into(at, &mut out[out_off..out_off + len], dst)
.map_err(Into::into)
},
)?;
}
}
DataLayoutMessage::Compact { data } => {
for_each_contiguous_run(
&dims,
starts,
counts,
element_size,
|src_off, out_off, len| {
let src = src_off as usize;
out[out_off..out_off + len].copy_from_slice(&data[src..src + len]);
Ok(())
},
)?;
}
DataLayoutMessage::ChunkedV3 {
chunk_dims,
b_tree_address,
} => {
let real_chunk_dims = &chunk_dims[..chunk_dims.len() - 1];
self.read_chunked_btree_v1(
name,
real_chunk_dims,
*b_tree_address,
ChunkReadRequest {
pipeline: pipeline.as_ref(),
target: ChunkTarget::Slice { starts, counts },
fill_value: fill_value.as_deref(),
dst,
},
out,
)?;
}
DataLayoutMessage::ChunkedV4 {
chunk_dims,
index_address,
index_type,
earray_params,
single_chunk_filter,
..
} => {
let real_chunk_dims = &chunk_dims[..chunk_dims.len() - 1];
self.read_chunked_v4(
name,
real_chunk_dims,
ChunkIndexDesc {
index_type: *index_type,
index_address: *index_address,
earray_params: earray_params.as_ref(),
single_chunk_filter: *single_chunk_filter,
},
ChunkReadRequest {
pipeline: pipeline.as_ref(),
target: ChunkTarget::Slice { starts, counts },
fill_value: fill_value.as_deref(),
dst,
},
out,
)?;
}
DataLayoutMessage::Virtual { .. } => {
let total = saturating_byte_len(&dims, element_size) as usize;
let mut full = alloc_tiled_fill(total, fill_value.as_deref())?;
self.read_virtual_into(name, &mut full, depth)?;
for_each_contiguous_run(
&dims,
starts,
counts,
element_size,
|src_off, out_off, len| {
let src_off = src_off as usize;
out[out_off..out_off + len].copy_from_slice(&full[src_off..src_off + len]);
Ok(())
},
)?;
}
}
Ok(())
}
pub fn read_hyperslab_into(
&mut self,
name: &str,
start: &[u64],
stride: &[u64],
count: &[u64],
block: &[u64],
out: &mut [u8],
) -> IoResult<()> {
if self.external_edge(name).is_some() {
let (owner, path, _) = self.external_owner(name, MAX_EXTERNAL_HOPS)?;
return owner.read_hyperslab_into(&path, start, stride, count, block, out);
}
let info = self
.dataset_info_local(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?;
let dims = info.dataspace.dims.clone();
let datatype = info.datatype.clone();
let element_size = datatype.element_size() as u64;
let rank = dims.len();
if start.len() != rank || stride.len() != rank || count.len() != rank || block.len() != rank
{
return Err(crate::io::IoError::InvalidState(
"start/stride/count/block length must match dataset rank".into(),
));
}
if stride.contains(&0) {
return Err(crate::io::IoError::InvalidState(
"hyperslab stride must be nonzero in every dimension".into(),
));
}
let out_bytes = saturating_byte_len(
&(0..rank)
.map(|d| count[d].saturating_mul(block[d]))
.collect::<Vec<u64>>(),
element_size,
);
if out.len() as u64 != out_bytes {
return Err(crate::io::IoError::InvalidState(format!(
"read_hyperslab_into: buffer is {} bytes but selection needs {out_bytes}",
out.len(),
)));
}
let src_sel = Selection::Hyperslab {
rank,
form: Hyperslab::Regular(RegularHyperslab {
start: start.to_vec(),
stride: stride.to_vec(),
count: count.to_vec(),
block: block.to_vec(),
}),
};
let out_dims: Vec<u64> = (0..rank)
.map(|d| count[d].saturating_mul(block[d]))
.collect();
let dst_sel = Selection::Hyperslab {
rank,
form: Hyperslab::Regular(RegularHyperslab {
start: vec![0u64; rank],
stride: block.to_vec(),
count: count.to_vec(),
block: block.to_vec(),
}),
};
fill_tiled_into(out, None);
copy_matched_selections(
|bstart, bcount, buf| {
self.read_slice_into_unconverted(name, bstart, bcount, buf, 0, ReadDst::Fresh)
},
&src_sel.resolve(&dims)?,
&dst_sel.resolve(&out_dims)?,
element_size,
out,
)?;
Self::apply_post_filter_conversion(out, &datatype)
}
pub fn read_points_into(
&mut self,
name: &str,
points: &[Vec<u64>],
out: &mut [u8],
dst: ReadDst,
) -> IoResult<()> {
if self.external_edge(name).is_some() {
let (owner, path, _) = self.external_owner(name, MAX_EXTERNAL_HOPS)?;
return owner.read_points_into(&path, points, out, dst);
}
let info = self
.dataset_info_local(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?;
let dims = info.dataspace.dims.clone();
let datatype = info.datatype.clone();
let element_size = datatype.element_size() as u64;
let rank = dims.len();
for p in points {
if p.len() != rank {
return Err(crate::io::IoError::InvalidState(format!(
"point coordinate has {} entries but the dataset has {} dimensions",
p.len(),
rank
)));
}
}
let sel = Selection::Points(PointSelection {
rank,
points: points.to_vec(),
});
let boxes = sel.to_boxes(&dims)?;
let es = element_size as usize;
if out.len() != points.len() * es {
return Err(crate::io::IoError::InvalidState(format!(
"read_points_into: buffer is {} bytes but {} points need {}",
out.len(),
points.len(),
points.len() * es
)));
}
for (i, (bstart, bcount)) in boxes.iter().enumerate() {
self.read_slice_into_unconverted(
name,
bstart,
bcount,
&mut out[i * es..(i + 1) * es],
0,
dst,
)?;
}
Self::apply_post_filter_conversion(out, &datatype)
}
pub fn read_chunk_raw_at(
&mut self,
name: &str,
chunk_coords: &[u64],
) -> IoResult<(Vec<u8>, u32)> {
if self.external_edge(name).is_some() {
let (owner, path, _) = self.external_owner(name, MAX_EXTERNAL_HOPS)?;
return owner.read_chunk_raw_at(&path, chunk_coords);
}
let info = self
.dataset_info_local(name)
.ok_or_else(|| crate::io::IoError::NotFound(name.to_string()))?;
let dims = info.dataspace.dims.clone();
let max_dims = info.dataspace.max_dims.clone();
let element_size = info.datatype.element_size() as u64;
let layout = info.layout.clone();
let ndims = dims.len();
if chunk_coords.len() != ndims {
return Err(crate::io::IoError::InvalidState(format!(
"chunk_coords has {} entries but the dataset has {} dimensions",
chunk_coords.len(),
ndims
)));
}
let not_written = || {
crate::io::IoError::InvalidState(format!(
"chunk at coordinates {chunk_coords:?} has not been written"
))
};
match &layout {
DataLayoutMessage::ChunkedV3 {
chunk_dims,
b_tree_address,
} => {
let real_chunk_dims = &chunk_dims[..chunk_dims.len() - 1];
if *b_tree_address == UNDEF_ADDR {
return Err(not_written());
}
let file_size = self.handle.file_size()?;
let mut entries = Vec::new();
self.collect_btree_v1_chunks(*b_tree_address, ndims, file_size, 0, &mut entries)?;
for (offsets, addr, chunk_size, mask) in &entries {
if *addr == UNDEF_ADDR {
continue;
}
let mut scaled = Vec::with_capacity(ndims);
for d in 0..ndims {
scaled.push(offsets[d].checked_div(real_chunk_dims[d]).unwrap_or(0));
}
if scaled.as_slice() == chunk_coords {
return Ok((self.handle.read_at(*addr, *chunk_size as usize)?, *mask));
}
}
Err(not_written())
}
DataLayoutMessage::ChunkedV4 {
chunk_dims,
index_type,
index_address,
earray_params,
single_chunk_filter,
..
} => {
let real_chunk_dims = &chunk_dims[..chunk_dims.len() - 1];
match index_type {
data_layout::ChunkIndexType::SingleChunk => {
if chunk_coords.iter().any(|&c| c != 0) {
return Err(crate::io::IoError::InvalidState(format!(
"chunk coordinates {chunk_coords:?} are outside the chunk \
grid (0..1): this dataset has a single-chunk index"
)));
}
if *index_address == UNDEF_ADDR {
return Err(not_written());
}
match single_chunk_filter {
Some(scf) => Ok((
self.handle.read_at(*index_address, scf.nbytes as usize)?,
scf.filter_mask,
)),
None => {
let total = saturating_byte_len(&dims, element_size);
Ok((self.handle.read_at(*index_address, total as usize)?, 0))
}
}
}
data_layout::ChunkIndexType::Implicit => {
if *index_address == UNDEF_ADDR {
return Err(not_written());
}
let linear = crate::io::chunk_grid::linear_index(
&dims,
max_dims.as_deref(),
real_chunk_dims,
chunk_coords,
)?;
let chunk_bytes = saturating_byte_len(real_chunk_dims, element_size);
let addr = index_address.saturating_add(linear.saturating_mul(chunk_bytes));
Ok((self.handle.read_at(addr, chunk_bytes as usize)?, 0))
}
data_layout::ChunkIndexType::FixedArray => {
let linear = crate::io::chunk_grid::linear_index(
&dims,
max_dims.as_deref(),
real_chunk_dims,
chunk_coords,
)?;
let entries = self.collect_fa_chunk_entries(
real_chunk_dims,
ndims,
element_size,
*index_address,
)?;
match entries.get(linear as usize) {
Some(&(addr, size, mask)) if addr != UNDEF_ADDR => {
Ok((self.handle.read_at(addr, size as usize)?, mask))
}
_ => Err(not_written()),
}
}
data_layout::ChunkIndexType::ExtensibleArray => {
let params = earray_params.as_ref().ok_or_else(|| {
crate::io::IoError::InvalidState("missing earray params".into())
})?;
let linear = crate::io::chunk_grid::linear_index(
&dims,
max_dims.as_deref(),
real_chunk_dims,
chunk_coords,
)?;
let entries = self.collect_ea_chunk_entries(
*index_address,
params,
&dims,
max_dims.as_deref(),
real_chunk_dims,
element_size,
)?;
match entries.get(linear as usize) {
Some(&(addr, size, mask)) if addr != UNDEF_ADDR => {
Ok((self.handle.read_at(addr, size as usize)?, mask))
}
_ => Err(not_written()),
}
}
data_layout::ChunkIndexType::BTreeV2 => {
let entries = self.collect_bt2_chunk_entries(
real_chunk_dims,
ndims,
element_size,
*index_address,
)?;
match entries
.iter()
.find(|(_, _, scaled, _)| scaled.as_slice() == chunk_coords)
{
Some(&(addr, size, _, mask)) if addr != UNDEF_ADDR => {
Ok((self.handle.read_at(addr, size)?, mask))
}
_ => Err(not_written()),
}
}
}
}
_ => Err(crate::io::IoError::InvalidState(
"read_chunk_raw_at is only for chunked datasets".into(),
)),
}
}
}
pub(crate) struct HandleBlockReader<'a> {
pub(crate) handle: &'a FileHandle,
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct ObjectAttributes {
entries: Vec<AttributeEntry>,
incomplete: Option<String>,
creation_order: CreationOrder,
storage: AttributeStorage,
}
impl ObjectAttributes {
fn push(&mut self, entry: AttributeEntry) {
self.entries.push(entry);
}
fn mark_incomplete(&mut self, reason: String) {
if self.incomplete.is_none() {
self.incomplete = Some(reason);
}
}
pub fn unreadable_reason(&self) -> Option<&str> {
self.incomplete.as_deref()
}
pub fn creation_order(&self) -> CreationOrder {
self.creation_order
}
pub fn storage(&self) -> AttributeStorage {
self.storage
}
pub fn header_count(&self, owner: &str) -> IoResult<u64> {
Ok(self.complete(owner)?.len() as u64)
}
pub(crate) fn into_complete(self, owner: &str) -> IoResult<Vec<AttributeEntry>> {
match self.incomplete {
Some(reason) => Err(incomplete_error(owner, &reason)),
None => Ok(self.entries),
}
}
fn complete(&self, owner: &str) -> IoResult<&[AttributeEntry]> {
match &self.incomplete {
Some(reason) => Err(incomplete_error(owner, reason)),
None => Ok(&self.entries),
}
}
pub(crate) fn ordered_names(&self, owner: &str) -> IoResult<Vec<String>> {
let mut ordered: Vec<&AttributeEntry> = self.complete(owner)?.iter().collect();
if !ordered.is_empty() && ordered.iter().all(|e| e.creation_index().is_some()) {
ordered.sort_by_key(|e| e.creation_index());
} else {
ordered.sort_by(|a, b| a.name().cmp(b.name()));
}
Ok(ordered.into_iter().map(|e| e.name().to_string()).collect())
}
}
fn incomplete_error(owner: &str, reason: &str) -> crate::io::IoError {
crate::io::IoError::Unsupported(format!(
"attributes of '{owner}' cannot be read whole: {reason}"
))
}
pub(crate) fn collect_object_attributes(
handle: &mut FileHandle,
ctx: &FormatContext,
header: &ObjectHeader,
) -> ObjectAttributes {
let mut attrs = ObjectAttributes {
creation_order: header.attribute_creation_order(),
..ObjectAttributes::default()
};
let tracked = header.has_creation_order();
for msg in &header.messages {
match msg.msg_type {
MSG_ATTRIBUTE => match AttributeEntry::parse(&msg.data, ctx) {
Ok(entry) => {
attrs.push(entry.with_creation_index(tracked.then_some(msg.creation_index)))
}
Err(e) => attrs.mark_incomplete(format!("an attribute message is unreadable: {e}")),
},
MSG_ATTR_INFO => match AttributeInfoMessage::decode(&msg.data, ctx) {
Ok((info, _)) => {
attrs.storage = if info.is_dense() {
AttributeStorage::Dense
} else {
AttributeStorage::Compact
};
let mut br = HandleBlockReader { handle };
match crate::format::dense_attr::read_dense_attributes(&info, ctx, &mut br) {
Ok(dense) => attrs.entries.extend(dense),
Err(e) => attrs
.mark_incomplete(format!("dense attribute storage is unreadable: {e}")),
}
}
Err(e) => {
attrs.mark_incomplete(format!("the attribute info message is unreadable: {e}"))
}
},
_ => {}
}
}
attrs
}
impl BlockReader for HandleBlockReader<'_> {
fn read_block(&mut self, offset: u64, len: usize) -> crate::format::FormatResult<Vec<u8>> {
self.handle.read_at_most(offset, len).map_err(|e| {
crate::format::FormatError::InvalidData(format!(
"metadata block read failed at {:#x}: {}",
offset, e
))
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::Write;
fn temp_path(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_reader_test_{}_{}_{}.h5",
name,
std::process::id(),
n
))
}
fn handle_over(name: &str, bytes: &[u8]) -> (std::path::PathBuf, FileHandle) {
let path = temp_path(name);
std::fs::File::create(&path)
.unwrap()
.write_all(bytes)
.unwrap();
let handle =
FileHandle::open_read_with_locking(&path, crate::io::locking::FileLocking::Disabled)
.unwrap();
(path, handle)
}
mod chunk_index_cache {
use super::*;
fn chunked_file(name: &str) -> (std::path::PathBuf, Hdf5Reader) {
let path = temp_path(name);
{
let file = crate::H5File::create(&path).unwrap();
let ds = file
.new_dataset::<u8>()
.shape([8usize])
.chunk(&[4])
.create("d")
.unwrap();
ds.write_raw(&(0u8..8).collect::<Vec<u8>>()).unwrap();
file.close().unwrap();
}
let reader = Hdf5Reader::open(&path).unwrap();
(path, reader)
}
fn index_address(reader: &Hdf5Reader, pos: usize) -> u64 {
match &reader.datasets[pos].layout {
DataLayoutMessage::ChunkedV4 { index_address, .. } => *index_address,
_ => panic!("expected a version-4 chunked layout"),
}
}
#[test]
fn a_read_leaves_its_decoded_index_for_the_next_read() {
let (path, mut r) = chunked_file("index_cache_hit");
let pos = r.dataset_position("d").unwrap();
let addr = index_address(&r, pos);
assert!(
r.datasets.chunk_index(pos, addr).is_none(),
"a reader that has read nothing holds a decoded index"
);
assert_eq!(r.read_slice("d", &[0], &[4]).unwrap(), vec![0u8, 1, 2, 3]);
let cached = r
.datasets
.chunk_index(pos, addr)
.expect("the read did not keep the index it decoded");
assert_eq!(cached.entries.len(), 2, "one entry per chunk");
assert_eq!(cached.coords, vec![0, 1], "slot 0 and slot 1 of the grid");
assert_eq!(r.read_slice("d", &[4], &[4]).unwrap(), vec![4u8, 5, 6, 7]);
let _ = std::fs::remove_file(&path);
}
#[test]
fn an_index_is_never_served_for_another_address() {
let (path, mut r) = chunked_file("index_cache_addr");
let pos = r.dataset_position("d").unwrap();
let addr = index_address(&r, pos);
r.read_slice("d", &[0], &[4]).unwrap();
assert!(
r.datasets.chunk_index(pos, addr.wrapping_add(1)).is_none(),
"an index decoded from {addr:#x} answered for another address"
);
let _ = std::fs::remove_file(&path);
}
#[test]
fn changing_an_entry_drops_the_index_decoded_against_it() {
let (path, mut r) = chunked_file("index_cache_entry_mut");
let pos = r.dataset_position("d").unwrap();
let addr = index_address(&r, pos);
r.read_slice("d", &[0], &[4]).unwrap();
assert!(r.datasets.chunk_index(pos, addr).is_some());
let _ = r.datasets.entry_mut(pos);
assert!(
r.datasets.chunk_index(pos, addr).is_none(),
"a mutable entry left its decoded index behind"
);
assert_eq!(r.read_slice("d", &[4], &[4]).unwrap(), vec![4u8, 5, 6, 7]);
let _ = std::fs::remove_file(&path);
}
}
mod chunk_image_cache {
use super::*;
#[cfg(feature = "deflate")]
fn dataset(name: &str, n: usize, chunk: usize, deflate: bool) -> std::path::PathBuf {
let path = temp_path(name);
let data: Vec<f64> = (0..n).map(|i| (i % 251) as f64).collect();
let file = crate::H5File::create(&path).unwrap();
let mut b = file.new_dataset::<f64>().shape([n]).chunk(&[chunk]);
if deflate {
b = b.deflate(6);
}
b.create("d").unwrap().write_raw(&data).unwrap();
file.close().unwrap();
path
}
fn dataset_2d(
name: &str,
rows: usize,
cols: usize,
chunk: [usize; 2],
deflate: bool,
) -> std::path::PathBuf {
let path = temp_path(name);
let data: Vec<f64> = (0..rows * cols).map(|i| (i % 251) as f64).collect();
let file = crate::H5File::create(&path).unwrap();
let mut b = file
.new_dataset::<f64>()
.shape([rows, cols])
.chunk(&chunk[..]);
if deflate {
b = b.deflate(6);
}
b.create("d").unwrap().write_raw(&data).unwrap();
file.close().unwrap();
path
}
fn held(reader: &Hdf5Reader) -> usize {
let pos = reader.dataset_position("d").unwrap();
let DataLayoutMessage::ChunkedV4 { index_address, .. } = &reader.datasets[pos].layout
else {
panic!("expected a version-4 chunked layout");
};
match reader.datasets.chunk_index(pos, *index_address) {
Some(index) => index.images.held(),
None => 0,
}
}
#[test]
#[cfg(feature = "deflate")]
fn a_partial_read_keeps_the_chunk_it_did_not_finish() {
let path = dataset("images_partial", 4096, 1024, true);
let mut r = Hdf5Reader::open(&path).unwrap();
let whole = {
let mut fresh = Hdf5Reader::open(&path).unwrap();
fresh.read_dataset_raw("d").unwrap()
};
for q in 0..4u64 {
let got = r.read_slice("d", &[q * 256], &[256]).unwrap();
let from = (q as usize) * 256 * 8;
assert_eq!(got, whole[from..from + 256 * 8], "quarter {q}");
assert_eq!(held(&r), 1, "quarter {q} left the wrong image count");
}
let _ = std::fs::remove_file(&path);
}
#[test]
#[cfg(feature = "deflate")]
fn a_whole_dataset_read_keeps_nothing() {
let path = dataset_2d("images_full", 256, 256, [64, 64], true);
let mut r = Hdf5Reader::open(&path).unwrap();
r.read_dataset_raw("d").unwrap();
assert_eq!(held(&r), 0, "a full read kept a chunk image");
let _ = std::fs::remove_file(&path);
}
#[test]
#[cfg(feature = "deflate")]
fn a_whole_read_of_a_ragged_extent_keeps_nothing() {
let path = dataset("images_ragged", 3000, 1024, true);
let mut r = Hdf5Reader::open(&path).unwrap();
r.read_dataset_raw("d").unwrap();
assert_eq!(held(&r), 0, "the ragged edge chunk was kept");
let _ = std::fs::remove_file(&path);
}
#[test]
fn an_unfiltered_partial_read_keeps_nothing() {
let path = dataset_2d("images_plain", 256, 256, [64, 64], false);
let mut r = Hdf5Reader::open(&path).unwrap();
let column = r.read_slice("d", &[0, 0], &[256, 1]).unwrap();
assert_eq!(column.len(), 256 * 8);
assert_eq!(held(&r), 0, "an unfiltered chunk reached the cache");
let _ = std::fs::remove_file(&path);
}
#[test]
#[cfg(feature = "deflate")]
fn a_filtered_column_walk_reuses_each_chunk() {
let path = dataset_2d("images_column", 256, 256, [64, 64], true);
let mut warm = Hdf5Reader::open(&path).unwrap();
for c in 0..4u64 {
let mut cold = Hdf5Reader::open(&path).unwrap();
assert_eq!(
warm.read_slice("d", &[0, c], &[256, 1]).unwrap(),
cold.read_slice("d", &[0, c], &[256, 1]).unwrap(),
"column {c} differed once served from the cache"
);
}
assert_eq!(held(&warm), 4, "the column walk re-inflated its chunks");
let _ = std::fs::remove_file(&path);
}
#[test]
#[cfg(feature = "deflate")]
fn a_chunk_over_the_budget_is_never_kept() {
let over = CHUNK_CACHE_BYTES / 8 + 1;
let path = dataset("images_oversize", over * 2, over, true);
let mut r = Hdf5Reader::open(&path).unwrap();
r.read_slice("d", &[0], &[1024]).unwrap();
assert_eq!(held(&r), 0, "a chunk over the budget was kept");
let _ = std::fs::remove_file(&path);
}
#[test]
fn an_image_over_the_budget_never_evicts_the_ones_under_it() {
let cache = ChunkImageCache::default();
let key = |addr| ChunkImageKey {
addr,
len: 1,
mask: 0,
};
for i in 0..3 {
cache.keep(key(i), vec![0u8; CHUNK_CACHE_BYTES / 4]);
}
assert_eq!(cache.held(), 3, "three quarter-budget images did not fit");
cache.keep(key(9), vec![0u8; CHUNK_CACHE_BYTES + 1]);
assert_eq!(
cache.held(),
3,
"an image that cannot fit evicted the ones that did"
);
cache.keep(key(8), vec![0u8; CHUNK_CACHE_BYTES]);
assert_eq!(
cache.held(),
1,
"an image of exactly the budget was refused"
);
}
#[test]
#[cfg(feature = "deflate")]
fn the_cache_holds_no_more_than_its_byte_budget() {
let chunk = 32 * 1024; let path = dataset("images_budget", chunk * 8, chunk, true);
let mut r = Hdf5Reader::open(&path).unwrap();
for c in 0..8u64 {
r.read_slice("d", &[c * chunk as u64], &[1024]).unwrap();
}
assert_eq!(
held(&r),
CHUNK_CACHE_BYTES / (chunk * 8),
"the cache outgrew its byte budget"
);
let _ = std::fs::remove_file(&path);
}
#[test]
#[cfg(feature = "deflate")]
fn changing_an_entry_drops_the_chunk_images() {
let path = dataset("images_entry_mut", 4096, 1024, true);
let mut r = Hdf5Reader::open(&path).unwrap();
r.read_slice("d", &[0], &[256]).unwrap();
assert_eq!(held(&r), 1, "the read kept no image to invalidate");
let pos = r.dataset_position("d").unwrap();
let _ = r.datasets.entry_mut(pos);
assert_eq!(held(&r), 0, "a mutable entry left its chunk images behind");
let _ = std::fs::remove_file(&path);
}
#[test]
#[cfg(feature = "deflate")]
fn a_cached_chunk_places_what_a_cold_read_places() {
let path = dataset("images_differential", 4096, 1024, true);
let mut warm = Hdf5Reader::open(&path).unwrap();
for &(start, count) in &[
(0u64, 100u64),
(100, 100),
(900, 300),
(1024, 1),
(1500, 2000),
(3000, 1096),
(4095, 1),
(0, 4096),
] {
let mut cold = Hdf5Reader::open(&path).unwrap();
assert_eq!(
warm.read_slice("d", &[start], &[count]).unwrap(),
cold.read_slice("d", &[start], &[count]).unwrap(),
"slice {start}+{count} differed once served from the cache"
);
}
let _ = std::fs::remove_file(&path);
}
}
mod chunk_fill_plan {
use super::*;
const DIMS: [u64; 1] = [8];
const CHUNKS: [u64; 1] = [4];
const FILL: [u8; 1] = [0x7E];
fn place(
name: &str,
jobs: Vec<Option<ChunkReadJob>>,
coords: &[u64],
) -> (std::path::PathBuf, Vec<u8>) {
let (path, handle) = handle_over(name, b"ABCDEFGH");
let geo = ChunkOutputGeometry {
dims: &DIMS,
chunk_dims: &CHUNKS,
element_size: 1,
};
let mut out = vec![0xAAu8; 8];
place_chunk_jobs(
&handle,
jobs,
coords,
ChunkReadRequest {
pipeline: None,
target: ChunkTarget::Full,
fill_value: Some(&FILL),
dst: ReadDst::Fresh,
},
&geo,
None,
&mut out,
)
.unwrap();
(path, out)
}
fn job(addr: u64, len: usize) -> Option<ChunkReadJob> {
Some(ChunkReadJob {
addr,
len,
at_most: true,
mask: 0,
})
}
#[test]
fn a_plan_that_covers_the_output_leaves_no_byte_to_fill() {
let geo = ChunkOutputGeometry {
dims: &DIMS,
chunk_dims: &CHUNKS,
element_size: 1,
};
let zeros = [0u64];
let placement = ChunkPlacement::resolve(&geo, ChunkTarget::Full, &zeros);
let jobs = vec![job(0, 4), job(4, 4)];
assert_eq!(
planned_coverage(&placement, &jobs, &[0, 1]),
8,
"a tiling plan must measure as covering the whole output"
);
let (path, out) = place("cover_full", jobs, &[0, 1]);
assert_eq!(out, b"ABCDEFGH");
let _ = std::fs::remove_file(path);
}
#[test]
fn a_slot_the_plan_skips_reads_back_as_fill() {
let (path, out) = place("cover_gap", vec![job(0, 4), None], &[0, 1]);
assert_eq!(out, b"ABCD~~~~");
let _ = std::fs::remove_file(path);
}
#[test]
fn duplicate_chunk_coordinates_still_leave_defined_output() {
let (path, out) = place("cover_dup", vec![job(0, 4), job(0, 4)], &[0, 0]);
assert_eq!(out, b"ABCD~~~~");
let _ = std::fs::remove_file(path);
}
#[test]
fn a_chunk_read_short_has_its_runs_filled_after_placement() {
let (path, out) = place("cover_short", vec![job(0, 4), job(4, 2)], &[0, 1]);
assert_eq!(out, b"ABCD~~~~");
let _ = std::fs::remove_file(path);
}
#[test]
fn a_chunk_address_near_u64_max_fails_the_read_not_the_arithmetic() {
let (path, handle) = handle_over("cover_wrap", b"ABEFCDGH");
let dims = [2u64, 4];
let chunks = [2u64, 2];
let geo = ChunkOutputGeometry {
dims: &dims,
chunk_dims: &chunks,
element_size: 1,
};
let run = |at_most: bool, out: &mut [u8]| {
place_chunk_jobs(
&handle,
vec![
job(0, 4),
Some(ChunkReadJob {
addr: u64::MAX - 1,
len: 4,
at_most,
mask: 0,
}),
],
&[0, 0, 0, 1],
ChunkReadRequest {
pipeline: None,
target: ChunkTarget::Full,
fill_value: Some(&FILL),
dst: ReadDst::Fresh,
},
&geo,
None,
out,
)
};
let mut out = vec![0xAAu8; 8];
run(true, &mut out).unwrap();
assert_eq!(out, b"AB~~EF~~");
let mut out = vec![0xAAu8; 8];
run(false, &mut out).expect_err("an exact read at u64::MAX - 1 cannot succeed");
let _ = std::fs::remove_file(path);
}
}
fn write_le(buf: &mut Vec<u8>, value: u64, n: usize) {
buf.extend_from_slice(&value.to_le_bytes()[..n]);
}
fn build_v0_file(dataset_name: &str, dims: &[u64], data: &[u8]) -> Vec<u8> {
let sa: usize = 8; let ss: usize = 8; let ndims = dims.len();
let element_size = data.len() as u64 / dims.iter().product::<u64>();
let mut file = Vec::new();
let sb_size = 8 + 8 + 4 + 4 * sa + (ss + sa + 4 + 4 + 16); let sb_size_aligned = (sb_size + 7) & !7;
let root_ohdr_addr = sb_size_aligned as u64;
let stab_msg_data_size = 2 * sa; let stab_msg_wire = 8 + stab_msg_data_size;
let stab_msg_wire_aligned = (stab_msg_wire + 7) & !7;
let root_ohdr_data_size = stab_msg_wire_aligned;
let root_ohdr_total = 16 + root_ohdr_data_size; let root_ohdr_total_aligned = (root_ohdr_total + 7) & !7;
let heap_hdr_addr = root_ohdr_addr + root_ohdr_total_aligned as u64;
let heap_hdr_size = 4 + 1 + 3 + ss + ss + sa;
let heap_hdr_size_aligned = (heap_hdr_size + 7) & !7;
let heap_data_addr = heap_hdr_addr + heap_hdr_size_aligned as u64;
let name_bytes = dataset_name.as_bytes();
let heap_data_content_size = 1 + name_bytes.len() + 1; let heap_data_size = (heap_data_content_size + 7) & !7;
let btree_addr = heap_data_addr + heap_data_size as u64;
let btree_size = 4 + 1 + 1 + 2 + 2 * sa + 2 * ss + sa;
let btree_size_aligned = (btree_size + 7) & !7;
let snod_addr = btree_addr + btree_size_aligned as u64;
let entry_size = ss + sa + 4 + 4 + 16;
let snod_size = 8 + entry_size;
let snod_size_aligned = (snod_size + 7) & !7;
let ds_ohdr_addr = snod_addr + snod_size_aligned as u64;
let ds_msg_data_size = 8 + ndims * ss;
let ds_msg_wire = 8 + ds_msg_data_size;
let ds_msg_wire_aligned = (ds_msg_wire + 7) & !7;
let dt_msg_data_size = 12;
let dt_msg_wire = 8 + dt_msg_data_size;
let dt_msg_wire_aligned = (dt_msg_wire + 7) & !7;
let dl_msg_data_size = 2 + sa + ss;
let dl_msg_wire = 8 + dl_msg_data_size;
let dl_msg_wire_aligned = (dl_msg_wire + 7) & !7;
let ds_ohdr_data_size = ds_msg_wire_aligned + dt_msg_wire_aligned + dl_msg_wire_aligned;
let ds_ohdr_total = 16 + ds_ohdr_data_size; let ds_ohdr_total_aligned = (ds_ohdr_total + 7) & !7;
let raw_data_addr = ds_ohdr_addr + ds_ohdr_total_aligned as u64;
let raw_data_size = data.len();
let eof = raw_data_addr + raw_data_size as u64;
let sig: [u8; 8] = [0x89, 0x48, 0x44, 0x46, 0x0d, 0x0a, 0x1a, 0x0a];
file.extend_from_slice(&sig);
file.push(0); file.push(0); file.push(0); file.push(0); file.push(0); file.push(sa as u8); file.push(ss as u8); file.push(0); file.extend_from_slice(&4u16.to_le_bytes()); file.extend_from_slice(&32u16.to_le_bytes()); file.extend_from_slice(&0u32.to_le_bytes()); write_le(&mut file, 0, sa);
write_le(&mut file, UNDEF_ADDR, sa);
write_le(&mut file, eof, sa);
write_le(&mut file, UNDEF_ADDR, sa);
write_le(&mut file, 0, ss); write_le(&mut file, root_ohdr_addr, sa); file.extend_from_slice(&1u32.to_le_bytes()); file.extend_from_slice(&0u32.to_le_bytes()); write_le(&mut file, btree_addr, sa);
write_le(&mut file, heap_hdr_addr, sa);
while file.len() < sb_size_aligned {
file.push(0);
}
assert_eq!(file.len(), root_ohdr_addr as usize);
file.push(1); file.push(0); file.extend_from_slice(&1u16.to_le_bytes()); file.extend_from_slice(&1u32.to_le_bytes()); file.extend_from_slice(&(root_ohdr_data_size as u32).to_le_bytes());
file.extend_from_slice(&[0u8; 4]); file.extend_from_slice(&0x0011u16.to_le_bytes()); file.extend_from_slice(&(stab_msg_data_size as u16).to_le_bytes()); file.push(0); file.extend_from_slice(&[0u8; 3]); write_le(&mut file, btree_addr, sa);
write_le(&mut file, heap_hdr_addr, sa);
while file.len() < (root_ohdr_addr as usize + root_ohdr_total_aligned) {
file.push(0);
}
assert_eq!(file.len(), heap_hdr_addr as usize);
file.extend_from_slice(b"HEAP");
file.push(0); file.extend_from_slice(&[0u8; 3]); write_le(&mut file, heap_data_size as u64, ss); write_le(&mut file, u64::MAX, ss); write_le(&mut file, heap_data_addr, sa); while file.len() < (heap_hdr_addr as usize + heap_hdr_size_aligned) {
file.push(0);
}
assert_eq!(file.len(), heap_data_addr as usize);
file.push(0); file.extend_from_slice(name_bytes); file.push(0); while file.len() < (heap_data_addr as usize + heap_data_size) {
file.push(0);
}
assert_eq!(file.len(), btree_addr as usize);
file.extend_from_slice(b"TREE");
file.push(0); file.push(0); file.extend_from_slice(&1u16.to_le_bytes()); write_le(&mut file, UNDEF_ADDR, sa); write_le(&mut file, UNDEF_ADDR, sa); write_le(&mut file, 0, ss);
write_le(&mut file, snod_addr, sa);
write_le(&mut file, 1, ss);
while file.len() < (btree_addr as usize + btree_size_aligned) {
file.push(0);
}
assert_eq!(file.len(), snod_addr as usize);
file.extend_from_slice(b"SNOD");
file.push(1); file.push(0); file.extend_from_slice(&1u16.to_le_bytes()); write_le(&mut file, 1, ss); write_le(&mut file, ds_ohdr_addr, sa); file.extend_from_slice(&0u32.to_le_bytes()); file.extend_from_slice(&0u32.to_le_bytes()); file.extend_from_slice(&[0u8; 16]); while file.len() < (snod_addr as usize + snod_size_aligned) {
file.push(0);
}
assert_eq!(file.len(), ds_ohdr_addr as usize);
file.push(1); file.push(0); file.extend_from_slice(&3u16.to_le_bytes()); file.extend_from_slice(&1u32.to_le_bytes()); file.extend_from_slice(&(ds_ohdr_data_size as u32).to_le_bytes());
file.extend_from_slice(&[0u8; 4]);
file.extend_from_slice(&0x0001u16.to_le_bytes());
file.extend_from_slice(&(ds_msg_data_size as u16).to_le_bytes());
file.push(0); file.extend_from_slice(&[0u8; 3]); file.push(1); file.push(ndims as u8);
file.push(0); file.push(0); file.extend_from_slice(&[0u8; 4]); for &d in dims {
write_le(&mut file, d, ss);
}
let target = ds_ohdr_addr as usize + 16 + ds_msg_wire_aligned;
while file.len() < target {
file.push(0);
}
file.extend_from_slice(&0x0003u16.to_le_bytes());
file.extend_from_slice(&(dt_msg_data_size as u16).to_le_bytes());
file.push(0); file.extend_from_slice(&[0u8; 3]); file.push(0x10); file.push(0x08); file.push(0); file.push(0); file.extend_from_slice(&(element_size as u32).to_le_bytes()); file.extend_from_slice(&0u16.to_le_bytes()); file.extend_from_slice(&((element_size * 8) as u16).to_le_bytes()); let target = ds_ohdr_addr as usize + 16 + ds_msg_wire_aligned + dt_msg_wire_aligned;
while file.len() < target {
file.push(0);
}
file.extend_from_slice(&0x0008u16.to_le_bytes());
file.extend_from_slice(&(dl_msg_data_size as u16).to_le_bytes());
file.push(0); file.extend_from_slice(&[0u8; 3]); file.push(3); file.push(1); write_le(&mut file, raw_data_addr, sa); write_le(&mut file, raw_data_size as u64, ss); let target = ds_ohdr_addr as usize + ds_ohdr_total_aligned;
while file.len() < target {
file.push(0);
}
assert_eq!(file.len(), raw_data_addr as usize);
file.extend_from_slice(data);
assert_eq!(file.len(), eof as usize);
file
}
#[test]
fn test_read_v0_file_with_one_dataset() {
let dims = [3u64, 4];
let values: Vec<i32> = (0..12).collect();
let raw_data: Vec<u8> = values.iter().flat_map(|v| v.to_le_bytes()).collect();
let file_bytes = build_v0_file("my_dataset", &dims, &raw_data);
let path = temp_path("v0_reader");
{
let mut f = std::fs::File::create(&path).unwrap();
f.write_all(&file_bytes).unwrap();
f.sync_all().unwrap();
}
let mut reader = Hdf5Reader::open(&path).unwrap();
let names = reader.dataset_names();
assert_eq!(names, vec!["my_dataset"]);
let shape = reader.dataset_shape("my_dataset").unwrap();
assert_eq!(shape, vec![3, 4]);
let data = reader.read_dataset_raw("my_dataset").unwrap();
assert_eq!(data, raw_data);
let read_values: Vec<i32> = data
.as_chunks::<4>()
.0
.iter()
.map(|c| i32::from_le_bytes(*c))
.collect();
assert_eq!(read_values, values);
std::fs::remove_file(&path).ok();
}
#[test]
fn test_read_v0_file_1d_dataset() {
let dims = [5u64];
let values: Vec<i32> = vec![100, 200, 300, 400, 500];
let raw_data: Vec<u8> = values.iter().flat_map(|v| v.to_le_bytes()).collect();
let file_bytes = build_v0_file("data_1d", &dims, &raw_data);
let path = temp_path("v0_1d");
{
let mut f = std::fs::File::create(&path).unwrap();
f.write_all(&file_bytes).unwrap();
}
let mut reader = Hdf5Reader::open(&path).unwrap();
assert_eq!(reader.dataset_names(), vec!["data_1d"]);
assert_eq!(reader.dataset_shape("data_1d").unwrap(), vec![5]);
let data = reader.read_dataset_raw("data_1d").unwrap();
let read_values: Vec<i32> = data
.as_chunks::<4>()
.0
.iter()
.map(|c| i32::from_le_bytes(*c))
.collect();
assert_eq!(read_values, values);
std::fs::remove_file(&path).ok();
}
#[test]
fn test_detect_v2v3_still_works() {
let path = temp_path("detect_v3");
{
use crate::io::writer::Hdf5Writer;
let writer = Hdf5Writer::create(&path).unwrap();
let datatype = crate::format::messages::datatype::DatatypeMessage::i32_type();
let idx = writer.create_dataset("test", datatype, &[4]).unwrap();
let data = [1i32, 2, 3, 4];
let raw: Vec<u8> = data.iter().flat_map(|v| v.to_le_bytes()).collect();
writer.write_dataset_raw(idx, &raw).unwrap();
writer.close().unwrap();
}
let mut reader = Hdf5Reader::open(&path).unwrap();
assert_eq!(reader.dataset_names(), vec!["test"]);
let shape = reader.dataset_shape("test").unwrap();
assert_eq!(shape, vec![4]);
let data = reader.read_dataset_raw("test").unwrap();
let vals: Vec<i32> = data
.as_chunks::<4>()
.0
.iter()
.map(|c| i32::from_le_bytes(*c))
.collect();
assert_eq!(vals, vec![1, 2, 3, 4]);
std::fs::remove_file(&path).ok();
}
fn collect_runs(
dims: &[u64],
starts: &[u64],
counts: &[u64],
es: u64,
) -> Vec<(u64, usize, usize)> {
let mut v = Vec::new();
for_each_contiguous_run(dims, starts, counts, es, |s, o, l| {
v.push((s, o, l));
Ok(())
})
.unwrap();
v
}
#[test]
fn coalesce_1d_is_single_run() {
assert_eq!(collect_runs(&[10], &[2], &[3], 1), vec![(2, 0, 3)]);
assert_eq!(collect_runs(&[10], &[2], &[3], 4), vec![(8, 0, 12)]);
}
#[test]
fn coalesce_full_last_dim_merges_into_one_run() {
assert_eq!(collect_runs(&[4, 5], &[1, 0], &[2, 5], 1), vec![(5, 0, 10)]);
}
#[test]
fn coalesce_partial_last_dim_keeps_one_run_per_row() {
assert_eq!(
collect_runs(&[4, 5], &[1, 1], &[2, 3], 1),
vec![(6, 0, 3), (11, 3, 3)]
);
assert_eq!(
collect_runs(&[4, 5], &[1, 1], &[2, 3], 4),
vec![(24, 0, 12), (44, 12, 12)]
);
}
#[test]
fn coalesce_3d_reported_case_one_run_per_outer_index() {
assert_eq!(
collect_runs(&[3, 4, 5], &[0, 1, 0], &[3, 2, 5], 1),
vec![(5, 0, 10), (25, 10, 10), (45, 20, 10)]
);
}
#[test]
fn coalesce_3d_full_inner_dims_is_single_run() {
assert_eq!(
collect_runs(&[3, 4, 5], &[1, 0, 0], &[2, 4, 5], 1),
vec![(20, 0, 40)]
);
}
#[test]
fn read_slice_contiguous_3d_matches_naive_extraction() {
let dims = [3u64, 4, 5];
let total: usize = (dims[0] * dims[1] * dims[2]) as usize;
let values: Vec<i32> = (0..total as i32).collect();
let raw_data: Vec<u8> = values.iter().flat_map(|v| v.to_le_bytes()).collect();
let file_bytes = build_v0_file("vol", &dims, &raw_data);
let path = temp_path("slice_3d_contig");
{
let mut f = std::fs::File::create(&path).unwrap();
f.write_all(&file_bytes).unwrap();
f.sync_all().unwrap();
}
let mut reader = Hdf5Reader::open(&path).unwrap();
let expect = |starts: [u64; 3], counts: [u64; 3]| -> Vec<i32> {
let mut out = Vec::new();
for i in 0..counts[0] {
for j in 0..counts[1] {
for k in 0..counts[2] {
let gi = starts[0] + i;
let gj = starts[1] + j;
let gk = starts[2] + k;
out.push(values[(gi * dims[1] * dims[2] + gj * dims[2] + gk) as usize]);
}
}
}
out
};
let decode = |raw: Vec<u8>| -> Vec<i32> {
raw.as_chunks::<4>()
.0
.iter()
.map(|c| i32::from_le_bytes(*c))
.collect()
};
let cases: &[([u64; 3], [u64; 3])] = &[
([0, 1, 0], [3, 2, 5]), ([1, 0, 0], [2, 4, 5]), ([0, 0, 1], [3, 4, 3]), ([1, 2, 1], [2, 2, 4]), ([0, 0, 0], [3, 4, 5]), ([2, 3, 4], [1, 1, 1]), ];
for &(starts, counts) in cases {
let got = decode(reader.read_slice("vol", &starts, &counts).unwrap());
assert_eq!(
got,
expect(starts, counts),
"slice starts={starts:?} counts={counts:?}"
);
}
let _ = std::fs::remove_file(&path);
}
fn collect_from(
messages: Vec<crate::format::object_header::ObjectHeaderMessage>,
) -> Result<Vec<String>, String> {
let path = temp_path("collect");
std::fs::File::create(&path).unwrap();
let mut handle = FileHandle::open_read(&path).unwrap();
let ctx = FormatContext {
sizeof_addr: 8,
sizeof_size: 8,
};
let header = ObjectHeader {
flags: 0x02,
times: None,
messages,
};
let attrs = collect_object_attributes(&mut handle, &ctx, &header);
drop(handle);
let _ = std::fs::remove_file(&path);
attrs
.complete("obj")
.map(|e| e.iter().map(|a| a.name().to_string()).collect())
.map_err(|e| e.to_string())
}
fn msg(msg_type: u8, data: Vec<u8>) -> crate::format::object_header::ObjectHeaderMessage {
crate::format::object_header::ObjectHeaderMessage {
msg_type,
flags: 0,
data,
creation_index: 0,
}
}
#[test]
fn an_undecodable_attribute_info_message_fails_the_listing() {
let err = collect_from(vec![msg(MSG_ATTR_INFO, vec![9, 0])]).unwrap_err();
assert!(
err.contains("attributes of 'obj' cannot be read whole")
&& err.contains("attribute info message"),
"{err}"
);
}
#[test]
fn an_unnameable_attribute_message_fails_the_listing() {
let err = collect_from(vec![msg(MSG_ATTRIBUTE, vec![1, 0, 0])]).unwrap_err();
assert!(
err.contains("attributes of 'obj' cannot be read whole")
&& err.contains("attribute message"),
"{err}"
);
}
#[test]
fn a_whole_attribute_set_still_lists() {
let ctx = FormatContext {
sizeof_addr: 8,
sizeof_size: 8,
};
let ainfo = crate::format::messages::attr_info::AttributeInfoMessage::compact();
let names = collect_from(vec![msg(MSG_ATTR_INFO, ainfo.encode(&ctx))]).unwrap();
assert!(names.is_empty(), "{names:?}");
}
#[test]
fn fixed_string_attr_value_honors_the_declared_pad() {
use crate::format::messages::dataspace::DataspaceMessage;
use crate::format::messages::datatype::DatatypeMessage;
let attr = |padding: u8, data: &[u8]| AttributeMessage {
name: "units".to_string(),
datatype: DatatypeMessage::FixedString {
size: data.len() as u32,
padding,
charset: 0,
},
dataspace: DataspaceMessage::scalar(),
data: data.to_vec(),
};
assert_eq!(
fixed_string_attr_value(&attr(2, b"volt ")).unwrap(),
"volt"
);
assert_eq!(
fixed_string_attr_value(&attr(0, b"volt \0\0")).unwrap(),
"volt "
);
assert_eq!(
fixed_string_attr_value(&attr(1, b"volt\0\0\0\0")).unwrap(),
"volt"
);
let err = fixed_string_attr_value(&attr(7, b"volt ")).unwrap_err();
assert!(
err.to_string().contains("padding rule 7"),
"unexpected error: {err}"
);
}
}
#[cfg(test)]
mod h5py_debug_tests {
use super::*;
#[test]
fn debug_read_h5py() {
let path = std::path::Path::new("/tmp/test_h5py_default.h5");
if !path.exists() {
return;
}
let handle = FileHandle::open_read(path).unwrap();
let sb_buf = handle.read_at_most(0, 1024).unwrap();
let version = detect_superblock_version(&sb_buf).unwrap();
eprintln!("Superblock version: {}", version);
let sb = SuperblockV0V1::decode(&sb_buf).unwrap();
eprintln!(
"sizeof_addr={}, sizeof_size={}",
sb.sizeof_offsets, sb.sizeof_lengths
);
let (ste_btree, ste_heap) = sb
.root_symbol_table_entry
.cached_symbol_table()
.unwrap_or((UNDEF_ADDR, UNDEF_ADDR));
eprintln!(
"STE: obj_header={}, cache={:?}, btree={}, heap={}",
sb.root_symbol_table_entry.obj_header_addr,
sb.root_symbol_table_entry.cache,
ste_btree,
ste_heap
);
let ctx = FormatContext {
sizeof_addr: sb.sizeof_offsets,
sizeof_size: sb.sizeof_lengths,
};
let heap_buf = handle.read_at_most(ste_heap, 128).unwrap();
let heap_hdr = LocalHeapHeader::decode(
&heap_buf,
ctx.sizeof_addr as usize,
ctx.sizeof_size as usize,
)
.unwrap();
eprintln!(
"Heap data_addr={}, data_size={}",
heap_hdr.data_addr, heap_hdr.data_size
);
let heap_data = handle
.read_at(heap_hdr.data_addr, heap_hdr.data_size as usize)
.unwrap();
eprintln!(
"Heap data bytes: {:?}",
&heap_data[..std::cmp::min(64, heap_data.len())]
);
let btree_buf = handle.read_at_most(ste_btree, 8192).unwrap();
let btree = BTreeV1Node::decode(
&btree_buf,
ctx.sizeof_addr as usize,
ctx.sizeof_size as usize,
BTreeV1Config::default().snode_max_entries(),
)
.unwrap();
eprintln!(
"BTree: type={}, level={}, entries={}, children={:?}",
btree.node_type, btree.level, btree.entries_used, btree.children
);
for &child in &btree.children {
let snod_buf = handle.read_at_most(child, 8192).unwrap();
let snod = SymbolTableNode::decode(
&snod_buf,
ctx.sizeof_addr as usize,
ctx.sizeof_size as usize,
BTreeV1Config::default().sym_leaf_max_entries(),
)
.unwrap();
eprintln!("SNOD at {}: {} entries", child, snod.entries.len());
for entry in &snod.entries {
let name = local_heap_get_string(&heap_data, entry.name_offset).unwrap();
eprintln!(
" entry: name='{}' (offset={}), obj_header={}, cache={:?}",
name, entry.name_offset, entry.obj_header_addr, entry.cache
);
}
}
let reader = Hdf5Reader::open(path).unwrap();
eprintln!("Datasets found: {:?}", reader.dataset_names());
}
const TEST_PYTHON: &str = "/Users/stevek/mamba/envs/bs2026.1/bin/python";
fn temp_path(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_gap_test_{}_{}_{}.h5",
name,
std::process::id(),
n
))
}
fn gen_fixture(script: &str) -> bool {
if !std::path::Path::new(TEST_PYTHON).exists() {
return false;
}
let status = std::process::Command::new(TEST_PYTHON)
.arg("-c")
.arg(script)
.status();
matches!(status, Ok(s) if s.success())
}
#[test]
fn gap1_v2_root_continuation_block() {
let path = temp_path("gap1_cont");
let p = path.display().to_string();
let script = format!(
"import h5py,numpy as np\n\
f=h5py.File(r'{p}','w',libver='latest')\n\
[f.create_dataset('ds_%d'%i,data=np.arange(i*10,i*10+10,dtype='int32')) for i in range(6)]\n\
f.close()"
);
if !gen_fixture(&script) {
eprintln!("skipping gap1: python unavailable");
return;
}
let mut reader = Hdf5Reader::open(&path).unwrap();
let mut names = reader.dataset_names();
names.sort();
assert_eq!(
names,
vec!["ds_0", "ds_1", "ds_2", "ds_3", "ds_4", "ds_5"],
"all 6 datasets must be found across the continuation block"
);
let raw = reader.read_dataset_raw("ds_3").unwrap();
let vals: Vec<i32> = raw
.as_chunks::<4>()
.0
.iter()
.map(|c| i32::from_le_bytes(*c))
.collect();
assert_eq!(vals, (30..40).collect::<Vec<i32>>());
let _ = std::fs::remove_file(&path);
}
#[test]
fn gap2_v2_dense_fractal_heap_links() {
let path = temp_path("gap2_dense");
let p = path.display().to_string();
let script = format!(
"import h5py,numpy as np\n\
f=h5py.File(r'{p}','w',libver='latest')\n\
g=f.create_group('dense')\n\
[g.create_dataset('d%02d'%i,data=np.full(4,i,dtype='float64')) for i in range(14)]\n\
f.close()"
);
if !gen_fixture(&script) {
eprintln!("skipping gap2: python unavailable");
return;
}
let mut reader = Hdf5Reader::open(&path).unwrap();
let mut names = reader.dataset_names();
names.sort();
let expected: Vec<String> = (0..14).map(|i| format!("dense/d{:02}", i)).collect();
assert_eq!(
names, expected,
"all 14 dense-stored links must be recovered from the fractal heap"
);
let raw = reader.read_dataset_raw("dense/d07").unwrap();
let vals: Vec<f64> = raw
.as_chunks::<8>()
.0
.iter()
.map(|c| f64::from_le_bytes(*c))
.collect();
assert_eq!(vals, vec![7.0; 4]);
let _ = std::fs::remove_file(&path);
}
#[test]
fn gap3_v0v1_legacy_subgroups() {
let path = temp_path("gap3_legacy");
let p = path.display().to_string();
let script = format!(
"import h5py,numpy as np\n\
f=h5py.File(r'{p}','w',libver='earliest')\n\
g1=f.create_group('grp1')\n\
g1.create_dataset('a',data=np.arange(5,dtype='int16'))\n\
g2=g1.create_group('sub')\n\
g2.create_dataset('b',data=np.arange(7,dtype='int64'))\n\
f.create_dataset('top',data=np.arange(3,dtype='int32'))\n\
f.close()"
);
if !gen_fixture(&script) {
eprintln!("skipping gap3: python unavailable");
return;
}
let mut reader = Hdf5Reader::open(&path).unwrap();
let mut names = reader.dataset_names();
names.sort();
assert_eq!(
names,
vec!["grp1/a", "grp1/sub/b", "top"],
"datasets nested in legacy symbol-table subgroups must be found"
);
let raw = reader.read_dataset_raw("grp1/sub/b").unwrap();
let vals: Vec<i64> = raw
.as_chunks::<8>()
.0
.iter()
.map(|c| i64::from_le_bytes(*c))
.collect();
assert_eq!(vals, (0..7).collect::<Vec<i64>>());
let _ = std::fs::remove_file(&path);
}
#[test]
fn nbit_chunked_post_filter_conversion() {
let path = temp_path("nbit_conv");
let p = path.display().to_string();
let script = format!(
"import h5py,numpy as np\n\
from h5py import h5t,h5p,h5s,h5d,h5f,h5z\n\
fid=h5f.create(r'{p}'.encode())\n\
def mk(name,bt,prec,off,npd,vals,chunk):\n\
\x20dt=bt.copy();dt.set_precision(prec);dt.set_offset(off)\n\
\x20arr=np.ascontiguousarray(np.asarray(vals,dtype=npd))\n\
\x20sp=h5s.create_simple(arr.shape)\n\
\x20dc=h5p.create(h5p.DATASET_CREATE);dc.set_chunk(chunk)\n\
\x20dc.set_filter(h5z.FILTER_NBIT,h5z.FLAG_OPTIONAL,())\n\
\x20ds=h5d.create(fid,name.encode(),dt,sp,dc)\n\
\x20ds.write(h5s.ALL,h5s.ALL,arr);ds.close()\n\
mk('u4_p17_o3',h5t.STD_U32LE,17,3,'u4',[0,1,1000,65535,131071,70000,42,99999],(4,))\n\
mk('i4_p13_o5',h5t.STD_I32LE,13,5,'i4',[-5,-1,0,1,7,-4096,4095,-77,42,100,-100,3],(4,))\n\
mk('i2_p9_o4',h5t.STD_I16LE,9,4,'i2',[-256,-1,0,1,255,-7,7,-200],(3,))\n\
mk('i4_2d_p11_o6',h5t.STD_I32LE,11,6,'i4',np.array([[-1024,-1,0,5],[1023,-77,88,-3]],dtype='i4'),(1,4))\n\
fid.close()"
);
if !gen_fixture(&script) {
eprintln!("skipping nbit_chunked_post_filter_conversion: python unavailable");
return;
}
let mut reader = Hdf5Reader::open(&path).unwrap();
let raw = reader.read_dataset_raw("u4_p17_o3").unwrap();
let got: Vec<u32> = raw
.as_chunks::<4>()
.0
.iter()
.map(|c| u32::from_le_bytes(*c))
.collect();
assert_eq!(
got,
vec![0u32, 1, 1000, 65535, 131071, 70000, 42, 99999],
"u4 N-bit dataset must decode to exact unsigned values"
);
let raw = reader.read_dataset_raw("i4_p13_o5").unwrap();
let got: Vec<i32> = raw
.as_chunks::<4>()
.0
.iter()
.map(|c| i32::from_le_bytes(*c))
.collect();
assert_eq!(
got,
vec![-5i32, -1, 0, 1, 7, -4096, 4095, -77, 42, 100, -100, 3],
"i4 N-bit dataset must sign-extend negative values"
);
let raw = reader.read_dataset_raw("i2_p9_o4").unwrap();
let got: Vec<i16> = raw
.as_chunks::<2>()
.0
.iter()
.map(|c| i16::from_le_bytes(*c))
.collect();
assert_eq!(
got,
vec![-256i16, -1, 0, 1, 255, -7, 7, -200],
"i2 N-bit dataset must sign-extend negative values"
);
let raw = reader.read_dataset_raw("i4_2d_p11_o6").unwrap();
let got: Vec<i32> = raw
.as_chunks::<4>()
.0
.iter()
.map(|c| i32::from_le_bytes(*c))
.collect();
assert_eq!(
got,
vec![-1024i32, -1, 0, 5, 1023, -77, 88, -3],
"2D i4 N-bit dataset must decode element-exact"
);
let raw = reader.read_slice("i4_p13_o5", &[4], &[3]).unwrap();
let got: Vec<i32> = raw
.as_chunks::<4>()
.0
.iter()
.map(|c| i32::from_le_bytes(*c))
.collect();
assert_eq!(got, vec![7i32, -4096, 4095], "read_slice must convert too");
let raw = reader.read_slice("i4_2d_p11_o6", &[1, 0], &[1, 4]).unwrap();
let got: Vec<i32> = raw
.as_chunks::<4>()
.0
.iter()
.map(|c| i32::from_le_bytes(*c))
.collect();
assert_eq!(
got,
vec![1023i32, -77, 88, -3],
"2D read_slice must convert"
);
let _ = std::fs::remove_file(&path);
}
#[test]
fn read_slice_chunked_all_index_types() {
let latest = temp_path("slice_chunk_latest");
let earliest = temp_path("slice_chunk_earliest");
let pl = latest.display().to_string();
let pe = earliest.display().to_string();
let script = format!(
"import h5py,numpy as np\n\
a=np.arange(5*4*6,dtype='int32').reshape(5,4,6)\n\
f=h5py.File(r'{pl}','w',libver='latest')\n\
f.create_dataset('single',data=a,chunks=(5,4,6))\n\
f.create_dataset('fa',data=a,chunks=(2,2,2))\n\
f.create_dataset('ea',data=a,chunks=(2,2,2),maxshape=(None,4,6))\n\
f.create_dataset('btv2',data=a,chunks=(2,2,2),maxshape=(None,None,6))\n\
f.close()\n\
g=h5py.File(r'{pe}','w',libver='earliest')\n\
g.create_dataset('btv1',data=a,chunks=(2,2,2))\n\
g.close()"
);
if !gen_fixture(&script) {
eprintln!("skipping read_slice_chunked_all_index_types: python unavailable");
return;
}
let dims = [5u64, 4, 6];
let val = |i: u64, j: u64, k: u64| (i * dims[1] * dims[2] + j * dims[2] + k) as i32;
let expect = |starts: [u64; 3], counts: [u64; 3]| -> Vec<i32> {
let mut out = Vec::new();
for i in 0..counts[0] {
for j in 0..counts[1] {
for k in 0..counts[2] {
out.push(val(starts[0] + i, starts[1] + j, starts[2] + k));
}
}
}
out
};
let decode = |raw: Vec<u8>| -> Vec<i32> {
raw.as_chunks::<4>()
.0
.iter()
.map(|c| i32::from_le_bytes(*c))
.collect()
};
let cases: &[([u64; 3], [u64; 3])] = &[
([0, 1, 0], [5, 2, 6]), ([1, 0, 0], [3, 4, 6]), ([0, 0, 2], [5, 4, 3]), ([2, 1, 3], [1, 2, 2]), ([0, 0, 0], [5, 4, 6]), ([4, 3, 5], [1, 1, 1]), ([1, 1, 1], [3, 3, 4]), ];
let mut reader_l = Hdf5Reader::open(&latest).unwrap();
for name in ["single", "fa", "ea", "btv2"] {
let full = decode(reader_l.read_dataset_raw(name).unwrap());
assert_eq!(full, expect([0, 0, 0], [5, 4, 6]), "{name} full read");
for &(starts, counts) in cases {
let got = decode(reader_l.read_slice(name, &starts, &counts).unwrap());
assert_eq!(
got,
expect(starts, counts),
"{name} slice starts={starts:?} counts={counts:?}"
);
}
}
let mut reader_e = Hdf5Reader::open(&earliest).unwrap();
let full = decode(reader_e.read_dataset_raw("btv1").unwrap());
assert_eq!(full, expect([0, 0, 0], [5, 4, 6]), "btv1 full read");
for &(starts, counts) in cases {
let got = decode(reader_e.read_slice("btv1", &starts, &counts).unwrap());
assert_eq!(
got,
expect(starts, counts),
"btv1 slice starts={starts:?} counts={counts:?}"
);
}
let _ = std::fs::remove_file(&latest);
let _ = std::fs::remove_file(&earliest);
}
#[test]
#[cfg(feature = "deflate")]
fn slice_reads_agree_however_the_chunk_reaches_the_output() {
let path = temp_path("slice_run_placement");
let (rows, cols) = (8usize, 4096usize); let data: Vec<f64> = (0..rows * cols).map(|i| i as f64).collect();
let (rows, cols) = (rows as u64, cols as u64);
{
let file = crate::H5File::create(&path).unwrap();
let ds = file
.new_dataset::<f64>()
.shape([rows as usize, cols as usize])
.chunk(&[2, cols as usize])
.create("plain")
.unwrap();
ds.write_raw(&data).unwrap();
let ds = file
.new_dataset::<f64>()
.shape([rows as usize, cols as usize])
.chunk(&[2, cols as usize])
.deflate(1)
.create("zipped")
.unwrap();
ds.write_raw(&data).unwrap();
file.close().unwrap();
}
let expect = |starts: [u64; 2], counts: [u64; 2]| -> Vec<f64> {
let mut out = Vec::new();
for i in 0..counts[0] {
for j in 0..counts[1] {
out.push(data[((starts[0] + i) * cols + starts[1] + j) as usize]);
}
}
out
};
let decode = |raw: Vec<u8>| -> Vec<f64> {
raw.as_chunks::<8>()
.0
.iter()
.map(|c| f64::from_le_bytes(*c))
.collect()
};
let cases: &[([u64; 2], [u64; 2])] = &[
([0, 0], [rows, cols]), ([3, 0], [4, cols]), ([1, 7], [5, 3]), ([2, 1000], [2, 2048]), ([7, 4095], [1, 1]), ];
let mut reader = Hdf5Reader::open(&path).unwrap();
for name in ["plain", "zipped"] {
let full = decode(reader.read_dataset_raw(name).unwrap());
assert_eq!(full, data, "{name} full read");
for &(starts, counts) in cases {
let got = decode(reader.read_slice(name, &starts, &counts).unwrap());
assert_eq!(
got,
expect(starts, counts),
"{name} slice starts={starts:?} counts={counts:?}"
);
}
}
std::fs::remove_file(&path).ok();
}
}