From af5eb6c37fdbe60950f2144d7fa5df86a85a59fa Mon Sep 17 00:00:00 2001 From: Francesc Alted Date: Fri, 11 Sep 2026 09:24:47 +0200 Subject: [PATCH 1/2] Make translate_refs_serializable walk nested values zarr v3 memory stores hand out Buffer objects, and this helper only converted the ones sitting at the top level of the reference dict. A Buffer nested in a dict or list value survived, and consumers then failed to serialize the reference set ("can not serialize 'Buffer' object"). Walk dicts and lists too, still stripping a leading slash from the keys of top-level buffers. --- kerchunk/utils.py | 29 ++++++++++++++++++++--------- tests/test_utils.py | 19 +++++++++++++++++++ 2 files changed, 39 insertions(+), 9 deletions(-) diff --git a/kerchunk/utils.py b/kerchunk/utils.py index 7d9dae6a..6a18468e 100644 --- a/kerchunk/utils.py +++ b/kerchunk/utils.py @@ -569,11 +569,25 @@ def templateize(strings, min_length=10, template_name="u"): return template, strings +def _serializable(value): + """Replace zarr buffers with bytes, however deeply they are nested.""" + if isinstance(value, zarr.core.buffer.cpu.Buffer): + return value.to_bytes() + if isinstance(value, dict): + return {k: _serializable(v) for k, v in value.items()} + if isinstance(value, list): + return [_serializable(v) for v in value] + return value + + def translate_refs_serializable(refs: dict): """Translate a reference set to a serializable form, given that zarr v3 memory stores store data in buffers by default. This modifies the input dictionary in place, and returns a reference to it. + Buffers can also hide inside nested values, so the conversion walks + them instead of only looking at the top level. + It also fixes keys that have a leading slash, which is not appropriate for zarr v3 keys @@ -587,14 +601,11 @@ def translate_refs_serializable(refs: dict): dict A serializable form of the reference set """ - keys_to_remove = [] - new_keys = {} - for k, v in refs.items(): + for k in list(refs): + v = refs[k] if isinstance(v, zarr.core.buffer.cpu.Buffer): - key = k.removeprefix("/") - new_keys[key] = v.to_bytes() - keys_to_remove.append(k) - for k in keys_to_remove: - del refs[k] - refs.update(new_keys) + del refs[k] + refs[k.removeprefix("/")] = v.to_bytes() + else: + refs[k] = _serializable(v) return refs diff --git a/tests/test_utils.py b/tests/test_utils.py index 5b686a67..0028abe7 100644 --- a/tests/test_utils.py +++ b/tests/test_utils.py @@ -183,3 +183,22 @@ def test_encode_fill_value(): assert kerchunk.utils.encode_fill_value(np.array(9999), np.dtype("int")) == 9999 assert kerchunk.utils.encode_fill_value([9999], np.dtype("int")) == 9999 assert kerchunk.utils.encode_fill_value(9999, np.dtype("int")) == 9999 + + +def test_translate_refs_serializable_nested_buffers(): + from zarr.core.buffer.core import default_buffer_prototype + + buffer = default_buffer_prototype().buffer.from_bytes(b"{}") + refs = { + "/a/.zarray": buffer, + "b": {"data": buffer}, + "c": [buffer, "text"], + } + + out = kerchunk.utils.translate_refs_serializable(refs) + + # the top-level buffer also loses its leading slash + assert out["a/.zarray"] == b"{}" + assert "/a/.zarray" not in out + assert out["b"] == {"data": b"{}"} + assert out["c"] == [b"{}", "text"] From 5a9d5bdc8bdfd84b3fb02339974fc64c7f249dde Mon Sep 17 00:00:00 2001 From: Francesc Alted Date: Fri, 11 Sep 2026 09:24:47 +0200 Subject: [PATCH 2/2] Do not quash a TimeoutError from the store sync bridge _translator catches every exception per node and, with the default error="warn", downgrades it to a warning and moves on. That is wrong for a timeout: it means the store (or zarr's sync bridge behind it) stalled, so the next call fails the same way and a caller who bounded the bridge still never gets a finished translation. Let TimeoutError through regardless of the error mode, so one bounded failure aborts the translation instead of multiplying. --- kerchunk/hdf.py | 7 +++++++ tests/test_hdf.py | 28 ++++++++++++++++++++++++++++ 2 files changed, 35 insertions(+) diff --git a/kerchunk/hdf.py b/kerchunk/hdf.py index 4b115dfa..5f9f2fda 100644 --- a/kerchunk/hdf.py +++ b/kerchunk/hdf.py @@ -604,6 +604,13 @@ def _translator( lggr.debug(f"HDF5 group: {h5obj.name}") zgrp = self._zroot.require_group(h5obj.name.lstrip("/")) self._transfer_attrs(h5obj, zgrp) + except TimeoutError: + # A timeout is not a property of this node: it means the store's + # sync bridge (or the transport behind it) stalled, and every + # later call will stall the same way. Quashing it here would + # turn one bounded failure into a hang, so let it abort the + # translation regardless of the error mode. + raise except Exception as e: import traceback diff --git a/tests/test_hdf.py b/tests/test_hdf.py index afaf6be1..60dacc56 100644 --- a/tests/test_hdf.py +++ b/tests/test_hdf.py @@ -448,3 +448,31 @@ def test_malicious_chunks(): store = ReferenceFileSystem(ref, target=fname, asynchronous=True).get_mapper() data = zarr.group(store)["data"][:] assert (data == np.arange(8, dtype=np.int32)).all() + + +def test_timeout_is_not_quashed(tmp_path, monkeypatch): + """A stalled store must abort the translation, not be warned away.""" + import h5py + import zarr.core.sync + + path = tmp_path / "timeout.h5" + with h5py.File(path, "w") as f: + f.attrs["title"] = "root" + f.create_dataset("a", data=np.arange(10)) + + real_sync = zarr.core.sync.SyncMixin._sync + calls = {"n": 0} + + def stalling_sync(self, coroutine, *args, **kwargs): + calls["n"] += 1 + if calls["n"] > 1: # call 1 transfers the root attributes + coroutine.close() + raise TimeoutError("store sync bridge stalled") + return real_sync(self, coroutine, *args, **kwargs) + + with fsspec.open(path, "rb") as f: + chunks = SingleHdf5ToZarr(f, url=str(path), error="warn") + monkeypatch.setattr(zarr.core.sync.SyncMixin, "_sync", stalling_sync) + calls["n"] = 0 + with pytest.raises(TimeoutError): + chunks.translate()