diff --git a/doc/getting_started/installation.rst b/doc/getting_started/installation.rst index 6bb842460..8d51e063c 100644 --- a/doc/getting_started/installation.rst +++ b/doc/getting_started/installation.rst @@ -44,9 +44,8 @@ grouped into *extras* that you opt into with the ``blosc2[extra]`` syntax: - Lazy Zarr sources and the ``b2nd-to-zarr`` (or ``blosc2-to-zarr``) converter (``zarr``). * - ``hdf5`` - - Reading HDF5 datasets lazily as virtual arrays: ``h5py`` for local - files, ``kerchunk`` for remote ones (``kerchunk``, ``h5py``, - ``hdf5plugin``). + - Reading HDF5 datasets lazily as virtual arrays, locally or through + fsspec (``h5py``, ``hdf5plugin``). * - ``fsspec`` - Reading and writing single-file containers through any `fsspec `_ URL. The HTTP(S) driver is diff --git a/doc/guides/b2view.rst b/doc/guides/b2view.rst index f21763346..c980e3f7b 100644 --- a/doc/guides/b2view.rst +++ b/doc/guides/b2view.rst @@ -97,7 +97,7 @@ objects and chunks. Small objects may fit entirely within a bounded opening read ``--profile`` and ``--endpoint-url`` are optional; when omitted, the S3 backend uses its normal credential and endpoint configuration. Install ``blosc2[tui]`` and ``blosc2[fsspec]`` plus ``s3fs`` for S3 access. Zarr requires -``blosc2[zarr]``; HDF5 requires ``blosc2[hdf5]`` (Kerchunk, h5py, and Zarr). +``blosc2[zarr]``; HDF5 requires ``blosc2[hdf5]`` (h5py and optional filter plugins). B2Z browsing does not require Zarr or HDF5 dependencies. Format details and limits: @@ -115,13 +115,12 @@ Format details and limits: support and LIST permission. Direct arrays do not require listing their parent. Empty groups and attributes are preserved. Unknown codecs and unsupported dtypes remain visible; preview support follows the existing Zarr array reader. -* **HDF5:** Kerchunk translates metadata once per session and all selected leaves - reuse those references. Translation can enumerate many chunk references and - inline small values; it avoids full-file localization, but is not a constant-cost - operation. Empty groups and attributes are preserved. Failed dataset translations - become unavailable nodes without hiding supported siblings. The view covers - Kerchunk's representation: hard-link aliases may be omitted, and soft/external - links and group cycles are not followed. +* **HDF5:** h5py builds a native metadata and chunk-range index once per session, + and all selected leaves reuse it. Indexing can enumerate many allocated chunks; + it avoids full-file localization, but is not a constant-cost operation. Empty + groups and attributes are preserved. Unsupported datasets remain unavailable + nodes without hiding supported siblings. Hard-link aliases may be omitted, and + soft/external links and group cycles are not followed. These internal browser adapters do not change the array-only contract of ``blosc2.open(..., lazy=True)`` or add a persisted RemoteArray hierarchy descriptor. diff --git a/doc/guides/remote_arrays.md b/doc/guides/remote_arrays.md index 6a71caa1b..0630d3262 100644 --- a/doc/guides/remote_arrays.md +++ b/doc/guides/remote_arrays.md @@ -36,7 +36,7 @@ The argument passed to {func}`blosc2.open` selects the route: | A URL string such as `s3://...` or `https://...` | fsspec | A byte-addressable, standalone `.b2nd` file | | A URL containing a `.b2z` path component, or `source_format="b2z"` | B2Z | One immutable external NDArray leaf in a `.b2z` archive | | A URL containing a `.zarr` path component | Zarr | One immutable Zarr v2 or v3 array | -| A URL containing a `.h5` or `.hdf5` path component, or `source_format="hdf5"` | HDF5 | One immutable HDF5 dataset (h5py locally, kerchunk remotely) | +| A URL containing a `.h5` or `.hdf5` path component, or `source_format="hdf5"` | HDF5 | One immutable HDF5 dataset (h5py locally, native range index remotely) | | A {ref}`URLPath` | Caterva2 | One array-like dataset on a Caterva2 server | ```python @@ -67,7 +67,7 @@ h1 = blosc2.open( store_b2z = blosc2.RemoteStore("s3://bucket/hierarchy.b2z") store_h5 = blosc2.RemoteStore("https://datasets.example.org/hierarchy.h5") -# Reopen an exported reference snapshot (.b2z): +# Reopen an exported portable snapshot (.b2z): store_snap = blosc2.open("snapshot.b2z") a1.shape, a1.dtype # metadata is available immediately @@ -89,13 +89,14 @@ Converted Blosc2 chunks are cached under an immutable source contract, so publis Remote HDF5 needs `pip install "blosc2[hdf5,fsspec]"`. HTTP/HTTPS works directly; cloud stores require their protocol driver (`s3fs` for S3, etc.). Datasets can be specified using standard slash syntax (`file.h5/d0/d1/a2`), the double-colon separator (`file.h5::d0/d1/a2`), or the `dataset="d0/d1/a2"` parameter. -Pre-indexing is performed via `kerchunk`. When opening a single {ref}`RemoteArray`, the resulting reference map is cached inside the array carrier (`schunk.vlmeta["hdf5-refs"]`). When using {ref}`RemoteStore`, indexing is performed once for the entire container and shared across all leaves and sessions. +Remote pre-indexing uses h5py to record dataset metadata and allocated chunk byte ranges. When opening a single {ref}`RemoteArray`, the native index is cached inside the array carrier (`schunk.vlmeta["hdf5-index"]`). When using {ref}`RemoteStore`, indexing is performed once for the entire container and shared across all leaves and sessions. Uncompressed, deflate, shuffle, and Blosc2 pipelines are decoded directly after fsspec range reads. Other pipelines use a retained h5py reader, including filters registered by `hdf5plugin`. Use `blosc2.available_datasets(url)` to inspect datasets in an HDF5 container. -Local HDF5 files use h5py directly, without kerchunk pre-indexing or Zarr/fsspec -dependencies. For example, `blosc2.open("hierarchy.h5::/d0/a2")` reads the selected +Local HDF5 files use h5py directly, without pre-indexing or an fsspec +dependency. For example, `blosc2.open("hierarchy.h5::/d0/a2")` reads the selected dataset through h5py and caches converted Blosc2 chunks in memory. Explicit -`refs=` inputs retain the reference-based reader, including for local files. +`hdf5_index=` accepts a native HDF5 index, including for local files. Legacy +HDF5 reference maps are rejected; omit `hdf5_index=` to regenerate the native index. `RemoteArray` assumes remote sources are immutable by default, avoiding a metadata request before every read. For a replaceable `.b2nd` or Caterva2 source, pass `assume_immutable=False` to refresh its identity and invalidate stale cached chunks before each operation. @@ -151,8 +152,9 @@ What differs between the transports is the types of remote objects each can open `lazy=True` changes *when* data is fetched; it does not expand the underlying storage formats supported by either route. -> [!TIP] -> **Browse remote hierarchies**: To explore groups, inspect attributes, or preview arrays in remote `.b2z`, `.zarr`, or `.h5` containers interactively in the terminal without downloading the complete container, use {doc}`b2view ` (e.g. `b2view s3://bucket/hierarchy.b2z`). To navigate containers programmatically in Python, use {ref}`RemoteStore`. +```{tip} +**Browse remote hierarchies**: To explore groups, inspect attributes, or preview arrays in remote `.b2z`, `.zarr`, or `.h5` containers interactively in the terminal without downloading the complete container, use {doc}`b2view ` (e.g. `b2view s3://bucket/hierarchy.b2z`). To navigate containers programmatically in Python, use {ref}`RemoteStore`. +``` ## Explore remote hierarchies with RemoteStore @@ -203,7 +205,7 @@ with blosc2.RemoteStore( ``` When reopening the same store later with the same `cache_dir`: -- Discovery metadata (such as B2Z member offsets or HDF5 Kerchunk reference maps) is restored from local disk, avoiding repeated remote translation scans. `store.metadata_bytes` reports the encoded manifest size. +- Discovery metadata (such as B2Z member offsets or native HDF5 indexes) is restored from local disk, avoiding repeated remote scans. `store.metadata_bytes` reports the encoded manifest size. - Retained leaf chunks are available immediately from disk without network transfers. - Single-owner locks ensure that concurrent processes do not corrupt the shared cache. @@ -468,7 +470,7 @@ with blosc2.RemoteStore("https://datasets.example.org/data.h5") as store: store["experiment"].save("experiment_sub.b2z") ``` -- **Portable reference**: The `.b2z` archive contains the discovered hierarchy, attributes, and source locators (such as the HDF5 Kerchunk reference map or B2Z member offsets), but no secrets or credentials. +- **Portable reference**: The `.b2z` archive contains the discovered hierarchy, attributes, and source locators (such as the native HDF5 index or B2Z member offsets), but no secrets or credentials. - **`include_cache=True` (default)**: Bundles warm cached chunks along with metadata so reading previously fetched slices requires zero network traffic. - **`include_cache=False`**: Omits cached payload chunks, producing a minimal reference archive for remote streaming. - **Subtree export**: Calling `save()` on a group view exports that subtree with relative child keys and the appropriate source root. @@ -517,8 +519,9 @@ store.save("writable.b2z", mutable=True) | **Immutable** (`mutable=False`, default) | Reads directly in-place from `.b2z` without disk writes. Safe on read-only media (`chmod 0o444`). | Fetched transiently into RAM to satisfy the read; never written to disk or the archive. | `fetch()`, `afetch()`, `trim_cache()`, and `refresh()` are disallowed. | | **Mutable** (`mutable=True`) | Staged into an independent writable runtime cache directory. Original `.b2z` stays untouched. | Fetched and cached to disk under standard LRU eviction rules. | Fully supported. Can be opened with a smaller budget, trimming excess chunks. | -> [!TIP] -> Use **immutable snapshots** (`mutable=False`) for sharing reproducible, read-only reference archives or distributing datasets that should never modify local storage. Use **mutable snapshots** (`mutable=True`) when users should be able to expand the local cache with newly fetched regions over time. +```{tip} +Use **immutable snapshots** (`mutable=False`) for sharing reproducible, read-only reference archives or distributing datasets that should never modify local storage. Use **mutable snapshots** (`mutable=True`) when users should be able to expand the local cache with newly fetched regions over time. +``` ## Retrieve scattered points diff --git a/doc/reference/hdf5ndsource.rst b/doc/reference/hdf5ndsource.rst index a70e04ef0..969db5606 100644 --- a/doc/reference/hdf5ndsource.rst +++ b/doc/reference/hdf5ndsource.rst @@ -4,16 +4,17 @@ HDF5NDSource ============ ``HDF5NDSource`` exposes an HDF5 dataset through :ref:`ProxyNDSource`. -Local files use ``h5py`` directly, without scanning the file with ``kerchunk``. -Remote files and explicit ``refs=`` inputs use ``kerchunk`` reference maps. +Local files use ``h5py`` directly. Remote files are scanned once with ``h5py`` +to build a native byte-range index, which can also be supplied through ``hdf5_index=``. Individual chunks are fetched on demand and converted to Blosc2-compressed chunks stored in the surrounding :ref:`Proxy` or :ref:`RemoteArray` cache. The source is assumed immutable (``assume_immutable=True``). It supports fixed-size boolean, integer, floating-point, complex, and fixed-length string arrays. -HDF5 filters such as Blosc2 (via ``hdf5plugin``), gzip, and uncompressed datasets -are supported. +Uncompressed, gzip/deflate, shuffle, and Blosc2 chunks use direct range reads. +Other filter pipelines fall back to retained ``h5py`` dataset reads; +``hdf5plugin`` enables its additional registered filters. Local reads require ``h5py``; ``hdf5plugin`` enables additional HDF5 filters. Install the full HDF5 support with ``pip install "blosc2[hdf5]"``. Remote datasets also diff --git a/doc/reference/remotearray.rst b/doc/reference/remotearray.rst index db11b2985..262485645 100644 --- a/doc/reference/remotearray.rst +++ b/doc/reference/remotearray.rst @@ -58,10 +58,10 @@ the double-colon separator (``.../file.h5::dataset``), or the ``dataset="dataset Zarr containers similarly accept all three forms (``.../file.zarr/dataset``, ``.../file.zarr::dataset``, or ``dataset="dataset"``). HDF5 datasets on a local path are read directly with ``h5py``; remote HDF5 files use -``kerchunk`` metadata pre-indexing. Like Zarr, HDF5 sources +native metadata pre-indexing with ``h5py`` and byte ranges through fsspec. Like Zarr, HDF5 sources are assumed immutable (``assume_immutable=True``); mutable HDF5 sources are not supported. -Pre-computed kerchunk references can be supplied via ``refs`` to avoid remote scanning -(or to use the reference reader for a local file). +Pre-computed native HDF5 indexes can be supplied via ``hdf5_index`` to avoid remote scanning +(including when opening a local file through the indexed reader). .. code-block:: python @@ -91,8 +91,8 @@ the same three addressing forms: Use ``source_format="b2z"`` for suffix-free archive URLs. The dataset is a logical tree key without the member's ``.b2nd`` suffix. The native Blosc2 reader preserves -source chunks, blocks, dtype, and compression parameters; no kerchunk, Zarr, or -HDF5 dependencies are needed. Install the fsspec extra and the protocol backend. +source chunks, blocks, dtype, and compression parameters; no Zarr or HDF5 +dependencies are needed. Install the fsspec extra and the protocol backend. Opening reads the ZIP directory and selected member's headers. Directory cost scales with archive member count. An 8 KiB archive tail and 16 KiB member prefix diff --git a/doc/reference/remotestore.rst b/doc/reference/remotestore.rst index 3fa77447d..d2b0cdac0 100644 --- a/doc/reference/remotestore.rst +++ b/doc/reference/remotestore.rst @@ -5,7 +5,7 @@ RemoteStore ``RemoteStore`` discovers a read-only B2Z, Zarr or HDF5 hierarchy and returns :ref:`RemoteArray` leaves. Groups and arrays share one source session: a B2Z -archive, an HDF5 reference map, or a Zarr store. Zarr listing remains lazy. +archive, a native HDF5 index, or a Zarr store. Zarr listing remains lazy. The default ``CachePolicy.MEMORY`` shares a 256 MiB allowance across all leaves. Set ``max_cache_bytes`` to a positive integer to change it. ``CachePolicy.NONE`` @@ -62,7 +62,7 @@ Closing a handle, or exiting its context, releases its ownership. Existing child handles remain usable until closed or garbage-collected. The last handle closes the owned archive/store wrappers and private HTTP/S3 transport sessions. Operations on an explicitly closed handle raise ``RuntimeError``. Standalone ``RemoteArray`` exports remain self-contained references, including -the HDF5 reference map when applicable. +the native HDF5 index when applicable. ``b2view`` uses ``RemoteStore`` for remote hierarchies with one 64 MiB MEMORY allowance, and ``RemoteArray`` for selected or directly opened leaves. Switching @@ -85,7 +85,7 @@ the operating system releases the lock after a process exits or crashes. Reopening restores all previously created leaf caches and trims them against the new aggregate allowance before returning. The manifest preserves B2Z directory -and bounded metadata reads, one HDF5 reference map, and lazily discovered Zarr +and bounded metadata reads, one native HDF5 index, and lazily discovered Zarr metadata. Metadata reads can contain small inline values or incidental bytes in bounded prefixes; they are separate from evictable payload. ``metadata_bytes`` is the encoded manifest size, and is zero without a disk manifest. Credentials diff --git a/plans/remote-hdf5.md b/plans/remote-hdf5.md new file mode 100644 index 000000000..8bd9c5b8d --- /dev/null +++ b/plans/remote-hdf5.md @@ -0,0 +1,391 @@ +# Native remote HDF5 reader + +## Status and objective + +Implemented on the current `remote-hdf5` branch. This document records the +architecture, compatibility policy, and verification scope. + +Replace the former HDF5 reference-based access with a native metadata index and range reader. +Keep h5py for metadata discovery and a compatibility fallback, fsspec for remote +I/O, and the existing Blosc2 chunk cache. Neither the legacy HDF5 reference layer nor Zarr should be +required or imported by the new HDF5 path. + +Agreed decoding policy: + +- Direct reads for uncompressed chunked datasets and pipelines composed entirely + of HDF5 deflate (filter 1), HDF5 shuffle (filter 2), and Blosc2 (filter 32026). +- Use `h5py.Dataset[selection]` for other filter pipelines, with `hdf5plugin` + supplying registered HDF5 filters where needed. +- Select the path for the complete dataset filter pipeline. An unsupported + filter selects the fallback even if some chunks skip that filter. +- Keep fallback handles open across chunk reads so HDF5 can reuse metadata. +- Both paths produce the same logical Blosc2 cache chunks. + +No new Python source modules or mandatory codec packages are needed. Implementation +belongs primarily in `src/blosc2/hdf5_source.py`, with integration in the existing +remote array/store modules. NumPy, Blosc2, and standard-library `zlib` provide +the direct decoding primitives. Protocol-specific fsspec backends remain optional. + +## Scope and explicit limits + +Preserve the public HDF5 URL syntax, dataset selection, immutable-source contract, +attributes, lazy expressions, references, cache policies, and hierarchy browsing. +Preserve direct local h5py access without fsspec or Zarr dependencies. + +Support fixed-size dtypes already accepted by `HDF5NDSource`; reject object and +variable-length data, null dataspaces, and oversized logical chunks clearly. +Check structured dtypes and non-native byte order explicitly, rather than assuming +that `dtype.str` fully describes every accepted dtype. + +Initial layout policy: + +- Chunked, self-contained datasets: native index and the pipeline policy above. +- Contiguous and compact datasets: h5py slicing with bounded logical Blosc2 + chunks, including scalars and zero-length axes. This keeps arbitrary + multidimensional contiguous slicing out of the first range-reader implementation. +- External raw storage and virtual datasets: mark unsupported initially; do not + follow additional files silently. +- Preserve existing link traversal behavior: ordinary object traversal, with + external and soft links omitted and hard-link aliases not promised as separate + leaves. Test the behavior and state it in hierarchy notices. + +Out of scope: writing remote HDF5, changing source mutation semantics, parsing +HDF5 structures ourselves, general legacy HDF5 reference execution, multi-file +aggregation, standalone LZF/Blosc1/bitshuffle decoders, in-memory HDF5 chunk +injection, and block-level reads within Blosc2-filter HDF5 chunks. + +## Current integration points + +| File | Required change | +| --- | --- | +| `src/blosc2/hdf5_source.py` | Replace translation, Zarr-backed arrays, filter monkey-patches, and sync-loop workarounds with indexing, direct decoding, and fallback ownership. | +| `src/blosc2/remote_array.py` | Persist native indexes, share container snapshots between leaves, restore carriers, and handle legacy metadata explicitly. | +| `src/blosc2/remote_store.py` | Discover nodes from native metadata, validate/restore manifests, share the index, and close HDF5 resources correctly. | +| `src/blosc2/remote_store_cache.py` | Adjust manifest compatibility checks only if the HDF5 metadata discriminator cannot be handled in the owner. | +| `src/blosc2/zarr_source.py` | Keep Zarr behavior intact; remove HDF5's reliance on its sync lock and conversion helper. | +| `pyproject.toml` | Remove the legacy HDF5 reference layer from HDF5, development, and test dependencies; retain Zarr for actual Zarr support. | +| Existing HDF5/store/fsspec/b2view tests | Replace legacy reference-layer-specific fixtures and assertions with behavioral checks. | +| HDF5 docs and examples | Document native indexing, fallback filters, dependencies, and legacy reference handling. | + +In particular, `RemoteStore._close_resources()` currently assumes every HDF5 +source has `source.array.store.close()`. Replace this with an explicit HDF5 source +resource lifecycle. `RemoteArray.close()` currently only releases store-derived +handles; do not assume it already closes standalone HDF5 sources. + +Use non-underscored names for helpers imported between modules. Tentative shared +entry points are `scan_hdf5_index()` and `validate_hdf5_index()`; private decoding +and serialization helpers remain in `hdf5_source.py`. + +## 1. Native metadata format and scanner + +Introduce a JSON-compatible, versioned container snapshot with an unambiguous +format discriminator, for example `format: "blosc2-hdf5-index", version: 1`. +Store one source identity and file size at container level; chunk records contain +offsets and lengths, not arbitrary remote URLs. + +Each dataset entry needs: + +- Dataset path, shape, complete dtype description, physical HDF5 layout/chunks, + fill value, and attributes. +- Ordered filter records from the dataset creation property list: filter ID, + flags, and client-data parameters. Do not use private `Dataset._filters`. +- Allocated chunk records: element coordinates, file byte offset, stored byte + length, and per-chunk filter mask. Omitted logical chunks mean unallocated + storage and must produce the HDF5 fill value. +- Enough layout information to select direct reads or fallback at runtime. + Availability of installed plugins is runtime state, not a persisted guarantee. + +Store group entries and isolated unsupported-leaf descriptions as well. Keep +physical metadata independent of caller-selected Blosc2 blocks and compression +parameters so one index serves different array carriers. + +Use tagged encodings for bytes, NumPy scalars/arrays, structured dtype +descriptions, and fill values where ordinary JSON would lose information. Specify +round trips for NaNs, complex values, fixed-size strings, and structured values. +Do not use pickle. Reuse suitable existing serialization helpers after inspection. + +Scanner procedure: + +1. Resolve/reuse the supplied filesystem and storage options. Open a seekable + binary handle with read-ahead disabled, following the existing HTTP convention + `block_size=1, cache_type="none"`. +2. Wrap reads for traffic accounting and open h5py read-only with scoped ownership. +3. Traverse groups/datasets; extract metadata through public h5py APIs. +4. Enumerate allocated chunks using `get_num_chunks()`/`get_chunk_info()` or a + supported efficient iterator. Confirm enumeration is not quadratic on large + indexes; record the minimum h5py version if relying on a newer API. +5. Never decompress dataset payloads while indexing chunked datasets. Attribute + and compact-layout metadata can themselves contain inline values. +6. Close HDF5 before the underlying file object, including all exception paths. + +Scan once per container and share the snapshot, preserving current semantics. +An all-fallback dataset does not need a complete chunk offset table; avoid +enumerating it if its metadata is sufficient for later h5py slicing. + +Validate snapshots before use: version, dataset paths, dtype/shape consistency, +chunk alignment and bounds, duplicate coordinates, nonnegative integer offsets +and sizes (excluding booleans), known filter-mask bits, and ranges within the +recorded file size. Bind the snapshot to the requested source identity and retain +existing credential/persistable-URL rules. Unknown versions must fail clearly. + +## 2. Direct range reader and decoding + +For each requested logical chunk: + +1. Validate the chunk number and translate it into the HDF5 chunk coordinates. +2. If no allocated record exists, construct values from the dataset fill value + without a remote read. Never mistake an absent index table for missing chunks. +3. Fetch exactly `[offset, offset + size)` using the shared filesystem's range + API. Avoid concurrent `seek`/`read` on one mutable file cursor. Honor the + existing concurrency limit and async filesystem ownership conventions. +4. Check for short reads. Decode filters in reverse creation order, skipping a + filter when its corresponding mask bit is set. +5. Validate the decoded byte count and interpret using the stored dtype and + physical chunk shape. Clip to the valid edge extent, then construct the + bounded logical cache chunk with the existing edge-padding convention. +6. Recompress into the caller's Blosc2 chunks/blocks/cparams and return the chunk + bytes through the existing source protocol. + +Decoder details: + +- Deflate: HDF5's zlib stream through `zlib`, with bounded expected output and + explicit truncated/corrupt-stream errors. +- Shuffle: reverse the HDF5 byte-plane shuffle using its element-size parameter; + handle itemsize 1 and validate buffer divisibility. This is separate from the + internal shuffle/bitshuffle used by a Blosc2 frame. +- Blosc2: preserve support for both ordinary compressed buffers and super-chunk + frames. Normalize `SChunk` bytes and any `NDArray` frame results correctly; + use logical values for shaped frames rather than assuming their internal block + ordering is NumPy C order. Test with real `hdf5plugin.Blosc2` outputs as well as + manually written raw chunks. + +Do not fall back silently after a known direct decoder reports corrupt data or +an inconsistent index. Fallback is a planned capability decision, not recovery +from arbitrary exceptions. Include dataset and chunk coordinates in failures. + +Keep conversion in `hdf5_source.py` initially. It must not enter `ZARR_SYNC_LOCK`. +Do not return the original Blosc2 frame as a cache chunk without conversion: +physical framing and requested logical cache geometry may differ. + +## 3. h5py fallback and resource lifecycle + +Open the fallback HDF5 file lazily on the first uncached read that needs it. +Use h5py slicing for the requested logical chunk; load `hdf5plugin` opportunistically +so installed plugin filters register with HDF5. Report unavailable filter IDs +and installation guidance without excluding readable siblings from a store. + +The same conversion code should normalize both direct and fallback values. +Cache hits, including persisted warm chunks, must not open the fallback file. +Fallback opening may reread HDF5 metadata even when the native snapshot is warm; +the snapshot is not an HDF5 metadata-cache image. + +Ownership requirements: + +- Standalone source: own the fallback handle, provide idempotent source cleanup, + and retain a finalizer for abandoned sources. Define how explicit standalone + `RemoteArray.close()` reaches this cleanup without changing unrelated sources. +- Store: share one lazily initialized fallback reader per container owner where + practical, so sibling leaves reuse metadata. Leaf closure must not close a + reader still used by another leaf. Owner close/refresh releases it once. +- Protect lazy initialization and close-versus-read races. Close h5py first, + then its file object, and only close filesystem sessions owned by this reader. +- Preserve generation checks and existing stale/closed-handle errors. + +h5py holds its global HDF5 lock during fallback network reads and decoding. +Concurrent fallback reads are therefore serialized; `max_concurrency` does not +promise parallel h5py I/O. Direct range requests and standalone decoders must +remain outside that lock and outside any new container-wide I/O lock. Reuse the +current fetch scheduling rather than adding a second thread pool per source. + +Measure traffic at one I/O boundary per path. With read-ahead disabled, h5py's +file read callbacks can be charged directly. Avoid counting an already counted +filesystem read again. Keep the established exclusion of non-payload `info`/HEAD +calls, and distinguish transport counts from h5py callback counts in benchmarks. + +## 4. Persistence, hierarchy integration, and compatibility + +Persist the native snapshot in carriers and in the shared container cache so +reopening a direct dataset needs no h5py rescan. Prefer a new `hdf5-index` vlmeta +key and `.hdf5-index.b2` sidecar rather than disguising the new format as a +legacy HDF5 reference dictionary. Update saving, loading, cloning, cache export, +lazy-expression reopening, and shared-index publication together. + +`RemoteStore` should build its node tree directly from snapshot groups and +datasets, replacing `.zarray`/`.zgroup` parsing. Keep one snapshot in the owner +and pass it to leaves. Restore and validate that format in manifests; update +the user-facing hierarchy notice and b2view integration. + +Compatibility policy for this implementation: + +- Keep the existing `hdf5_index=` argument as an input slot for a native snapshot or + its JSON path, with its new accepted format documented. Do not add another + public keyword merely to rename it in this change. +- Explicit legacy HDF5 reference dictionaries/JSON files: detect and reject with a + precise migration message explaining how to omit `hdf5_index` and rescan the source. + Do not silently reinterpret arbitrary multi-file reference maps or fetch their + targets. This is a documented compatibility break. +- Implicit legacy snapshots found in reusable array caches: rescan the original + HDF5 source and check geometry, dtype, and conversion identity before retaining + any warm chunks. A failed rescan or mismatch must not relabel old cached data + as valid. Document that this migration requires source access. +- Legacy store manifests: report the established incompatible-cache error and + advise a new cache directory; do not silently delete or rewrite generations. +- Version native metadata independently of logical chunk encoding. Retain the + existing encoding stamp only where conversion semantics and geometry are + demonstrably identical; otherwise bump it and reject incompatible cache reuse. + Correct structured-dtype identity at the same boundary if needed. + +Retire the old internal scanner after updating callers. The public +`scan_hdf5_index()` helper returns the native format. `available_datasets()` must +accept native indexes and retain direct local-file discovery. + +## 5. Dependencies and removal of workarounds + +- Remove the legacy HDF5 reference layer from the `hdf5` extra and development/test dependency groups. +- Retain `h5py` and the existing convenience installation of `hdf5plugin` in the + HDF5 extra; runtime plugin use remains optional. Keep fsspec in its existing + extra, so remote installation remains `blosc2[hdf5,fsspec]`. +- Ensure CI installs `hdf5plugin` for actual plugin-filter tests, rather than + allowing all such tests to skip. Keep platform markers consistent with h5py. +- Remove HDF5's numcodecs registration, legacy HDF5 filter monkey-patch, Zarr store adapter, + recursive Buffer conversion, scan timeout retries, and global Zarr-loop reset. +- Keep Zarr dependencies and sync machinery needed by genuine Zarr readers. +- Test HDF5 functionality with the legacy HDF5 reference layer, Zarr, and numcodecs imports blocked; + test package import without any HDF5 optional dependencies. + +Update `doc/guides/remote_arrays.md`, `doc/reference/remotearray.rst`, +`doc/getting_started/installation.rst`, `doc/guides/b2view.rst`, HDF5 source +docstrings, and relevant examples. Explain pipeline selection, warm snapshots, +fallback locking, and legacy migration. Do not describe installing `hdf5plugin` +as supplying standalone decoders. + +## 6. Verification matrix + +Extend existing test modules; no new source modules are required. + +### Deterministic offline correctness + +Use generated HDF5 fixtures on `memory://` and the existing local HTTP range +server. Compare complete results and slices against NumPy/h5py ground truth. + +- Direct pipelines: no filters, shuffle only, deflate only, shuffle + deflate, + Blosc2, and valid combined pipelines. Include big-endian and fixed-size string + dtypes, supported compound dtypes, and itemsize 1. +- Per-chunk masks: use raw chunk writes to construct chunks that skip individual + filters; validate order and masks independently of high-level writer defaults. +- Sparse allocation: missing chunks, custom/nonzero/NaN fill values, fully + unallocated datasets, partial edge chunks, and empty axes. +- Local and fallback layouts: scalars, contiguous arrays, compact storage, and + bounded conversion of a large contiguous array. +- Blosc2: ordinary compressed buffers, SChunk frames, shaped frames where + supported, and plugin-generated multidimensional edge chunks. +- Fallback: LZF where available, Fletcher32 pipelines, scale-offset, and at least + one `hdf5plugin` filter outside the direct set (such as Blosc1 or bitshuffle). + Verify every requested chunk uses h5py and the handle is reused. +- Failures: unsupported/null/object datasets, missing plugins, corrupt streams, + short ranges, invalid chunk metadata, and malformed/unknown-version snapshots. + +### Cache, ownership, and concurrency + +- Native snapshot round trips preserve dtype, fill, attrs, filters, and offsets. +- Cold scan once per container; opening sibling leaves does not rescan. +- Persisted direct snapshots reopen without scanning; warm cached chunks require + no data requests. An uncached allocated direct chunk makes one range fetch. +- Missing chunks make no range fetch. Fallback metadata reads occur only when + needed, and repeated fallback access reuses the handle. +- Assert overlapping direct requests using synchronization barriers in a fake + range backend, rather than brittle wall-clock thresholds. Verify the fetch + cap and concurrent correctness without involving h5py in direct decoding. +- Source/store close, partial initialization failures, refresh, stale handles, + finalizers, sibling ownership, and externally owned filesystem lifetimes. +- Native/legacy carrier cases, explicit legacy indexes, store manifests, shared + sidecars, saved expressions, and cache export/reopen. +- Keep HTTP range behavior and moto/S3 coverage: a memory filesystem alone does + not exercise fsspec's real async transport/session ownership. +- Replace tests asserting HDF5 reference translation counts with scanner counts, and + remove tests whose only purpose was resetting Zarr resources for HDF5 scans. + +Relevant suites: `tests/test_hdf5_source.py`, `tests/test_remote_store.py`, +`tests/test_fsspec.py`, `tests/test_fsspec_s3.py`, `tests/test_remote_array.py`, +and `tests/b2view/test_hierarchy.py`. Also run Zarr regression tests because the +old conversion helper and resource assumptions are shared. + +### Public HTTPS fixture + +Use the credential-free URL supplied for this work: + +`https://f001.backblazeb2.com/file/blosc2/hierarchy.h5` + +Metadata inspected on 2026-09-17: file size 1,037,240 bytes. The groups `d0`, +`d0/d1`, and `d0/d1/d2` each contain: + +| Leaf | Shape | dtype | HDF5 chunks | Filters | +| --- | --- | --- | --- | --- | +| `a0` | `()` | int32 | contiguous | none | +| `a1` | `(10000,)` | float32 | `(2500,)` | Blosc2, ID 32026 | +| `a2` | `(1000, 1000)` | int32 | `(500, 500)` | Blosc2, ID 32026 | +| `a3` | `(10, 1000, 1000)` | int32 | `(2, 500, 500)` | Blosc2, ID 32026 | + +Verified h5py sample reads: `d0/a0` is 0, and both `d0/d1/a2[:3, :3]` +and `d0/d1/d2/a3[0, :3, :3]` return rows `[0, 1, 2]`, +`[1000, 1001, 1002]`, and `[2000, 2001, 2002]`. The first raw chunk of +`d0/d1/a2`, fetched separately by byte range, opens with `blosc2.from_cframe()` +as an `NDArray` of shape `(500, 500)`. This is a concrete shaped-frame decoder +case, not just a hypothetical extension. These checks validate the fixture and +primitives; they do not constitute a test of the proposed reader. + +Existing S3 tests assume `d0/d1/a2` is 3-D. Do not copy that assumption to the +HTTPS tests, or assume the two endpoints necessarily serve identical versions. + +Add separately marked network tests for hierarchy discovery, scalar reads, +1-D/2-D/3-D slices, a slice crossing chunk boundaries, persistent reopening, +and no extra traffic on a repeated cached slice. Use h5py reference reads or a +temporary local copy for value comparison; exclude reference traffic from the +adapter's measurements. Generated fixtures remain the authoritative coverage +for gzip/shuffle and plugin fallbacks, absent from this public file. + +Keep normal CI independent of this endpoint. Live failures should distinguish +transport/fixture changes from incorrect values, rather than skip any exception. + +### Commands and performance evidence + +Run all Python/test/build commands in the `blosc2` conda environment. If the +conda launcher fails due to its installed solver plugin, use that environment's +Python executable directly; do not substitute the base interpreter. + +After implementation, run targeted suites first, Ruff on changed Python files, +then the default suite (warnings are errors). Run public tests explicitly, e.g.: + +```sh +conda run -n blosc2 pytest tests/test_hdf5_source.py tests/test_remote_store.py tests/test_fsspec.py tests/b2view/test_hierarchy.py +conda run -n blosc2 pytest -m network tests/test_hdf5_source.py -k https +conda run -n blosc2 pytest +``` + +Record cold discovery, first chunk, adjacent/distant uncached chunks, cached +rereads, and persisted reopening separately. Report elapsed time, payload bytes, +data requests, and index size. Compare direct and fallback paths with identical +values/chunk geometry and several concurrency settings. Do not turn external +network timing into a CI pass/fail threshold. + +## Implementation sequence and completion criteria + +1. **Index contract and fixtures:** implement serialization/validation and scanner; + establish dtype, fill, layout, and filter-mask tests before replacing callers. +2. **Direct reader:** implement exact range reads, three filter decoders, and + cache-chunk conversion; prove no h5py calls during direct chunk retrieval. +3. **Fallback and lifecycle:** add lazy retained h5py readers, plugin errors, + resource ownership, and fallback layout handling. +4. **Integration and persistence:** switch HDF5NDSource, RemoteArray, RemoteStore, + discovery, manifests, and b2view; implement the stated compatibility policy. +5. **Dependency cleanup and documentation:** remove legacy HDF5 reference/Zarr machinery, + update packaging and examples, and verify missing-dependency behavior. +6. **Validation:** run offline suites, live HTTPS tests, and request/concurrency + measurements; resolve regressions before declaring the replacement complete. + +Completion means HDF5 reads and browsing work without the legacy HDF5 reference layer, Zarr, or numcodecs; +all agreed direct filters use indexed range reads; other supported HDF5 filters +work through the retained h5py fallback; native indexes survive cache round trips; +warm cached reads avoid remote I/O; ownership and legacy errors are tested; and +the public fixture and generated test matrix pass. The first release does not +require optimizing contiguous reads or parallelizing the fallback. diff --git a/pyproject.toml b/pyproject.toml index 124d37930..25c4c8ae1 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -52,7 +52,7 @@ documentation = "https://www.blosc.org/python-blosc2/python-blosc2.html" [project.optional-dependencies] parquet = ["pyarrow"] zarr = ["zarr>=3.0.9"] -hdf5 = ["kerchunk", "h5py", "hdf5plugin"] +hdf5 = ["h5py", "hdf5plugin"] # The b2view terminal viewer (the `b2view` script) is opt-in: most users want # blosc2 only as a compression library, and the TUI stack has no use under # wasm32 (no TTY). Install with `pip install "blosc2[tui]"`. This also pulls @@ -79,7 +79,6 @@ dev = [ "h5py", "hdf5plugin", "jupyterlab", - "kerchunk", "matplotlib", "pandas", "plotly", @@ -100,8 +99,10 @@ test = [ # Exercise the optional remote array sources instead of silently skipping # tests/test_zarr_source.py and tests/test_hdf5_source.py. "zarr>=3.0.9; platform_machine != 'wasm32'", - "kerchunk; platform_machine != 'wasm32'", "h5py; platform_machine != 'wasm32'", + # tests/test_hdf5_source.py exercises the optional Blosc2 HDF5 filter through + # it; without the plugin those paths skip instead of running in the default job. + "hdf5plugin; platform_machine != 'wasm32'", # tests/test_fsspec_s3.py needs a real S3 endpoint (moto, served locally, so # still offline) and a real *async* backend (s3fs). memory:// is neither, and # cannot see the class of bug that lives there: awaiting an async filesystem's diff --git a/src/blosc2/hdf5_source.py b/src/blosc2/hdf5_source.py index 89620c649..611522050 100644 --- a/src/blosc2/hdf5_source.py +++ b/src/blosc2/hdf5_source.py @@ -4,11 +4,11 @@ # # SPDX-License-Identifier: BSD-3-Clause ####################################################################### - -"""An immutable HDF5 source using h5py locally and kerchunk for remote files.""" +"""Immutable local and remote HDF5 sources backed by h5py and fsspec.""" from __future__ import annotations +import base64 import contextlib import hashlib import io @@ -17,158 +17,125 @@ import os import threading import weakref +import zlib from urllib.parse import urlsplit import numpy as np import blosc2 from blosc2.proxy_source import REMOTE_MAX_CONCURRENCY, ProxyNDSource, Traffic -from blosc2.zarr_source import ZARR_SYNC_LOCK, counting_store, zarr_chunk_to_blosc2 -_HDF5_SCAN_LOCK = threading.Lock() +HDF5_INDEX_FORMAT = "blosc2-hdf5-index" +HDF5_INDEX_VERSION = 1 +_DIRECT_FILTERS = {1, 2, 32026} # deflate, shuffle, Blosc2 def check_hdf5_dependencies() -> None: - """Validate that kerchunk and h5py are available.""" + """Validate the dependencies needed for remote HDF5 access.""" + _check_h5py_dependencies() try: - import kerchunk.hdf # noqa: F401 + import fsspec # noqa: F401 except ImportError as exc: raise ImportError( - "HDF5 support requires kerchunk; install it with 'pip install blosc2[hdf5]'" + "Remote HDF5 support requires fsspec; install it with 'pip install blosc2[fsspec]'" ) from exc - _check_h5py_dependencies() - _ensure_blosc2_filter_registered() - def _check_h5py_dependencies() -> None: try: import h5py # noqa: F401 except ImportError as exc: raise ImportError("HDF5 support requires h5py; install it with 'pip install blosc2[hdf5]'") from exc - with contextlib.suppress(ImportError): - import hdf5plugin # noqa: F401 - - -def _register_numcodecs_blosc2() -> None: - try: - import numcodecs - import numcodecs.abc - except ImportError: - return - - if hasattr(numcodecs, "Blosc2"): - return - - class Blosc2Codec(numcodecs.abc.Codec): - codec_id = "blosc2" - - def __init__(self, **kwargs): - self.kwargs = kwargs - - def encode(self, buf): - return buf - - def decode(self, buf, out=None): - try: - decomp = blosc2.decompress(buf) - except Exception: - # A Blosc2-filter chunk is a whole super-chunk frame, so - # from_cframe() may hand back an SChunk whose slice is bytes. - decomp = bytes(blosc2.from_cframe(buf)[:]) - if out is not None: - np.frombuffer(out, dtype=np.uint8)[:] = np.frombuffer(decomp, dtype=np.uint8) - return out - return decomp - - def get_config(self): - return {"id": self.codec_id, **self.kwargs} - - numcodecs.register_codec(Blosc2Codec) - numcodecs.Blosc2 = Blosc2Codec - - -def _patch_kerchunk_decode_filters() -> None: - try: - import kerchunk.hdf - import numcodecs - except ImportError: - return - - if getattr(kerchunk.hdf.SingleHdf5ToZarr, "_blosc2_patched", False): - return - - orig_decode = kerchunk.hdf.SingleHdf5ToZarr._decode_filters - - def patched_decode(self, h5obj): - filters = [] - saved = {} - for filter_id, props in list(h5obj._filters.items()): - if str(filter_id) == "32026": - saved[filter_id] = props - filters.append(numcodecs.Blosc2()) - for fid in saved: - h5obj._filters.pop(fid, None) - try: - filters.extend(orig_decode(self, h5obj)) - finally: - h5obj._filters.update(saved) - return filters - - kerchunk.hdf.SingleHdf5ToZarr._decode_filters = patched_decode - kerchunk.hdf.SingleHdf5ToZarr._blosc2_patched = True - - -def _ensure_blosc2_filter_registered() -> None: - _register_numcodecs_blosc2() - _patch_kerchunk_decode_filters() - - -def check_zarr_fsspec_dependencies() -> None: - """Validate that zarr and fsspec are available.""" - try: - import zarr # noqa: F401 - except ImportError as exc: - raise ImportError( - "HDF5NDSource requires Zarr-Python; install it with 'pip install blosc2[zarr]'" - ) from exc - - try: - import fsspec # noqa: F401 - except ImportError as exc: - raise ImportError( - "HDF5NDSource requires fsspec; install it with 'pip install blosc2[fsspec]'" - ) from exc + import hdf5plugin # noqa: F401 # registers optional filters with HDF5 + + +def dtype_value(dtype): + dtype = np.dtype(dtype) + return {"descr": dtype.descr} if dtype.fields is not None else {"str": dtype.str} + + +def dtype_from_value(value): + if "str" in value: + return np.dtype(value["str"]) + return np.dtype([_dtype_field_from_json(field) for field in value["descr"]]) + + +def _dtype_field_from_json(field): + name, spec, *shape = field + if isinstance(name, list): + # dtype.descr writes a titled field as (title, name); JSON and msgpack + # round trips turn that tuple into a list again. + name = tuple(name) + if isinstance(spec, list): + spec = [_dtype_field_from_json(item) for item in spec] + return (name, spec, tuple(shape[0])) if shape else (name, spec) + + +def _json_value(value): + """Encode HDF5 metadata without pickle or lossy byte coercion.""" + if isinstance(value, np.ndarray): + if value.dtype.hasobject: + # Object arrays box Python references, so tobytes() would persist + # pointers instead of the elements they name. + return { + "__object_ndarray__": [_json_value(item) for item in value.ravel().tolist()], + "shape": list(value.shape), + } + return { + "__ndarray__": base64.b64encode(value.tobytes()).decode(), + "dtype": dtype_value(value.dtype), + "shape": list(value.shape), + } + if isinstance(value, np.generic): + return {"__scalar__": base64.b64encode(value.tobytes()).decode(), "dtype": dtype_value(value.dtype)} + if isinstance(value, bytes): + return {"__bytes__": base64.b64encode(value).decode()} + if isinstance(value, float) and not math.isfinite(value): + return {"__float__": repr(value)} + if isinstance(value, (str, int, float, bool)) or value is None: + return value + if isinstance(value, (list, tuple)): + return [_json_value(item) for item in value] + if isinstance(value, dict): + return {str(key): _json_value(item) for key, item in value.items()} + return str(value) -def _reset_zarr_sync_resources() -> None: - """Stop and detach Zarr's process-global synchronous event-loop resources.""" - from zarr.core import sync - - with ZARR_SYNC_LOCK: - loop = sync.loop[0] - thread = sync.iothread[0] - executor = getattr(sync, "_executor", None) - sync.loop[0] = None - sync.iothread[0] = None - if hasattr(sync, "_executor"): - sync._executor = None - if loop is not None: - if loop.is_running(): - with contextlib.suppress(RuntimeError): - loop.call_soon_threadsafe(loop.stop) - if thread is not None: - thread.join(timeout=0.2) - if not thread or not thread.is_alive(): - with contextlib.suppress(RuntimeError): - loop.close() - if executor is not None: - executor.shutdown(wait=False, cancel_futures=True) +def _from_json_value(value): + if isinstance(value, list): + return [_from_json_value(item) for item in value] + if not isinstance(value, dict): + return value + if "__bytes__" in value: + return base64.b64decode(value["__bytes__"]) + if "__float__" in value: + return float(value["__float__"]) + if "__scalar__" in value: + return np.frombuffer(base64.b64decode(value["__scalar__"]), dtype=dtype_from_value(value["dtype"]))[ + 0 + ] + if "__object_ndarray__" in value: + items = [_from_json_value(item) for item in value["__object_ndarray__"]] + result = np.empty(value["shape"], dtype=object) + flat = result.reshape(-1) + for index, item in enumerate(items): + flat[index] = item + return result + if "__ndarray__" in value: + return np.frombuffer( + base64.b64decode(value["__ndarray__"]), dtype=dtype_from_value(value["dtype"]) + ).reshape(value["shape"]) + return {key: _from_json_value(item) for key, item in value.items()} + + +def decode_hdf5_value(value): + """Decode a JSON-compatible value stored in a native HDF5 index.""" + return _from_json_value(value) class _CountingFile(io.IOBase): - """Count h5py's read/readinto calls without fsspec read-ahead.""" + """Count h5py file-object reads without adding read-ahead.""" def __init__(self, file, traffic): self.file, self.traffic = file, traffic @@ -190,193 +157,342 @@ def tell(self): return self.file.tell() -def scan_hdf5_refs(urlpath, storage_options=None, *, unsupported=None, traffic=None, _filesystem=None): - """Translate once, closing owned handles and optionally isolating bad leaves. - - Disable fsspec read-ahead: h5py requests metadata ranges itself. Translation - may still read small inline values and scales with the file's chunk index. - """ - check_hdf5_dependencies() +def _filesystem_and_path(urlpath, storage_options=None, filesystem=None): import fsspec - import zarr - fs, path = ( - fsspec.core.url_to_fs(urlpath, **(storage_options or {})) - if _filesystem is None - else (_filesystem, _filesystem._strip_protocol(urlpath)) - ) - # kerchunk builds a zarr hierarchy through zarr's sync bridge, which runs - # its event loop in another thread and waits with no timeout by default - # (async.timeout is None). A lost wake-up there hung CI for hours, on - # every platform, inside this translation. Bound it -- and if it fires, - # reset the process-global loop before the first attempt and after a timeout - # so a previous failed translation cannot poison this one. - with _HDF5_SCAN_LOCK: - _reset_zarr_sync_resources() - for attempt in range(2): - try: - with zarr.config.set({"async.timeout": 120}): - return _translate_hdf5(fs, path, urlpath, unsupported, traffic) - except TimeoutError: - if attempt: - raise - _reset_zarr_sync_resources() - raise TimeoutError("HDF5 translation timed out twice") # pragma: no cover -- retry re-raises - - -def _plain_hdf5_refs(value): - """Make zarr's Buffer objects safe to keep in a manifest. - - kerchunk writes the zarr hierarchy through MemoryStore, which holds Buffer - objects; when a translation is cut short (zarr's sync bridge timing out) - some of those buffers survive in the returned references, and the manifest - cannot serialize them. Buffer bytes are what a reference stores anyway. - """ - if isinstance(value, dict): - return {key: _plain_hdf5_refs(item) for key, item in value.items()} - if isinstance(value, list): - return [_plain_hdf5_refs(item) for item in value] - to_bytes = getattr(value, "to_bytes", None) - if type(value).__module__.startswith("zarr.") and callable(to_bytes): - return bytes(to_bytes()) - return value - - -def _translate_hdf5(fs, path, urlpath, unsupported, traffic): - import kerchunk.hdf - - # HTTP uses block_size=0 for non-seekable streaming; cache_type disables read-ahead. - with fs.open(path, "rb", block_size=1, cache_type="none") as file: - if traffic is not None: - file = _CountingFile(file, traffic) - translator = kerchunk.hdf.SingleHdf5ToZarr( - file, url=urlpath, error="raise" if unsupported is not None else "warn" + if filesystem is not None: + return filesystem, filesystem._strip_protocol(urlpath) + # Owned filesystems must never be the process-wide fsspec instance: closing + # one source's session must not invalidate another source for the same URL. + # Force privacy even if the caller passed skip_instance_cache=False, since + # close() treats the filesystem as owned. + options = dict(storage_options or {}) + options["skip_instance_cache"] = True + return fsspec.core.url_to_fs(urlpath, **options) + + +def _close_owned_filesystem(filesystem): + """Release the async session of an fsspec filesystem this code created.""" + close = getattr(filesystem, "close_session", None) + session = getattr(filesystem, "_s3creator", None) or getattr(filesystem, "_session", None) + if close is not None and session is not None: + close(filesystem.loop, session) + + +def _dataset_metadata(dataset): + dcpl = dataset.id.get_create_plist() + filters = [] + for index in range(dcpl.get_nfilters()): + filter_id, flags, values, name = dcpl.get_filter(index) + filters.append( + { + "id": int(filter_id), + "flags": int(flags), + "values": [int(v) for v in values], + "name": bytes(name).decode(errors="replace"), + } ) - if unsupported is not None: - translate_node = translator._translator + chunks = None if dataset.chunks is None else [int(v) for v in dataset.chunks] + direct = chunks is not None and all(item["id"] in _DIRECT_FILTERS for item in filters) + allocated = [] + if direct: + for index in range(dataset.id.get_num_chunks()): + info = dataset.id.get_chunk_info(index) + allocated.append( + { + "offset": [int(v) for v in info.chunk_offset], + "filter_mask": int(info.filter_mask), + "byte_offset": int(info.byte_offset), + "size": int(info.size), + } + ) + return { + "shape": [int(v) for v in dataset.shape], + "dtype": dtype_value(dataset.dtype), + "chunks": chunks, + "fill_value": _json_value(dataset.fillvalue), + "attrs": {key: _json_value(value) for key, value in dataset.attrs.items()}, + "filters": filters, + "direct": direct, + "allocated": allocated, + } + + +def scan_hdf5_index(urlpath, storage_options=None, *, unsupported=None, traffic=None, _filesystem=None): + """Build a versioned native index for one local or remote HDF5 container.""" + import h5py + + urlpath = blosc2.core.normalize_urlpath(os.fspath(urlpath)) + local = _filesystem is None and (not urlsplit(urlpath).scheme or os.path.isabs(urlpath)) + if local: + _check_h5py_dependencies() + fs = None + path = urlpath + else: + check_hdf5_dependencies() + fs, path = _filesystem_and_path(urlpath, storage_options, _filesystem) + groups, datasets = {"": {"attrs": {}}}, {} + try: + with contextlib.ExitStack() as stack: + if local: + raw = stack.enter_context(open(path, "rb")) + else: + raw = stack.enter_context(fs.open(path, "rb", block_size=1, cache_type="none")) + fileobj = _CountingFile(raw, traffic) if traffic is not None else raw + with h5py.File(fileobj, "r") as h5file: + groups[""]["attrs"] = {key: _json_value(value) for key, value in h5file.attrs.items()} + + def visit(name, obj): + try: + if isinstance(obj, h5py.Group): + groups[name] = { + "attrs": {key: _json_value(value) for key, value in obj.attrs.items()} + } + elif isinstance(obj, h5py.Dataset): + if obj.is_virtual: + raise TypeError("HDF5 virtual datasets are not supported") + if obj.external: + raise TypeError("HDF5 externally stored datasets are not supported") + if obj.shape is None: + raise TypeError("HDF5 null datasets are not supported") + dtype = np.dtype(obj.dtype) + if dtype.hasobject or dtype.itemsize == 0: + raise TypeError(f"HDF5NDSource only supports fixed-size dtypes, got {dtype}") + datasets[name] = _dataset_metadata(obj) + except Exception as exc: + if unsupported is None: + raise + unsupported[name] = f"{type(exc).__name__}: {exc}" + + h5file.visititems(visit) + with contextlib.suppress(Exception): + size = os.path.getsize(path) if local else int(fs.info(path)["size"]) + finally: + if fs is not None and _filesystem is None: + _close_owned_filesystem(fs) + if "size" not in locals(): + size = None + return { + "format": HDF5_INDEX_FORMAT, + "version": HDF5_INDEX_VERSION, + "urlpath": os.fspath(urlpath), + "size": size, + "groups": groups, + "datasets": datasets, + } + + +def validate_hdf5_index(index, urlpath=None): + """Validate and return a native HDF5 index.""" + if not isinstance(index, dict): + raise ValueError("Invalid HDF5 index") + if index.get("format") != HDF5_INDEX_FORMAT: + legacy_keys = {".zarray", ".zgroup"} + if "refs" in index or any( + str(key) in legacy_keys or str(key).endswith(("/.zarray", "/.zgroup")) for key in index + ): + raise ValueError( + "Legacy HDF5 reference maps are unsupported; omit hdf5_index and rescan the source" + ) + raise ValueError("Invalid HDF5 index format") + if index.get("version") != HDF5_INDEX_VERSION: + raise ValueError(f"Unsupported HDF5 index version {index.get('version')!r}") + if urlpath is not None and index.get("urlpath") != os.fspath(urlpath): + raise ValueError("HDF5 index specification does not match the requested URL") + if not isinstance(index.get("groups"), dict) or not isinstance(index.get("datasets"), dict): + raise ValueError("Invalid HDF5 index contents") + size = index.get("size") + for path, meta in index["datasets"].items(): + _validate_dataset_entry(path, meta, size) + return index + + +def _validate_dataset_entry(path, meta, file_size): + """Validate one dataset entry in a native index.""" + if not isinstance(path, str) or not isinstance(meta, dict): + raise ValueError("Invalid HDF5 dataset entry") + required = {"shape", "dtype", "chunks", "fill_value", "attrs", "filters", "direct", "allocated"} + if not required.issubset(meta): + missing = sorted(required - set(meta)) + raise ValueError(f"Incomplete HDF5 dataset entry for {path!r}: missing {missing}") + if not isinstance(meta["attrs"], dict): + raise ValueError(f"Invalid HDF5 attributes for {path!r}") + shape, chunks = tuple(meta["shape"]), meta["chunks"] + dtype = dtype_from_value(meta["dtype"]) + if len(shape) > blosc2.MAX_DIM or any( + isinstance(v, bool) or not isinstance(v, int) or v < 0 for v in shape + ): + raise ValueError(f"Invalid HDF5 shape for {path!r}") + if chunks is not None and ( + not isinstance(chunks, list) + or len(chunks) != len(shape) + or any(isinstance(v, bool) or not isinstance(v, int) or v <= 0 for v in chunks) + ): + raise ValueError(f"Invalid HDF5 chunks for {path!r}") + if dtype.hasobject or dtype.itemsize == 0: + raise ValueError(f"Invalid HDF5 dtype for {path!r}") + filters = meta.get("filters") + if not isinstance(filters, list) or any( + not isinstance(item, dict) or isinstance(item.get("id"), bool) or not isinstance(item.get("id"), int) + for item in filters + ): + raise ValueError(f"Invalid HDF5 filters for {path!r}") + if not isinstance(meta.get("direct"), bool): + raise ValueError(f"Invalid HDF5 read mode for {path!r}") + if meta["direct"]: + _validate_direct_filters(path, chunks, filters) + allocated = meta.get("allocated") + if not isinstance(allocated, list): + raise ValueError(f"Invalid HDF5 allocation table for {path!r}") + _validate_allocated_records(path, allocated, shape, chunks, filters, file_size) + + +def _validate_direct_filters(path, chunks, filters): + """Validate the pipeline a direct chunk reader will decode.""" + if chunks is None or any(item["id"] not in _DIRECT_FILTERS for item in filters): + raise ValueError(f"Invalid direct HDF5 filter pipeline for {path!r}") + for item in filters: + values = item.get("values") + if not isinstance(values, list) or any( + isinstance(value, bool) or not isinstance(value, int) for value in values + ): + raise ValueError(f"Invalid HDF5 filter values for {path!r}") + # Shuffle records its element size as the only client value; a bogus + # value would make the decoder skip unscrambling and return wrong data. + if item["id"] == 2 and (len(values) != 1 or values[0] <= 0): + raise ValueError(f"Invalid HDF5 shuffle filter for {path!r}") + + +def _validate_allocated_records(path, allocated, shape, chunks, filters, file_size): + """Validate the chunk byte ranges recorded in a native index.""" + seen = set() + for record in allocated: + coord = tuple(record.get("offset", ())) + byte_offset, length = record.get("byte_offset"), record.get("size") + if coord in seen or len(coord) != len(shape): + raise ValueError(f"Invalid HDF5 chunk coordinates for {path!r}") + seen.add(coord) + if chunks is None or any( + value % chunk or value >= extent + for value, chunk, extent in zip(coord, chunks, shape, strict=True) + ): + raise ValueError(f"Misaligned HDF5 chunk coordinates for {path!r}") + mask = record.get("filter_mask") + if isinstance(mask, bool) or not isinstance(mask, int) or mask < 0 or mask >> len(filters): + raise ValueError(f"Invalid HDF5 filter mask for {path!r}") + if any( + isinstance(v, bool) or not isinstance(v, int) or v < 0 for v in (*coord, byte_offset, length) + ): + raise ValueError(f"Invalid HDF5 chunk range for {path!r}") + if file_size is not None and byte_offset + length > file_size: + raise ValueError(f"HDF5 chunk range exceeds the file for {path!r}") + + +# Kept while callers migrate from the old internal name. +def available_datasets(url, storage_options: dict | None = None) -> list[str]: + """Return all dataset paths in an HDF5 file or native index.""" + if isinstance(url, dict): + return sorted(validate_hdf5_index(url)["datasets"]) + if not isinstance(url, (str, os.PathLike)): + raise TypeError("url must be a URL string, path-like object, or HDF5 index") + url_str = blosc2.core.normalize_urlpath(os.fspath(url)) + if "::" in url_str and "://" not in url_str.split("::", 1)[1]: + url_str = url_str.split("::", 1)[0].rstrip("/") + lower = url_str.lower() + for ext in (".h5/", ".hdf5/"): + index = lower.find(ext) + if index != -1: + url_str = url_str[: index + len(ext) - 1] + break + if url_str.endswith(".json"): + if not urlsplit(url_str).scheme or os.path.isabs(url_str): + with open(url_str) as file: + return sorted(validate_hdf5_index(json.load(file))["datasets"]) + import fsspec - def isolated_node(name, obj): - try: - return translate_node(name, obj) - except TimeoutError: - # Not this node's fault: zarr's sync bridge is wedged, and - # swallowing it would just wedge again on the next call. - # Let scan_hdf5_refs reset the loop and start over. - raise - except Exception as exc: - unsupported[name] = f"{type(exc).__name__}: {exc}" - return None - - translator._translator = isolated_node - try: - return _plain_hdf5_refs(translator.translate()) - finally: - translator.close() + with fsspec.open(url_str, "r", **(storage_options or {})) as file: + return sorted(validate_hdf5_index(json.load(file))["datasets"]) + if not urlsplit(url_str).scheme or os.path.isabs(url_str): + _check_h5py_dependencies() + import h5py + datasets = [] + with h5py.File(url_str, "r") as file: + file.visititems( + lambda name, obj: datasets.append(name) if isinstance(obj, h5py.Dataset) else None + ) + return sorted(datasets) + return sorted(scan_hdf5_index(url_str, storage_options)["datasets"]) + + +def _selection(nchunk, shape, chunks): + grid = tuple(math.ceil(size / chunk) for size, chunk in zip(shape, chunks, strict=True)) + total = math.prod(grid) + if isinstance(nchunk, bool) or not isinstance(nchunk, int) or nchunk < 0 or nchunk >= total: + raise IndexError(f"nchunk must be in range [0, {total}), got {nchunk}") + coords = np.unravel_index(nchunk, grid) + offsets = tuple(int(coord) * chunk for coord, chunk in zip(coords, chunks, strict=True)) + selection = tuple( + slice(offset, min(offset + chunk, size)) + for offset, chunk, size in zip(offsets, chunks, shape, strict=True) + ) + return offsets, selection -def available_datasets(url, storage_options: dict | None = None) -> list[str]: - """Return all dataset paths within an HDF5 file or reference dictionary. - - Parameters - ---------- - url : str, os.PathLike, or dict - Path or URL to an HDF5 file, a JSON reference file, or an in-memory - kerchunk reference dictionary. - storage_options : dict, optional - Options passed to fsspec or kerchunk for remote URLs. - - Returns - ------- - list[str] - Sorted list of dataset paths (e.g. ``['d0/a0', 'd0/d1/a2']``). - """ - if isinstance(url, dict): - ref_dict = url.get("refs", url) - elif isinstance(url, (str, os.PathLike)): - url_str = blosc2.core.normalize_urlpath(os.fspath(url)) - if "::" in url_str: - parts = url_str.split("::", 1) - if "://" not in parts[1]: - url_str = parts[0].rstrip("/") - lower = url_str.lower() - for ext in (".h5/", ".hdf5/"): - idx = lower.find(ext) - if idx != -1: - url_str = url_str[: idx + len(ext) - 1] - break - if url_str.endswith(".json"): - try: - import fsspec - with fsspec.open(url_str, "r", **(storage_options or {})) as f: - refs = json.load(f) - except Exception: - with open(url_str) as f: - refs = json.load(f) - ref_dict = refs.get("refs", refs) - elif not urlsplit(url_str).scheme or os.path.isabs(url_str): - _check_h5py_dependencies() - import h5py - - datasets = [] - with h5py.File(url_str, "r") as file: - file.visititems( - lambda name, obj: datasets.append(name) if isinstance(obj, h5py.Dataset) else None - ) - return sorted(datasets) - else: - check_hdf5_dependencies() +def _unshuffle(data, itemsize): + if itemsize <= 1: + return data + if len(data) % itemsize: + raise ValueError("Invalid HDF5 shuffle buffer length") + return np.frombuffer(data, dtype=np.uint8).reshape(itemsize, -1).T.copy().tobytes() + - refs = scan_hdf5_refs(url_str, storage_options) - ref_dict = refs.get("refs", refs) +def _decode_blosc2(data): + try: + return blosc2.decompress(data) + except Exception: + values = blosc2.from_cframe(data)[:] + return values if isinstance(values, bytes) else np.ascontiguousarray(values).tobytes() + + +def _decompress_deflate(data, size): + """Decode one deflate chunk with a hard bound on the output size.""" + decompressor = zlib.decompressobj() + values = decompressor.decompress(data, size + 1) + if len(values) > size or decompressor.unconsumed_tail or decompressor.unused_data: + raise ValueError("Invalid HDF5 deflate chunk") + values += decompressor.flush() + if not decompressor.eof or len(values) != size: + raise ValueError("Invalid HDF5 deflate chunk") + return values + + +def _values_to_chunk(values, chunks, blocks, dtype, cparams): + values = np.asarray(values, dtype=dtype) + buffer = np.zeros(chunks, dtype=dtype) + if values.shape: + values = np.ascontiguousarray(values) + buffer[tuple(slice(0, size) for size in values.shape)] = values else: - raise TypeError("url must be a URL string, path-like object, or reference dict") + buffer[()] = values + converted = blosc2.asarray(buffer, chunks=chunks, blocks=blocks, cparams=cparams) + return converted.schunk.get_chunk(0) - datasets = [] - for k in ref_dict: - if k.endswith("/.zarray"): - datasets.append(k[: -len("/.zarray")]) - elif k == ".zarray": - datasets.append("/") - return sorted(datasets) + +def _close_hdf5_file(h5file, raw): + """Close h5py before the file object it calls into.""" + with contextlib.suppress(Exception): + h5file.close() + with contextlib.suppress(Exception): + raw.close() class HDF5NDSource(ProxyNDSource): - """Read an immutable HDF5 dataset as Blosc2-compressed logical chunks. - - Local files use h5py directly. Remote files and explicit reference maps use - kerchunk and Zarr. Local file handles are released when the source is - garbage-collected. - - Replacing data beneath the same store identity violates this adapter's - contract and may leave previously converted chunks stale. Peak working - memory includes concurrently decoded Zarr chunks and their Blosc2 - conversion buffers; ``max_cache_bytes`` only limits retained compressed - chunks. - - Parameters - ---------- - urlpath : str or path-like - URL or file path to the HDF5 file. - dataset : str - Path to the dataset within the HDF5 file (e.g. ``"d0/d1/a2"``). - refs : dict, str, or path-like, optional - Pre-computed kerchunk reference dictionary or path to a JSON reference - file. If omitted, local files use h5py directly; remote files are scanned - using kerchunk. - storage_options : dict, optional - Parameters passed to fsspec or kerchunk when accessing remote files. - max_concurrency : int, optional - Maximum number of concurrent remote requests. - blocks : tuple, optional - Blosc2 block shape for chunk caching. - cparams : dict or CParams, optional - Blosc2 compression parameters for chunk conversion. - _traffic : Traffic, optional - Traffic monitor instance. - """ + """Read one immutable HDF5 dataset as Blosc2-compressed logical chunks.""" serves_blocks = False + # The logical Blosc2 cache chunks are unchanged from the former reader, so + # compatible warm carrier chunks remain reusable after rebuilding the index. encoding_version = 1 def __init__( @@ -384,112 +500,146 @@ def __init__( urlpath, dataset: str, *, - refs: dict | str | os.PathLike | None = None, - storage_options: dict | None = None, - max_concurrency: int = REMOTE_MAX_CONCURRENCY, + hdf5_index=None, + storage_options=None, + max_concurrency=REMOTE_MAX_CONCURRENCY, blocks=None, cparams=None, _traffic: Traffic | None = None, _filesystem=None, ): - if isinstance(urlpath, os.PathLike): - urlpath = os.fspath(urlpath) - urlpath = blosc2.core.normalize_urlpath(urlpath) - if isinstance(urlpath, str): - if "::" in urlpath: - parts = urlpath.split("::", 1) - if "://" not in parts[1]: - if dataset is not None and dataset != parts[1].strip("/"): - raise ValueError("Cannot specify dataset in both URL path and dataset parameter") - urlpath = parts[0].rstrip("/") - dataset = parts[1].strip("/") - lower = urlpath.lower() - for ext in (".h5/", ".hdf5/"): - idx = lower.find(ext) - if idx != -1: - base_len = idx + len(ext) - 1 - sub = urlpath[base_len + 1 :].strip("/") - if dataset is not None and dataset != sub: - raise ValueError("Cannot specify dataset in both URL path and dataset parameter") - dataset = sub - urlpath = urlpath[:base_len] - break - - if dataset is None: - raise ValueError("HDF5 sources require a dataset path (e.g., dataset='d0/d1/a2')") - if not isinstance(dataset, str): - raise TypeError("dataset must be a string") - - self.urlpath = urlpath if isinstance(urlpath, str) else str(urlpath) - self.dataset = dataset.strip("/") - self.max_concurrency = max_concurrency - - remote = isinstance(self.urlpath, str) and bool(urlsplit(self.urlpath).scheme) + urlpath, dataset = self._parse_url(urlpath, dataset) + self.urlpath, self.dataset, self.max_concurrency = urlpath, dataset, max_concurrency + self._storage_options, self._external_filesystem = storage_options, _filesystem + self._fallback_lock = threading.RLock() + self._fallback_h5 = self._fallback_file = None + self._lifecycle = threading.Condition() + self._active_reads = 0 + self._closed = False + remote = bool(urlsplit(urlpath).scheme) self.traffic = _traffic if _traffic is not None else Traffic() if remote else None - self._local = ( - (not remote or os.path.isabs(self.urlpath)) - and "::" not in self.urlpath - and refs is None - and _filesystem is None - ) - self._refs = None - with contextlib.ExitStack() as stack: + self._local = (not remote or os.path.isabs(urlpath)) and _filesystem is None + self._hdf5_index = None + try: if self._local: - file = self._open_local_array(stack) + self._open_local() + self._metadata = None + if hdf5_index is not None: + # Local reads stay on h5py; still validate and retain the + # explicit index without needing a remote-only dependency. + self._hdf5_index = self._load_or_scan_index(hdf5_index) + shape, physical_chunks, dtype = self.array.shape, self.array.chunks, self.array.dtype else: check_hdf5_dependencies() - check_zarr_fsspec_dependencies() - self._refs = self._load_or_scan_refs(refs, storage_options) + self._filesystem, self._path = _filesystem_and_path(urlpath, storage_options, _filesystem) + if self._external_filesystem is None: + # An abandoned source must not leave its HTTP/S3 session alive. + self._filesystem_finalizer = weakref.finalize( + self, _close_owned_filesystem, self._filesystem + ) + self._hdf5_index = self._load_or_scan_index(hdf5_index) self._validate_dataset_presence(dataset) - self.array = self._open_array(storage_options, _filesystem) - - self._init_geometry(blocks, cparams) - if self._local: - self._file_finalizer = weakref.finalize(self, file.close) - stack.pop_all() - identity = { - "encoding_version": self.encoding_version, - "urlpath": self.urlpath, - "dataset": self.dataset, - "shape": self._shape, - "chunks": self._chunks, - "blocks": self._blocks, - "dtype": self._dtype.str, - } - self.stamp = hashlib.sha256( - json.dumps(identity, sort_keys=True, separators=(",", ":")).encode() - ).hexdigest() + self._metadata = self._hdf5_index["datasets"][self.dataset] + shape, physical_chunks = tuple(self._metadata["shape"]), self._metadata["chunks"] + dtype = dtype_from_value(self._metadata["dtype"]) + self._chunk_records = {tuple(item["offset"]): item for item in self._metadata["allocated"]} + self._init_geometry(shape, physical_chunks, dtype, blocks, cparams) + identity = { + "encoding_version": self.encoding_version, + "urlpath": self.urlpath, + "dataset": self.dataset, + "shape": self._shape, + "chunks": self._chunks, + "blocks": self._blocks, + "dtype": self._dtype.descr if self._dtype.fields else self._dtype.str, + } + self.stamp = hashlib.sha256( + json.dumps(identity, sort_keys=True, separators=(",", ":")).encode() + ).hexdigest() + except BaseException: + # A failed initialization must not leave the h5py file open (and + # locked on Windows) or an owned fsspec session alive. + self.close() + raise + + @staticmethod + def _parse_url(urlpath, dataset): + urlpath = blosc2.core.normalize_urlpath(os.fspath(urlpath)) + if "::" in urlpath and "://" not in urlpath.split("::", 1)[1]: + base, embedded = urlpath.split("::", 1) + embedded = embedded.strip("/") + if dataset is not None and dataset != embedded: + raise ValueError("Cannot specify dataset in both URL path and dataset parameter") + urlpath, dataset = base.rstrip("/"), embedded + base, embedded = blosc2.core.split_h5_url(urlpath) + if embedded is not None: + if dataset is not None and dataset != embedded: + raise ValueError("Cannot specify dataset in both URL path and dataset parameter") + urlpath, dataset = base, embedded + if dataset is None: + raise ValueError("HDF5 sources require a dataset path (e.g., dataset='d0/d1/a2')") + if not isinstance(dataset, str): + raise TypeError("dataset must be a string") + return urlpath, dataset.strip("/") - def _open_local_array(self, stack): + def _open_local(self): _check_h5py_dependencies() import h5py - file = stack.enter_context(h5py.File(self.urlpath, "r")) - if (self.dataset or "/") not in file: - raise ValueError(f"dataset {self.dataset!r} not found in {self.urlpath!r}") - self.array = file[self.dataset or "/"] - if not isinstance(self.array, h5py.Dataset): - raise ValueError(f"{self.dataset!r} is an HDF5 group; pass the path of a dataset") - if self.array.shape is None: - raise TypeError("HDF5 null datasets are not supported") - return file - - def _init_geometry(self, blocks, cparams): - self._shape = tuple(int(value) for value in self.array.shape) - self._chunks = self.array.chunks + file = h5py.File(self.urlpath, "r") try: - self._dtype = np.dtype(self.array.dtype) - except TypeError as exc: - raise TypeError(f"HDF5NDSource only supports fixed-size dtypes, got {self.array.dtype}") from exc + if (self.dataset or "/") not in file: + raise ValueError(f"dataset {self.dataset!r} not found in {self.urlpath!r}") + self.array = file[self.dataset or "/"] + if not isinstance(self.array, h5py.Dataset): + raise ValueError(f"{self.dataset!r} is an HDF5 group; pass the path of a dataset") + if self.array.shape is None: + raise TypeError("HDF5 null datasets are not supported") + except Exception: + file.close() + raise + self._file_finalizer = weakref.finalize(self, file.close) + + def _load_or_scan_index(self, hdf5_index): + if hdf5_index is None: + return scan_hdf5_index( + self.urlpath, self._storage_options, traffic=self.traffic, _filesystem=self._filesystem + ) + if isinstance(hdf5_index, (str, os.PathLike)): + hdf5_index_str = os.fspath(hdf5_index) + if urlsplit(hdf5_index_str).scheme: + import fsspec + + with fsspec.open(hdf5_index_str, "r", **(self._storage_options or {})) as file: + hdf5_index = json.load(file) + else: + with open(hdf5_index_str) as file: + hdf5_index = json.load(file) + if not isinstance(hdf5_index, dict): + raise TypeError("hdf5_index must be a dict, string, or path-like object") + return validate_hdf5_index(hdf5_index, self.urlpath) + + def _validate_dataset_presence(self, raw_dataset): + if self.dataset in self._hdf5_index["groups"]: + raise ValueError( + f"{raw_dataset!r} is an HDF5 group; pass the path of a dataset. Available datasets: {available_datasets(self._hdf5_index)}" + ) + if self.dataset not in self._hdf5_index["datasets"]: + raise ValueError( + f"dataset {raw_dataset!r} not found in {self.urlpath!r}. Available datasets: {available_datasets(self._hdf5_index)}" + ) + + def _init_geometry(self, shape, physical_chunks, dtype, blocks, cparams): + self._shape, self._dtype, self._chunks = ( + tuple(int(v) for v in shape), + np.dtype(dtype), + physical_chunks, + ) if self._chunks is None: - # Contiguous datasets need bounded logical chunks, including for empty axes. self._chunks, _ = blosc2.compute_chunks_blocks( - tuple(max(1, size) for size in self._shape), - blocks=blocks, - dtype=self._dtype, - cparams=cparams, + tuple(max(1, v) for v in self._shape), blocks=blocks, dtype=self._dtype, cparams=cparams ) - self._chunks = tuple(int(value) for value in self._chunks) + self._chunks = tuple(int(v) for v in self._chunks) self._validate_metadata() _, computed_blocks = blosc2.compute_chunks_blocks( self._shape, chunks=self._chunks, blocks=blocks, dtype=self._dtype, cparams=cparams @@ -503,126 +653,151 @@ def _init_geometry(self, blocks, cparams): else cparams ) - def _load_or_scan_refs(self, refs, storage_options) -> dict: - if refs is not None: - if isinstance(refs, (str, os.PathLike)): - refs_str = os.fspath(refs) - if isinstance(refs_str, str) and bool(urlsplit(refs_str).scheme): - import fsspec - - with fsspec.open(refs_str, "r", **(storage_options or {})) as f: - return json.load(f) - with open(refs_str) as f: - return json.load(f) - if isinstance(refs, dict): - return refs - raise TypeError("refs must be a dict, string, or path-like object") - - return scan_hdf5_refs(self.urlpath, storage_options) - - def _validate_dataset_presence(self, raw_dataset: str) -> None: - ref_dict = self._refs.get("refs", self._refs) - clean = self.dataset - - is_group = f"{clean}/.zgroup" in ref_dict or (clean == "" and ".zgroup" in ref_dict) - if is_group: - available = available_datasets(self._refs) - raise ValueError( - f"{raw_dataset!r} is an HDF5 group; pass the path of a dataset. " - f"Available datasets: {available}" - ) - - is_array = f"{clean}/.zarray" in ref_dict or (clean == "" and ".zarray" in ref_dict) - if not is_array: - available = available_datasets(self._refs) - raise ValueError( - f"dataset {raw_dataset!r} not found in {self.urlpath!r}. Available datasets: {available}" - ) - - def _open_array(self, storage_options, filesystem=None): - import fsspec - import zarr - - rfs_kwargs = {} - if storage_options: - rfs_kwargs["target_options"] = storage_options - rfs_kwargs["remote_options"] = storage_options - if filesystem is not None: - rfs_kwargs.update(fs=filesystem, skip_instance_cache=True) - fs = fsspec.filesystem("reference", fo=self._refs, **rfs_kwargs) - mapper = fs.get_mapper(self.dataset) - try: - if filesystem is None: - open_store = zarr.storage.FsspecStore.from_mapper(mapper, read_only=True) - else: - from blosc2.zarr_source import owned_fsspec_store - - open_store = owned_fsspec_store(zarr, fs, mapper.root) - except ValueError: - from fsspec.implementations.asyn_wrapper import AsyncFileSystemWrapper - - wrapped_fs = AsyncFileSystemWrapper(fs, asynchronous=True) - open_store = zarr.storage.FsspecStore(wrapped_fs, path=mapper.root, read_only=True) - if self.traffic is not None: - open_store = counting_store(zarr, open_store, self.traffic) - with ZARR_SYNC_LOCK: - return zarr.open_array(store=open_store, mode="r") - - def _validate_metadata(self) -> None: + def _validate_metadata(self): if len(self._shape) > blosc2.MAX_DIM: raise ValueError(f"HDF5 arrays may have at most {blosc2.MAX_DIM} dimensions") if len(self._chunks) != len(self._shape) or any(size <= 0 for size in self._chunks): raise ValueError("HDF5 chunk extents must be positive and match the array dimensions") if self._dtype.hasobject or self._dtype.itemsize == 0: raise TypeError(f"HDF5NDSource only supports fixed-size dtypes, got {self._dtype}") - chunk_nbytes = math.prod(self._chunks) * self._dtype.itemsize - if chunk_nbytes > blosc2.MAX_BUFFERSIZE: - raise ValueError( - f"HDF5 chunks must be at most {blosc2.MAX_BUFFERSIZE} bytes, got {chunk_nbytes}" - ) + size = math.prod(self._chunks) * self._dtype.itemsize + if size > blosc2.MAX_BUFFERSIZE: + raise ValueError(f"HDF5 chunks must be at most {blosc2.MAX_BUFFERSIZE} bytes, got {size}") - @property - def shape(self) -> tuple: - return self._shape - - @property - def chunks(self) -> tuple: - return self._chunks - - @property - def blocks(self) -> tuple: - return self._blocks - - @property - def dtype(self) -> np.dtype: - return self._dtype - - @property - def cparams(self): - return self._cparams + def _open_fallback(self): + if self._fallback_h5 is not None: + return self._fallback_h5[self.dataset or "/"] + import h5py - @property - def attrs(self) -> dict: - """The user attributes of the remote dataset.""" - return self.vlmeta + if self._closed: + raise RuntimeError("HDF5 source is closed") + raw = self._filesystem.open(self._path, "rb", block_size=1, cache_type="none") + fileobj = _CountingFile(raw, self.traffic) if self.traffic is not None else raw + try: + h5file = h5py.File(fileobj, "r") + except Exception: + raw.close() + raise + self._fallback_file, self._fallback_h5 = raw, h5file + self._fallback_finalizer = weakref.finalize(self, _close_hdf5_file, h5file, raw) + return h5file[self.dataset or "/"] + + def _direct_values(self, offsets, selection): + valid_shape = tuple(item.stop - item.start for item in selection) + # Count in-flight reads so close() can wait for them without serializing + # independent direct fetches against each other. The closed check comes + # first so sparse fill chunks obey the same contract as allocated ones. + with self._lifecycle: + if self._closed: + raise RuntimeError("HDF5 source is closed") + record = self._chunk_records.get(offsets) + if record is None: + return np.full(valid_shape, _from_json_value(self._metadata["fill_value"]), dtype=self.dtype) + filesystem = self._filesystem + self._active_reads += 1 + try: + data = filesystem.cat_file( + self._path, start=record["byte_offset"], end=record["byte_offset"] + record["size"] + ) + finally: + with self._lifecycle: + self._active_reads -= 1 + self._lifecycle.notify_all() + if self.traffic is not None: + self.traffic.charge(len(data)) + if len(data) != record["size"]: + raise OSError(f"Short HDF5 chunk read for {self.dataset!r} at {offsets}") + expected = math.prod(self.chunks) * self.dtype.itemsize + try: + for position in range(len(self._metadata["filters"]) - 1, -1, -1): + if record["filter_mask"] & (1 << position): + continue + info = self._metadata["filters"][position] + if info["id"] == 1: + data = _decompress_deflate(data, expected) + elif info["id"] == 2: + data = _unshuffle(data, info["values"][0] if info["values"] else self.dtype.itemsize) + elif info["id"] == 32026: + data = _decode_blosc2(data) + else: + raise ValueError(f"Unsupported direct HDF5 filter {info['id']}") + except Exception as exc: + raise OSError(f"Cannot decode HDF5 chunk {self.dataset!r} at {offsets}") from exc + if len(data) != expected: + raise ValueError( + f"Decoded HDF5 chunk {self.dataset!r} at {offsets} has {len(data)} bytes, expected {expected}" + ) + values = np.frombuffer(data, dtype=self.dtype).reshape(self.chunks) + return values[tuple(slice(0, size) for size in valid_shape)] + + def close(self): + with self._fallback_lock: + fallback_finalizer = getattr(self, "_fallback_finalizer", None) + if fallback_finalizer is not None: + fallback_finalizer() + self._fallback_h5 = self._fallback_file = None + fs_finalizer = getattr(self, "_filesystem_finalizer", None) + if fs_finalizer is not None: + fs_finalizer.detach() + self._filesystem_finalizer = None + # Reject new reads, then wait for local and direct reads to finish + # before closing the file or filesystem they are using. + with self._lifecycle: + self._closed = True + filesystem = getattr(self, "_filesystem", None) + self._filesystem = None + while self._active_reads: + self._lifecycle.wait() + finalizer = getattr(self, "_file_finalizer", None) + if finalizer is not None: + finalizer() + if filesystem is not None and self._external_filesystem is None: + # fsspec's HTTP and S3 clients keep an async session alive after + # the file objects it feeds are gone. + _close_owned_filesystem(filesystem) + + shape = property(lambda self: self._shape) + chunks = property(lambda self: self._chunks) + blocks = property(lambda self: self._blocks) + dtype = property(lambda self: self._dtype) + cparams = property(lambda self: self._cparams) + attrs = property(lambda self: self.vlmeta) @property - def vlmeta(self) -> dict: + def vlmeta(self): try: if self._local: return dict(self.array.attrs) - # Kerchunk adds dimension metadata to the translated Zarr attributes. - return {key: value for key, value in self.array.attrs.items() if key != "_ARRAY_DIMENSIONS"} + return {key: _from_json_value(value) for key, value in self._metadata["attrs"].items()} except Exception: return {} def get_chunk(self, nchunk: int) -> bytes: - return zarr_chunk_to_blosc2( - self.array, - nchunk, - self.shape, - self.chunks, - self.blocks, - self.dtype, - self.cparams, - ) + offsets, selection = _selection(nchunk, self.shape, self.chunks) + if self._local: + with self._lifecycle: + if self._closed: + raise RuntimeError("HDF5 source is closed") + self._active_reads += 1 + try: + values = self.array[selection] + finally: + with self._lifecycle: + self._active_reads -= 1 + self._lifecycle.notify_all() + elif self._metadata["direct"]: + values = self._direct_values(offsets, selection) + else: + with self._fallback_lock: + try: + values = self._open_fallback()[selection] + except OSError as exc: + if "filter" not in str(exc).lower(): + # Transport, permission or corruption errors keep their cause. + raise + filters = [item["id"] for item in self._metadata["filters"]] + raise OSError( + f"Cannot decode HDF5 dataset {self.dataset!r} with filters {filters}; " + "install hdf5plugin if the file uses an optional HDF5 filter" + ) from exc + return _values_to_chunk(values, self.chunks, self.blocks, self.dtype, self.cparams) diff --git a/src/blosc2/msgpack_utils.py b/src/blosc2/msgpack_utils.py index 7db0fa331..8802d864b 100644 --- a/src/blosc2/msgpack_utils.py +++ b/src/blosc2/msgpack_utils.py @@ -28,6 +28,9 @@ _BLOSC2_STRUCTURED_VERSION = 1 _BLOSC2_COMPLEX_EXT_CODE = 44 _BLOSC2_SET_EXT_CODE = 45 +# Reserve code 46 for NumPy arrays: the payload is a msgpack mapping holding either +# the dtype, shape and raw bytes, or the elements of an object-dtype array. +_BLOSC2_NDARRAY_EXT_CODE = 46 def _encode_structured_reference(obj): @@ -60,6 +63,34 @@ def _decode_structured_reference(data): raise ValueError(f"Unsupported structured Blosc2 msgpack payload kind: {kind!r}") +def _encode_ndarray(value): + from blosc2.hdf5_source import dtype_value + + if value.dtype.hasobject: + payload = {"values": value.ravel().tolist(), "shape": list(value.shape)} + else: + payload = { + "dtype": dtype_value(value.dtype), + "shape": list(value.shape), + "data": value.tobytes(), + } + return ExtType(_BLOSC2_NDARRAY_EXT_CODE, msgpack_packb(payload)) + + +def _decode_ndarray(data): + from blosc2.hdf5_source import dtype_from_value + + payload = msgpack_unpackb(data) + shape = payload["shape"] + if "values" in payload: + result = np.empty(shape, dtype=object) + flat = result.reshape(-1) + for index, item in enumerate(payload["values"]): + flat[index] = item + return result + return np.frombuffer(payload["data"], dtype=dtype_from_value(payload["dtype"])).reshape(shape) + + def _encode_msgpack_ext(obj): import blosc2 @@ -76,6 +107,8 @@ def _encode_msgpack_ext(obj): structured = _encode_structured_reference(obj) if structured is not None: return structured + if isinstance(obj, np.ndarray): + return _encode_ndarray(obj) if isinstance(obj, (complex, np.complexfloating)): return ExtType(_BLOSC2_COMPLEX_EXT_CODE, struct.pack(">dd", float(obj.real), float(obj.imag))) if isinstance(obj, (set, frozenset)): @@ -109,6 +142,8 @@ def _decode_msgpack_ext(code, data): if code == _BLOSC2_COMPLEX_EXT_CODE: real, imag = struct.unpack(">dd", data) return complex(real, imag) + if code == _BLOSC2_NDARRAY_EXT_CODE: + return _decode_ndarray(data) if code == _BLOSC2_SET_EXT_CODE: return set(msgpack_unpackb(data)) return ExtType(code, data) diff --git a/src/blosc2/remote_array.py b/src/blosc2/remote_array.py index 53b3e0a9e..72fffedc2 100644 --- a/src/blosc2/remote_array.py +++ b/src/blosc2/remote_array.py @@ -126,42 +126,44 @@ def _validate_assume_immutable(value, name="assume_immutable"): return value -def _hdf5_refs_from_carrier(carrier): - """Decode the kerchunk reference snapshot a carrier stores, if present.""" +def _hdf5_index_from_carrier(carrier): + """Decode the native HDF5 index a carrier stores, if present.""" if carrier is None: return None - raw_refs = getattr(carrier, "schunk", carrier).vlmeta.get("hdf5-refs") - if raw_refs is None: + raw_hdf5_index = getattr(carrier, "schunk", carrier).vlmeta.get("hdf5-index") + if raw_hdf5_index is None: return None try: import ujson as json_mod except ImportError: import json as json_mod - return json_mod.loads(blosc2.decompress(raw_refs).decode("utf-8")) + return json_mod.loads(blosc2.decompress(raw_hdf5_index).decode("utf-8")) -def _store_hdf5_refs(carrier, refs): - """Keep the kerchunk snapshot on the carrier so a reopen needs no rescan.""" - if refs is None: +def _store_hdf5_index(carrier, hdf5_index): + """Keep the native HDF5 index on the carrier so a reopen needs no rescan.""" + if hdf5_index is None: return try: import ujson as json_mod except ImportError: import json as json_mod - carrier.schunk.vlmeta["hdf5-refs"] = blosc2.compress(json_mod.dumps(refs).encode("utf-8"), typesize=1) + carrier.schunk.vlmeta["hdf5-index"] = blosc2.compress( + json_mod.dumps(hdf5_index).encode("utf-8"), typesize=1 + ) -def _publish_hdf5_refs(path, carrier, *, scanned): +def _publish_hdf5_index(path, carrier, *, scanned): if path is None or carrier is None: return - raw_refs = getattr(carrier, "schunk", carrier).vlmeta.get("hdf5-refs") - if raw_refs is None: - return # Local h5py readers have no reference snapshot to share. + raw_hdf5_index = getattr(carrier, "schunk", carrier).vlmeta.get("hdf5-index") + if raw_hdf5_index is None: + return # Local h5py readers have no native index to share. if scanned or not path.exists(): from blosc2.remote_store_cache import atomic_write # Also seed the shared snapshot when reopening an older leaf cache. - atomic_write(path, raw_refs) + atomic_write(path, raw_hdf5_index) def _b2z_seed_from_carrier(carrier): @@ -320,7 +322,7 @@ def _open_url_source( source_format=None, assume_immutable=True, dataset=None, - refs=None, + hdf5_index=None, seed=None, blocks=None, cparams=None, @@ -351,7 +353,7 @@ def _open_url_source( src = blosc2.HDF5NDSource( urlpath, dataset, - refs=refs, + hdf5_index=hdf5_index, _traffic=traffic, blocks=blocks, cparams=cparams, @@ -443,7 +445,18 @@ def _parse_source_from_payload(source): return source_kind, urlpath -def _resolve_init_dataset_and_url(urlpath, dataset, source_format): +def _source_format_with_index(urlpath, source_format, hdf5_index): + """Fold an explicit HDF5 index into the requested source format.""" + if hdf5_index is None: + return source_format + if isinstance(urlpath, (blosc2.URLPath, blosc2.C2Array)): + raise ValueError("hdf5_index is not supported for Caterva2 inputs") + if source_format not in {None, "hdf5"}: + raise ValueError("hdf5_index is only supported for HDF5 sources") + return "hdf5" + + +def _resolve_init_dataset_and_url(urlpath, dataset, source_format, hdf5_index=None): if isinstance(urlpath, (blosc2.URLPath, blosc2.C2Array)) and source_format is not None: raise ValueError("source_format is not supported for Caterva2 inputs") urlpath, parsed_dataset, detected_format = parse_container_url(urlpath, dataset) @@ -451,6 +464,7 @@ def _resolve_init_dataset_and_url(urlpath, dataset, source_format): dataset = parsed_dataset if source_format is None: source_format = detected_format + source_format = _source_format_with_index(urlpath, source_format, hdf5_index) resolved_format = _normalize_source_format(urlpath, source_format) if dataset is not None and resolved_format not in {"hdf5", "zarr", "b2z"}: raise ValueError("dataset is only supported for HDF5 and Zarr sources or B2Z archives") @@ -517,7 +531,7 @@ class RemoteArray(blosc2.Operand): cache_dir: str or path-like, optional Directory in which a source-derived persistent cache filename is made. Only valid with ``DISK``. - HDF5 leaves also share a container reference snapshot here, avoiding + HDF5 leaves also share a native container index here, avoiding repeated discovery scans when opening sibling datasets. max_cache_bytes: int or None, optional Post-operation compressed-payload bound. It defaults to 256 MiB for @@ -534,6 +548,10 @@ class RemoteArray(blosc2.Operand): dataset: str, optional Array path within an HDF5, Zarr, or B2Z container. B2Z supports external NDArray leaves in immutable archives, e.g. ``dataset="d0/a3"``. + hdf5_index: dict, str, or path-like, optional + Pre-computed native HDF5 index for the dataset, or the path to a JSON + encoding of one. It must match the source URL. Legacy HDF5 reference + maps are rejected; omit it to rescan and build a native index. assume_immutable: bool, optional Skip remote identity checks before reads. Defaults to ``True``. Set to ``False`` when the object at the URL may be replaced. @@ -552,7 +570,7 @@ def __init__( source_format: str | None = None, assume_immutable: bool = True, dataset: str | None = None, - refs=None, + hdf5_index=None, _carrier=None, _runtime_cache_path=None, _source_descriptor=None, @@ -573,19 +591,19 @@ def __init__( self._cache_limit = normalize_cache_limit(cache_policy, max_cache_bytes) self._max_concurrency = _validate_max_concurrency(max_concurrency) urlpath, self._dataset, self._source_format = _resolve_init_dataset_and_url( - urlpath, dataset, source_format + urlpath, dataset, source_format, hdf5_index ) self._authorized_source = _source_descriptor is not None - shared_refs_path = None + shared_index_path = None if ( not self._authorized_source and self._source_format == "hdf5" and cache_policy is blosc2.CachePolicy.DISK and cache_dir is not None - and refs is None + and hdf5_index is None ): - shared_refs_path = Path( - fsspec_cache_path(urlpath, cache_dir, ".hdf5-refs.b2", storage_options=storage_options) + shared_index_path = Path( + fsspec_cache_path(urlpath, cache_dir, ".hdf5-index.b2", storage_options=storage_options) ) if self._authorized_source: self.src, self._source = _validate_authorized_source( @@ -596,10 +614,10 @@ def __init__( _zarr_metadata_from_carrier if self._source_format == "zarr" else _b2z_seed_from_carrier ) seed = read_seed(_carrier) if self._source_format in {"b2z", "zarr"} else None - if refs is None and _carrier is not None: - refs = _hdf5_refs_from_carrier(_carrier) + if hdf5_index is None and _carrier is not None: + hdf5_index = _hdf5_index_from_carrier(_carrier) elif ( - refs is None + hdf5_index is None and _carrier is None and cache_policy is blosc2.CachePolicy.DISK and (cache_dir is not None or cache_path is not None) @@ -612,13 +630,16 @@ def __init__( with contextlib.suppress(Exception): cached = blosc2.blosc2_ext.open(path, "r", 0, dparams=blosc2.DParams(nthreads=1)) if self._source_format == "hdf5": - refs = _hdf5_refs_from_carrier(cached) + hdf5_index = _hdf5_index_from_carrier(cached) else: seed = read_seed(cached) - if refs is None and shared_refs_path is not None: + if hdf5_index is None and shared_index_path is not None: # Disposable metadata: an absent or damaged snapshot needs a fresh scan. + from blosc2.hdf5_source import validate_hdf5_index + with contextlib.suppress(OSError, ValueError, RuntimeError): - refs = json.loads(blosc2.decompress(shared_refs_path.read_bytes())) + candidate = json.loads(blosc2.decompress(shared_index_path.read_bytes())) + hdf5_index = validate_hdf5_index(candidate, urlpath) self.src, self._source = self._open_source( urlpath, self._max_concurrency, @@ -627,7 +648,7 @@ def __init__( source_format=self._source_format, assume_immutable=assume_immutable, dataset=self._dataset, - refs=refs, + hdf5_index=hdf5_index, seed=seed, blocks=_source_blocks, cparams=_source_cparams, @@ -639,6 +660,7 @@ def __init__( self._expected_cparams = self.src.cparams self._refresh_lock = threading.Lock() self._operation_lock = threading.RLock() + self._closed = False self._proxy = None self._carrier = _carrier self._runtime_cache = _carrier if cache_policy is blosc2.CachePolicy.DISK else None @@ -653,7 +675,7 @@ def __init__( self._initialize_runtime_cache(cache_dir, cache_path, _runtime_cache_path) - _publish_hdf5_refs(shared_refs_path, self._carrier, scanned=refs is None) + _publish_hdf5_index(shared_index_path, self._carrier, scanned=hdf5_index is None) if self._carrier is not None: if self._cached_meta is None: @@ -692,6 +714,8 @@ def _initialize_runtime_cache(self, cache_dir, cache_path, _runtime_cache_path): _store_b2z_seed(self._carrier, self.src._seed) def _check_open(self): + if getattr(self, "_closed", False): + raise RuntimeError("RemoteArray handle is closed") finalizer = getattr(self, "_store_finalizer", None) if finalizer is not None and not finalizer.alive: raise RuntimeError("RemoteArray handle is closed") @@ -699,13 +723,27 @@ def _check_open(self): raise RuntimeError("RemoteArray handle is stale; look it up again after refresh") def close(self): - """Release this store-derived handle; standalone handles retain their existing lifetime.""" + """Release this handle and any HDF5 file resources it owns. + + Closing is idempotent. Standalone HDF5 handles reject further + operations; other standalone handles keep their no-op close behavior. + """ + if getattr(self, "_closed", False): + return owner = getattr(self, "_store_owner", None) if owner is not None: with owner.lock, self._operation_lock: + self._closed = True self._store_finalizer() self._proxy = None self._runtime_cache = None + else: + hdf5_source = getattr(self, "src", None) + if isinstance(hdf5_source, blosc2.HDF5NDSource): + with self._operation_lock: + # Serialize against in-flight reads before closing the source. + self._closed = True + hdf5_source.close() def _runtime_source(self, original): """Keep credentials in live process state, outside the descriptor.""" @@ -754,8 +792,8 @@ def _open_or_create_carrier(self, cache_dir, cache_path): "open legacy Proxy caches directly with blosc2.open(cache_path), " "or choose a new cache_path" ) - if self._source.get("kind") == "hdf5" and "hdf5-refs" not in carrier.schunk.vlmeta: - _store_hdf5_refs(carrier, getattr(self.src, "_refs", None)) + if self._source.get("kind") == "hdf5" and "hdf5-index" not in carrier.schunk.vlmeta: + _store_hdf5_index(carrier, getattr(self.src, "_hdf5_index", None)) if self._source.get("kind") == "zarr" and "zarr-metadata" not in carrier.schunk.vlmeta: carrier.schunk.vlmeta["zarr-metadata"] = self.src._metadata stored = carrier.schunk.vlmeta.get("proxy-stamp") @@ -1037,7 +1075,7 @@ def _open_source( source_format: str | None = None, assume_immutable: bool = True, dataset: str | None = None, - refs=None, + hdf5_index=None, seed=None, blocks=None, cparams=None, @@ -1085,7 +1123,7 @@ def _open_source( source_format=source_format, assume_immutable=assume_immutable, dataset=dataset, - refs=refs, + hdf5_index=hdf5_index, seed=seed, blocks=blocks, cparams=cparams, @@ -1161,7 +1199,7 @@ def _prepare_read(self): source_format=self._source_format, assume_immutable=self._assume_immutable, dataset=self.dataset, - refs=getattr(self.src, "_refs", None), + hdf5_index=getattr(self.src, "_hdf5_index", None), ) if current_stamp is None and not isinstance(fresh, blosc2.C2Array): # No stable validator means cached bytes cannot safely be @@ -1582,13 +1620,13 @@ def _to_b2object_carrier(self, mutable=None, **kwargs): if user_vlmeta: write_b2object_user_vlmeta(array, user_vlmeta) if self._source.get("kind") == "hdf5": - refs = getattr(self.src, "_refs", None) - if refs is not None: - _store_hdf5_refs(array, refs) + hdf5_index = getattr(self.src, "_hdf5_index", None) + if hdf5_index is not None: + _store_hdf5_index(array, hdf5_index) elif self._carrier is not None: carrier_schunk = getattr(self._carrier, "schunk", self._carrier) - if "hdf5-refs" in carrier_schunk.vlmeta: - array.schunk.vlmeta["hdf5-refs"] = carrier_schunk.vlmeta["hdf5-refs"] + if "hdf5-index" in carrier_schunk.vlmeta: + array.schunk.vlmeta["hdf5-index"] = carrier_schunk.vlmeta["hdf5-index"] elif self._source.get("kind") == "zarr": array.schunk.vlmeta["zarr-metadata"] = self.src._metadata elif self._source.get("kind") == "b2z": @@ -1720,7 +1758,7 @@ def _from_payload(cls, payload, carrier): expected = (carrier.shape, carrier.dtype, carrier.chunks, carrier.blocks) kwargs = {} if policy is blosc2.CachePolicy.NONE else {"max_cache_bytes": limit} carrier_arg = carrier if policy is blosc2.CachePolicy.DISK else None - refs = _hdf5_refs_from_carrier(carrier) if source_kind == "hdf5" else None + hdf5_index = _hdf5_index_from_carrier(carrier) if source_kind == "hdf5" else None carrier_mode = getattr(carrier.schunk, "mode", "r") if carrier is not None else "r" is_disk_file = carrier is not None and bool(getattr(carrier.schunk, "urlpath", None)) is_runtime_mutable = mutable and (carrier_mode != "r" if is_disk_file else True) @@ -1729,7 +1767,7 @@ def _from_payload(cls, payload, carrier): cache_policy=policy, source_format=source_kind if source_kind in {"zarr", "hdf5", "b2z"} else None, dataset=source.get("dataset") if source_kind in {"hdf5", "b2z"} else None, - refs=refs, + hdf5_index=hdf5_index, assume_immutable=source["assume_immutable"], _carrier=carrier_arg, _source_blocks=carrier.blocks if source_kind in {"zarr", "hdf5"} else None, diff --git a/src/blosc2/remote_store.py b/src/blosc2/remote_store.py index 61e7472c1..3aca6502f 100644 --- a/src/blosc2/remote_store.py +++ b/src/blosc2/remote_store.py @@ -3,7 +3,6 @@ from __future__ import annotations import contextlib -import json import os import shutil import tempfile @@ -171,33 +170,24 @@ def _restore_manifest(self, manifest): _filesystem=self.filesystem, ) elif self.format == "hdf5": - self.refs = self.metadata - self._validate_refs() + self.hdf5_index = self.metadata + self._validate_hdf5_index() elif self.format == "zarr" and self.zstore is None: self._open_zarr() self._check_node_limit() - def _validate_refs(self): - refs = self.refs.get("refs", self.refs) - if not isinstance(refs, dict) or self.refs.get("templates"): - raise ValueError("Invalid HDF5 manifest references") - for value in refs.values(): - if isinstance(value, list): - if ( - len(value) != 3 - or value[0] != self.urlpath - or any(isinstance(n, bool) or not isinstance(n, int) or n < 0 for n in value[1:]) - ): - raise ValueError("HDF5 manifest contains an unsafe reference") - validate_persistable_url(value[0]) + def _validate_hdf5_index(self): + from blosc2.hdf5_source import validate_hdf5_index + + validate_hdf5_index(self.hdf5_index, self.urlpath) def save_manifest(self): if self.disk is None or self.restoring or not self.is_mutable: return # ponytail: publish whole discovery snapshots; add dirty tracking if large maps make this costly. if self.format == "hdf5": - self._validate_refs() - self.metadata = self.refs + self._validate_hdf5_index() + self.metadata = self.hdf5_index elif self.format == "b2z": self.metadata = self.archive.metadata nodes = { @@ -369,30 +359,30 @@ def _open_b2z(self): self.archive.capture_metadata = False def _open_hdf5(self): - from blosc2.hdf5_source import scan_hdf5_refs + from blosc2.hdf5_source import decode_hdf5_value, scan_hdf5_index unsupported = {} - self.refs = scan_hdf5_refs( + self.hdf5_index = scan_hdf5_index( self.urlpath, self.storage_options, unsupported=unsupported, traffic=self.traffic, _filesystem=self.filesystem, ) - refs = self.refs.get("refs", self.refs) - for key, value in refs.items(): - name = key.rsplit("/", 1)[-1] - path = key.rpartition("/")[0] - if name in {".zgroup", ".zarray"}: - self._add(path, "group" if name == ".zgroup" else "ndarray", json.loads(value)) - elif name == ".zattrs": - self.attrs[path] = {k: v for k, v in json.loads(value).items() if k != "_ARRAY_DIMENSIONS"} + for path, metadata in self.hdf5_index["groups"].items(): + self._add(path, "group") + self.attrs[path] = { + key: decode_hdf5_value(value) for key, value in metadata.get("attrs", {}).items() + } + for path, metadata in self.hdf5_index["datasets"].items(): + self._add(path, "ndarray", metadata) + self.attrs[path] = { + key: decode_hdf5_value(value) for key, value in metadata.get("attrs", {}).items() + } for path, message in unsupported.items(): self._validate(path) self.nodes[path] = ("unsupported", message) - self.notice = ( - "HDF5 view includes objects represented by Kerchunk; external and group links are omitted." - ) + self.notice = "HDF5 view includes indexed groups and datasets; external and soft links are omitted." def _open_zarr(self): import zarr @@ -518,7 +508,7 @@ def open_source(self, path): source = HDF5NDSource( self.urlpath, full, - refs=self.refs, + hdf5_index=self.hdf5_index, storage_options=self.storage_options, _traffic=self.traffic, _filesystem=self.filesystem, @@ -660,13 +650,13 @@ def _close_resources(self): self.zstore.close() for source in self.sources.values(): if isinstance(source, blosc2.HDF5NDSource): - source.array.store.close() + source.close() self.sources.clear() self.caches.clear() self.nodes.clear() self.attrs.clear() self.listed.clear() - self.refs = None + self.hdf5_index = None if self.filesystem is not None and self._external_filesystem is None: # fsspec's HTTP and S3 clients expose their own synchronous close hook. close = getattr(self.filesystem, "close_session", None) @@ -1174,8 +1164,8 @@ def save( ) if self._owner.format == "hdf5": - self._owner._validate_refs() - metadata = self._owner.refs + self._owner._validate_hdf5_index() + metadata = self._owner.hdf5_index elif self._owner.format == "b2z": metadata = self._owner.archive.metadata else: diff --git a/src/blosc2/schunk.py b/src/blosc2/schunk.py index 1355420d3..6ce2e4a8b 100644 --- a/src/blosc2/schunk.py +++ b/src/blosc2/schunk.py @@ -1909,18 +1909,18 @@ def _reconstruct_legacy_proxy(proxy_cache, proxy_src): ) return blosc2.Proxy(src, _cache=proxy_cache, _refresh_source=False) if source_kind == "hdf5": - refs = None - raw_refs = getattr(proxy_cache, "schunk", proxy_cache).vlmeta.get("hdf5-refs") - if raw_refs is not None: + hdf5_index = None + raw_index = getattr(proxy_cache, "schunk", proxy_cache).vlmeta.get("hdf5-index") + if raw_index is not None: try: import ujson as json_mod except ImportError: import json as json_mod - refs = json_mod.loads(blosc2.decompress(raw_refs).decode("utf-8")) + hdf5_index = json_mod.loads(blosc2.decompress(raw_index).decode("utf-8")) src = blosc2.HDF5NDSource( proxy_src["urlpath"], proxy_src["dataset"], - refs=refs, + hdf5_index=hdf5_index, blocks=proxy_cache.blocks, cparams=proxy_cache.cparams, ) @@ -2069,7 +2069,7 @@ def _remote_array_options( source_format=None, assume_immutable=True, dataset=None, - refs=None, + hdf5_index=None, ): """Return the explicit RemoteArray options, or None when remote access was not requested.""" policy_present = "cache_policy" in kwargs @@ -2100,16 +2100,16 @@ def _remote_array_options( options["source_format"] = source_format if dataset is not None: options["dataset"] = dataset - if refs is not None: - options["refs"] = refs + if hdf5_index is not None: + options["hdf5_index"] = hdf5_index return options def _validate_c2_urlpath_options(kwargs: dict): if kwargs.pop("dataset", None) is not None: raise ValueError("dataset is not supported for Caterva2 inputs") - if kwargs.pop("refs", None) is not None: - raise ValueError("refs is not supported for Caterva2 inputs") + if kwargs.pop("hdf5_index", None) is not None: + raise ValueError("hdf5_index is not supported for Caterva2 inputs") source_format = kwargs.pop("source_format", None) if source_format not in {None, "blosc2", "zarr", "hdf5"}: raise ValueError("source_format must be None, 'blosc2', 'zarr', or 'hdf5'") @@ -2184,7 +2184,7 @@ def _validate_non_lazy_fsspec_options(immutable_present, remote_array_options, c raise NotImplementedError("max_concurrency is only supported with lazy=True") -def _open_localized_fsspec(localized, mode, offset, kwargs, source_format, dataset, refs): +def _open_localized_fsspec(localized, mode, offset, kwargs, source_format, dataset, hdf5_index): """Dispatch an already-localized fsspec container with its original selection.""" if source_format == "b2z": # The localized archive is a plain local TreeStore now, and an @@ -2192,8 +2192,8 @@ def _open_localized_fsspec(localized, mode, offset, kwargs, source_format, datas from blosc2.tree_store import TreeStore return _open_treestore_root_object(TreeStore(localized, mode=mode), localized, mode) - if dataset is not None or refs is not None: - raise NotImplementedError("dataset and refs are only supported with lazy=True") + if dataset is not None or hdf5_index is not None: + raise NotImplementedError("dataset and hdf5_index are only supported with lazy=True") return open(localized, mode, offset, **kwargs) @@ -2216,6 +2216,20 @@ def _resolve_lazy(lazy, dataset, source_format, urlpath): return _infer_lazy(False, dataset, source_format, urlpath) if lazy is None else lazy +def _resolve_fsspec_format(urlpath, dataset, source_format, hdf5_index): + """Infer the container format; an explicit HDF5 index forces HDF5.""" + urlpath, parsed_dataset, detected_format = parse_container_url(urlpath, dataset) + if dataset is None: + dataset = parsed_dataset + if source_format is None: + source_format = detected_format + if hdf5_index is not None: + if source_format not in {None, "hdf5"}: + raise ValueError("hdf5_index is only supported for HDF5 sources") + source_format = "hdf5" + return urlpath, dataset, source_format + + def _open_fsspec_url(urlpath: str, mode: str, offset: int, kwargs: dict): """Open a container living behind an fsspec URL. @@ -2232,22 +2246,20 @@ def _open_fsspec_url(urlpath: str, mode: str, offset: int, kwargs: dict): storage_options = kwargs.pop("storage_options", None) source_format = kwargs.pop("source_format", None) dataset = kwargs.pop("dataset", None) - refs = kwargs.pop("refs", None) + hdf5_index = kwargs.pop("hdf5_index", None) max_concurrency = kwargs.pop("max_concurrency", None) immutable_present = "assume_immutable" in kwargs assume_immutable = kwargs.pop("assume_immutable", True) lazy = kwargs.pop("lazy", None) - urlpath, parsed_dataset, detected_format = parse_container_url(urlpath, dataset) - if dataset is None: - dataset = parsed_dataset - if source_format is None: - source_format = detected_format + urlpath, dataset, source_format = _resolve_fsspec_format(urlpath, dataset, source_format, hdf5_index) # Local-file options require a complete local copy; cache placement alone # does not choose between lazy access and eager localization. if lazy is None and (offset != 0 or kwargs.get("mmap_mode") is not None): lazy = False + if lazy is False and hdf5_index is not None: + raise NotImplementedError("hdf5_index is only supported with lazy=True") # Auto-infer lazy=True only when the caller left the choice unspecified. lazy = _resolve_lazy(lazy, dataset, source_format, urlpath) @@ -2262,7 +2274,7 @@ def _open_fsspec_url(urlpath: str, mode: str, offset: int, kwargs: dict): source_format=source_format, assume_immutable=assume_immutable, dataset=dataset, - refs=refs, + hdf5_index=hdf5_index, ) if lazy: if offset != 0: @@ -2276,7 +2288,7 @@ def _open_fsspec_url(urlpath: str, mode: str, offset: int, kwargs: dict): if cache_dir is not None: localized = localize_fsspec_url(urlpath, cache_dir, storage_options=storage_options) - return _open_localized_fsspec(localized, mode, offset, kwargs, source_format, dataset, refs) + return _open_localized_fsspec(localized, mode, offset, kwargs, source_format, dataset, hdf5_index) if source_format == "b2z": raise NotImplementedError( @@ -2300,7 +2312,7 @@ def _open_fsspec_url(urlpath: str, mode: str, offset: int, kwargs: dict): def _is_hdf5_open_request(urlpath: str, kwargs: dict) -> bool: - if kwargs.get("source_format") == "hdf5" or "refs" in kwargs: + if kwargs.get("source_format") == "hdf5" or "hdf5_index" in kwargs: return True if not isinstance(urlpath, str): return False @@ -2311,7 +2323,7 @@ def _is_hdf5_open_request(urlpath: str, kwargs: dict) -> bool: def _is_container_open_request(urlpath: str, kwargs: dict) -> bool: if os.path.isfile(urlpath) and urlpath.endswith((".b2nd", ".b2frame")) and "dataset" not in kwargs: return False - if kwargs.get("source_format") in {"hdf5", "zarr"} or "refs" in kwargs: + if kwargs.get("source_format") in {"hdf5", "zarr"} or "hdf5_index" in kwargs: return True if not isinstance(urlpath, str): return False @@ -2342,11 +2354,11 @@ def _try_open_aliased_store(urlpath: str, mode: str, offset: int, kwargs: dict): return special, special_path -def _normalize_open_target(urlpath, kwargs, dataset, refs): +def _normalize_open_target(urlpath, kwargs, dataset, hdf5_index): if dataset is not None: kwargs["dataset"] = dataset - if refs is not None: - kwargs["refs"] = refs + if hdf5_index is not None: + kwargs["hdf5_index"] = hdf5_index if isinstance(urlpath, pathlib.PurePath): urlpath = str(urlpath) urlpath = normalize_urlpath(urlpath) @@ -2370,7 +2382,7 @@ def open( mode: str = "r", offset: int = 0, dataset: str | None = None, - refs: dict | str | os.PathLike | None = None, + hdf5_index: dict | str | os.PathLike | None = None, **kwargs: dict, ) -> ( blosc2.SChunk @@ -2498,9 +2510,8 @@ def open( Array path within HDF5, Zarr, or B2Z containers (e.g. ``dataset="d0/d1/a2"``). B2Z supports external NDArray leaves in immutable archives. Requires ``lazy=True``. - refs: dict | str | PathLike, optional - Pre-computed kerchunk reference dictionary or path to a JSON reference - file for HDF5 sources. + hdf5_index: dict | str | PathLike, optional + Pre-computed native HDF5 index or path to a JSON index file. source_format: {None, "blosc2", "zarr", "hdf5", "b2z"}, optional Format of a lazy remote source. A ``.zarr`` URL path component selects Zarr automatically; a ``.h5`` or ``.hdf5`` path selects HDF5 automatically; @@ -2521,9 +2532,11 @@ def open( ----- * Returned objects can be used as context managers for API consistency. For objects with an explicit ``close()`` implementation, exiting the - context will close/flush them; for logical handles such as regular - :class:`SChunk`, :class:`NDArray`, :class:`C2Array`, standalone :class:`RemoteArray`, - :class:`Proxy`, and :class:`LazyArray`, exiting the context is currently a + context will close/flush them. Standalone HDF5 :class:`RemoteArray` + handles close their HDF5 source, so the handle rejects further reads. + Other logical handles such as regular :class:`SChunk`, :class:`NDArray`, + :class:`C2Array`, standalone non-HDF5 :class:`RemoteArray`, + :class:`Proxy`, and :class:`LazyArray` currently treat context exit as a no-op. Store-derived :class:`RemoteArray` handles release their shared source ownership when closed; other handles from that store remain usable. @@ -2607,12 +2620,12 @@ def open( if offset != 0 and not is_fsspec_url(urlpath): local_path = normalize_urlpath(os.fspath(urlpath)) if os.path.isfile(local_path): - if dataset is not None or refs is not None: - raise ValueError("dataset and refs cannot be combined with an embedded frame offset") + if dataset is not None or hdf5_index is not None: + raise ValueError("dataset and hdf5_index cannot be combined with an embedded frame offset") _set_default_dparams(kwargs) return process_opened_object(blosc2_ext.open(local_path, mode, offset, **kwargs)) - urlpath = _normalize_open_target(urlpath, kwargs, dataset, refs) + urlpath = _normalize_open_target(urlpath, kwargs, dataset, hdf5_index) if is_fsspec_url(urlpath) or _is_container_open_request(urlpath, kwargs): return _open_fsspec_url(urlpath, mode, offset, kwargs) diff --git a/tests/b2view/test_hierarchy.py b/tests/b2view/test_hierarchy.py index 5d9619118..6cba00331 100644 --- a/tests/b2view/test_hierarchy.py +++ b/tests/b2view/test_hierarchy.py @@ -126,7 +126,8 @@ def test_zarr_hierarchy(version, consolidated): def test_hdf5_hierarchy(tmp_path, monkeypatch): h5py = pytest.importorskip("h5py") - kerchunk = pytest.importorskip("kerchunk.hdf") + import blosc2.hdf5_source as hdf5_source + path = tmp_path / "hierarchy.h5" data = np.arange(600).reshape(30, 20) with h5py.File(path, "w") as root: @@ -139,13 +140,13 @@ def test_hdf5_hierarchy(tmp_path, monkeypatch): url = "memory://v12/hierarchy.h5" fsspec.filesystem("memory").pipe_file(url, path.read_bytes()) calls = [] - translate = kerchunk.SingleHdf5ToZarr.translate + scan = hdf5_source.scan_hdf5_index - def counted(self, *args, **kwargs): + def counted(*args, **kwargs): calls.append(1) - return translate(self, *args, **kwargs) + return scan(*args, **kwargs) - monkeypatch.setattr(kerchunk.SingleHdf5ToZarr, "translate", counted) + monkeypatch.setattr(hdf5_source, "scan_hdf5_index", counted) with StoreBrowser(url) as browser: assert browser.get_info("/").user_attrs == {"title": "root"} assert browser.get_info("/group").user_attrs == {"title": "child"} @@ -313,7 +314,6 @@ async def counted(self, key, *args, **kwargs): def test_hdf5_unsupported_and_links(tmp_path): h5py = pytest.importorskip("h5py") - pytest.importorskip("kerchunk") path = tmp_path / "links.h5" with h5py.File(path, "w") as file: file.create_dataset("a", data=np.arange(10)) @@ -332,7 +332,7 @@ def test_hdf5_unsupported_and_links(tmp_path): assert children["empty"] == "group" assert "external" not in children assert "cycle" not in children - assert "Kerchunk" in browser.get_info("/").metadata["notice"] + assert "indexed groups" in browser.get_info("/").metadata["notice"] np.testing.assert_array_equal(browser.preview("/a", start=0, stop=3)["data"]["value"], np.arange(3)) @@ -426,7 +426,6 @@ def test_group_decoder_fresh_process(tmp_path, format): else: h5py = pytest.importorskip("h5py") plugin = pytest.importorskip("hdf5plugin") - pytest.importorskip("kerchunk") with h5py.File(path, "w") as file: file.create_dataset("a", data=np.arange(100, dtype="i4"), chunks=(20,), **plugin.Blosc()) script = """ @@ -474,7 +473,7 @@ def test_b2z_optional_dependency_isolation(tmp_path): import sys class BlockOptional(importlib.abc.MetaPathFinder): def find_spec(self, fullname, path=None, target=None): - if fullname.split('.')[0] in {'zarr', 'h5py', 'kerchunk'}: + if fullname.split('.')[0] in {'zarr', 'h5py'}: raise ImportError('optional dependency intentionally unavailable') sys.meta_path.insert(0, BlockOptional()) from pathlib import Path diff --git a/tests/test_b2z_source.py b/tests/test_b2z_source.py index 679620dc6..eb9ea4922 100644 --- a/tests/test_b2z_source.py +++ b/tests/test_b2z_source.py @@ -490,7 +490,7 @@ def test_optional_dependencies(): import builtins original = builtins.__import__ def blocked(name, *args, **kwargs): - if name.split('.')[0] in {'zarr', 'kerchunk', 'h5py', 'hdf5plugin'}: + if name.split('.')[0] in {'zarr', 'h5py', 'hdf5plugin'}: raise ImportError('blocked optional dependency') return original(name, *args, **kwargs) builtins.__import__ = blocked diff --git a/tests/test_fsspec.py b/tests/test_fsspec.py index a3c6877a2..92c12ebe1 100644 --- a/tests/test_fsspec.py +++ b/tests/test_fsspec.py @@ -8,6 +8,7 @@ import contextlib import functools +import gc import hashlib import http.server import os @@ -770,8 +771,6 @@ def do_HEAD(self): def test_http_hdf5_scan_and_warm_slice(tmp_path): h5py = pytest.importorskip("h5py") - pytest.importorskip("kerchunk") - pytest.importorskip("zarr") data = np.arange(10_000, dtype="int32") path = tmp_path / "seekable.h5" with h5py.File(path, "w") as file: @@ -786,10 +785,108 @@ def test_http_hdf5_scan_and_warm_slice(tmp_path): assert len(requests) == count +def test_http_hdf5_source_close_closes_session(tmp_path): + h5py = pytest.importorskip("h5py") + data = np.arange(10_000, dtype="int32") + path = tmp_path / "standalone.h5" + with h5py.File(path, "w") as file: + file.create_dataset("data", data=data, chunks=(1000,)) + file.create_dataset("sparse", shape=(100,), chunks=(10,), dtype="int32") + with _ranged_server(tmp_path) as (urlbase, _requests): + source = blosc2.HDF5NDSource( + f"{urlbase}/{path.name}", + "data", + storage_options={"skip_instance_cache": False}, + ) + filesystem = source._filesystem + cached, _ = fsspec.core.url_to_fs(f"{urlbase}/{path.name}") + assert filesystem is not cached # Sources own a private filesystem. + session = filesystem._session + assert not session.closed + source.close() + assert session.closed + with pytest.raises(RuntimeError, match="closed"): + source.get_chunk(0) + # Sparse fill chunks obey the same closed contract as allocated chunks. + sparse = blosc2.HDF5NDSource(f"{urlbase}/{path.name}", "sparse") + assert sparse._metadata["direct"] is True + sparse.close() + with pytest.raises(RuntimeError, match="closed"): + sparse.get_chunk(0) + + +def test_http_hdf5_scan_closes_owned_session(tmp_path, monkeypatch): + h5py = pytest.importorskip("h5py") + import blosc2.hdf5_source as hdf5_source + + data = np.arange(10_000, dtype="int32") + path = tmp_path / "scanned.h5" + with h5py.File(path, "w") as file: + file.create_dataset("data", data=data, chunks=(1000,)) + created = [] + original = hdf5_source._filesystem_and_path + + def tracking(urlpath, storage_options=None, filesystem=None): + fs, path = original(urlpath, storage_options, filesystem) + created.append(fs) + return fs, path + + monkeypatch.setattr(hdf5_source, "_filesystem_and_path", tracking) + with _ranged_server(tmp_path) as (urlbase, _requests): + hdf5_source.scan_hdf5_index(f"{urlbase}/{path.name}") + assert created + assert created[0]._session is not None + assert created[0]._session.closed + + +def test_http_hdf5_failed_init_closes_owned_session(tmp_path, monkeypatch): + h5py = pytest.importorskip("h5py") + import blosc2.hdf5_source as hdf5_source + + data = np.arange(10_000, dtype="int32") + path = tmp_path / "failed-init.h5" + with h5py.File(path, "w") as file: + file.create_dataset("data", data=data, chunks=(1000,)) + created = [] + original = hdf5_source._filesystem_and_path + + def tracking(urlpath, storage_options=None, filesystem=None): + fs, path = original(urlpath, storage_options, filesystem) + if filesystem is None: # record only filesystems created by the source + created.append(fs) + return fs, path + + monkeypatch.setattr(hdf5_source, "_filesystem_and_path", tracking) + with _ranged_server(tmp_path) as (urlbase, _requests): + url = f"{urlbase}/{path.name}" + with pytest.raises(ValueError, match="Invalid HDF5 index"): + blosc2.HDF5NDSource(url, "data", hdf5_index={"format": "bad"}) + with pytest.raises(ValueError, match="not found"): + blosc2.HDF5NDSource(url, "missing") + assert len(created) == 2 + assert all(fs._session is None or fs._session.closed for fs in created) + # The scan in the second attempt did open a session, and it must be closed. + assert created[1]._session is not None + assert created[1]._session.closed + + +def test_http_hdf5_source_finalizer_closes_session(tmp_path): + h5py = pytest.importorskip("h5py") + data = np.arange(10_000, dtype="int32") + path = tmp_path / "finalizer.h5" + with h5py.File(path, "w") as file: + file.create_dataset("data", data=data, chunks=(1000,)) + with _ranged_server(tmp_path) as (urlbase, _requests): + source = blosc2.HDF5NDSource(f"{urlbase}/{path.name}", "data") + session = source._filesystem._session + assert not session.closed + del source + gc.collect() + assert session.closed + + def test_http_store_disk_reopen_and_transport_close(tmp_path): h5py = pytest.importorskip("h5py") - pytest.importorskip("kerchunk") - pytest.importorskip("zarr") data = np.arange(10_000, dtype="int32") path = tmp_path / "store.h5" with h5py.File(path, "w") as file: @@ -1573,6 +1670,26 @@ def test_fsspec_ndsource_and_remote_array_storage_options(): assert np.array_equal(proxy[:], a[:]) +def test_fsspec_hdf5_index_selects_hdf5_without_suffix(tmp_path): + import h5py + + from blosc2.hdf5_source import scan_hdf5_index + + path = tmp_path / "indexed.h5" + with h5py.File(path, "w") as file: + file.create_dataset("data", data=np.arange(10, dtype="i4"), chunks=(5,)) + fsspec.filesystem("memory").pipe_file("hdf5-index/container", path.read_bytes()) + url = "memory://hdf5-index/container" + index = scan_hdf5_index(url) + + proxy = blosc2.open(url, dataset="data", hdf5_index=index) + assert isinstance(proxy, blosc2.RemoteArray) + np.testing.assert_array_equal(proxy[:], np.arange(10, dtype="i4")) + + with pytest.raises(ValueError, match="hdf5_index"): + blosc2.open("memory://hdf5-index/other.zarr", lazy=True, hdf5_index=index) + + def test_non_lazy_cache_dir_preserves_explicit_b2z(tmp_path): archive = tmp_path / "hierarchy.b2z" with blosc2.TreeStore(archive, mode="w", threshold=0) as root: @@ -1587,9 +1704,23 @@ def test_non_lazy_cache_dir_preserves_explicit_b2z(tmp_path): np.testing.assert_array_equal(store["/group/a"][:], np.arange(10, dtype="i4")) -def test_non_lazy_cache_dir_rejects_refs(tmp_path): - fsspec.filesystem("memory").pipe_file("nonlazy-refs/data", b"whatever") - with pytest.raises(NotImplementedError, match="refs"): +def test_non_lazy_cache_dir_rejects_hdf5_index(tmp_path): + fsspec.filesystem("memory").pipe_file("nonlazy-index/data", b"whatever") + with pytest.raises(NotImplementedError, match="hdf5_index"): blosc2.open( - "memory://nonlazy-refs/data", lazy=False, cache_dir=tmp_path / "cache", refs={"refs": {}} + "memory://nonlazy-index/data", + lazy=False, + cache_dir=tmp_path / "cache", + hdf5_index={"format": "invalid"}, + ) + # The rejection must not depend on the cache options that follow it. + with pytest.raises(NotImplementedError, match="hdf5_index"): + blosc2.open( + "memory://nonlazy-index/data", + lazy=False, + cache_dir=tmp_path / "cache", + cache_policy=blosc2.CachePolicy.DISK, + hdf5_index={"format": "invalid"}, ) + with pytest.raises(NotImplementedError, match="hdf5_index"): + blosc2.open("memory://nonlazy-index/data", lazy=False, hdf5_index={"format": "invalid"}) diff --git a/tests/test_hdf5_source.py b/tests/test_hdf5_source.py index c5b9205ad..2b6fa6d95 100644 --- a/tests/test_hdf5_source.py +++ b/tests/test_hdf5_source.py @@ -11,6 +11,7 @@ import gc import io import json +import zlib from pathlib import Path import numpy as np @@ -20,8 +21,6 @@ from blosc2.hdf5_source import check_hdf5_dependencies h5py = pytest.importorskip("h5py") -kerchunk = pytest.importorskip("kerchunk") -zarr = pytest.importorskip("zarr") fsspec = pytest.importorskip("fsspec") @@ -42,24 +41,90 @@ def make_memory_h5(name: str = "test.h5", **datasets) -> str: return f"memory://{name}" -def test_hdf5_scan_resets_zarr_resources_after_timeout(monkeypatch): - import blosc2.hdf5_source as hdf5_source +def test_hdf5_native_index(): + from blosc2.hdf5_source import HDF5_INDEX_FORMAT, scan_hdf5_index, validate_hdf5_index + + url = make_memory_h5("native-index.h5", data=(np.arange(12, dtype="i4"), (4,))) + index = scan_hdf5_index(url) + assert index["format"] == HDF5_INDEX_FORMAT + assert index["datasets"]["data"]["direct"] is True + assert len(index["datasets"]["data"]["allocated"]) == 3 + + malformed = json.loads(json.dumps(index)) + malformed["datasets"]["data"]["allocated"] = "not-a-list" + with pytest.raises(ValueError, match="allocation table"): + validate_hdf5_index(malformed) + incomplete = json.loads(json.dumps(index)) + del incomplete["datasets"]["data"]["allocated"] + with pytest.raises(ValueError, match="Incomplete"): + validate_hdf5_index(incomplete) + for field in ("shape", "chunks", "fill_value", "attrs", "dtype"): + incomplete = json.loads(json.dumps(index)) + del incomplete["datasets"]["data"][field] + with pytest.raises(ValueError, match="Incomplete"): + validate_hdf5_index(incomplete) + + +def test_hdf5_object_ndarray_reconstruction(): + from blosc2.hdf5_source import _from_json_value, _json_value + + value = np.empty(2, dtype=object) + value[0] = np.array([1, 2]) + value[1] = np.array([3, 4]) + restored = _from_json_value(_json_value(value)) + assert restored.shape == (2,) + np.testing.assert_array_equal(restored[0], [1, 2]) + np.testing.assert_array_equal(restored[1], [3, 4]) + + ragged = np.empty(2, dtype=object) + ragged[0] = np.array([1, 2]) + ragged[1] = np.array([3, 4, 5]) + restored = _from_json_value(_json_value(ragged)) + assert restored.shape == (2,) + np.testing.assert_array_equal(restored[0], [1, 2]) + np.testing.assert_array_equal(restored[1], [3, 4, 5]) + - attempts = [] - resets = [] +def test_hdf5_titled_dtype_roundtrip(): + from blosc2.hdf5_source import dtype_from_value, dtype_value - def translate(*args): - attempts.append(None) - if len(attempts) == 1: - raise TimeoutError("wedged sync bridge") - return {"refs": {}} + dtype = np.dtype([(("title", "value"), "i4"), np.arange(20, dtype=">i4")), + ( + np.dtype([("value", "