Skip to content

Commit 0de6077

Browse files
authored
fix(indexing): delegate source tokenization to Dask (#4348)
* docs(indexing): ground design and integration claims in current behavior Assisted-by: Codex:GPT-6 * docs(indexing): correct reader lazy-array and cache contracts Assisted-by: Codex:GPT-6 * fix(indexing): validate wire boundaries and clarify format contracts Assisted-by: Codex:GPT-6 * fix(indexing): validate selector bounds and shared dependencies Correct mathematical API documentation to match supported coordinate, grid, and chunk projection contracts. Assisted-by: Codex:GPT-6 * docs(indexing): reconcile reader contracts and record audit fixes Assisted-by: Codex:GPT-6 * docs(indexing): reconcile audit with current partition implementation Retain the existing unsigned selector fix and update the unsupported mixed-dependency error assertion for general intersection routing. Assisted-by: Codex:GPT-6 * docs(indexing): clarify planning coverage and benchmark measurement boundaries Assisted-by: Codex:GPT-6 * fix(indexing): group signed chunk coordinates without collisions Use lexicographic tuple grouping when chunk indices contain negative values. Cover shared one-axis and two-axis array dependencies, repeated points, and extreme signed coordinates. Assisted-by: Codex:GPT-6 * docs(indexing): state remaining planner limits precisely Assisted-by: Codex:GPT-6 * fix(indexing): define explicit source token contract Assisted-by: Codex:GPT-6 * fix(indexing): reject mmap-backed token buffers Assisted-by: Codex:GPT-6 * docs(indexing): number audit changelog entries for PR 4345 Assisted-by: Codex:GPT-6 * docs(indexing): describe current contracts in docstrings Remove implementation history and unsupported historical claims from source and test docstrings. Distinguish immutable coordinate mappings from mutable source values. Assisted-by: Codex:GPT-6 * fix(indexing): delegate source tokenization to Dask Assisted-by: Codex:GPT-6
1 parent 34dd17c commit 0de6077

5 files changed

Lines changed: 88 additions & 96 deletions

File tree

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
Delegate LazyArray source tokenization to Dask, honoring its registered normalizers, source hooks, and deterministic-token requirements. Remove local content-hashing and UUID fallbacks. Dask remains optional for indexing and reading.

‎packages/zarr-indexing/docs/guide/integrations.md‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -243,3 +243,14 @@ that implementation or reproduce its full worker/GPU lifecycle.
243243
·
244244
**API:** [API reference](../api/index.md)
245245
</nav>
246+
247+
## Dask tokenization
248+
249+
`LazyArray.__dask_tokenize__()` combines Dask's token for the wrapped source
250+
with the serialized view transform. Dask owns source hashing, registered
251+
normalizers, custom source hooks, and deterministic-token requirements.
252+
Tokenization may read or hash source values. Dask is optional for indexing
253+
and reading, but required when requesting a Dask token.
254+
255+
The reader and partitioning are omitted because they must preserve values.
256+
Changing a source after graph construction does not update existing Dask keys.

‎packages/zarr-indexing/pyproject.toml‎

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -99,9 +99,8 @@ extend = "../../pyproject.toml"
9999
target-version = "py312"
100100

101101
[tool.ruff.lint.per-file-ignores]
102-
# Chunk discovery and __dask_tokenize__ deliberately catch Exception: a
103-
# foreign source's attributes or token hooks may fail, so these paths provide
104-
# fallback metadata or tokens when ordinary exceptions occur. Configured here (not as
102+
# Chunk discovery deliberately catches Exception: a foreign source's
103+
# attributes may fail, so this path provides fallback metadata. Configured here (not as
105104
# noqa comments) because different ruff versions
106105
# have differed on whether these rules fire; RUF100 can remove unused noqa comments.
107106
"src/zarr_indexing/lazy_array.py" = ["BLE001", "S110"]

‎packages/zarr-indexing/src/zarr_indexing/lazy_array.py‎

Lines changed: 10 additions & 72 deletions
Original file line numberDiff line numberDiff line change
@@ -12,8 +12,8 @@
1212
```
1313
1414
Selection construction does not read source values: `result()`, `__array__`,
15-
and eager `__getitem__` perform reads. Tokenization can also read or hash source
16-
data, depending on the source and tokenization path. `.lazy` operations inspect
15+
and eager `__getitem__` perform reads. Source tokenization is delegated to
16+
Dask and may inspect source values. `.lazy` operations inspect
1717
selection metadata and may copy or process supplied index arrays. Composition
1818
does not accumulate wrapper layers: a view of a view is still a single transform
1919
and retains its reader.
@@ -138,7 +138,6 @@
138138
import json
139139
import math
140140
import operator
141-
import uuid
142141
from collections.abc import Sequence
143142
from dataclasses import dataclass, field
144143
from typing import TYPE_CHECKING, Any, Protocol, cast
@@ -173,10 +172,6 @@
173172

174173
__all__ = ["LazyArray", "Partition"]
175174

176-
# Above this declared byte count, the no-Dask token fallback adds a fresh UUID
177-
# instead of digesting contents. See `_wrapped_token`.
178-
_TOKEN_DIGEST_LIMIT = 1 << 20
179-
180175

181176
def _invoke_reader(
182177
reader: Reader,
@@ -525,66 +520,6 @@ def _validate_prepared_parts(parts: Sequence[Partition], out_shape: tuple[int, .
525520
raise ValueError("prepared parts do not tile the view exactly")
526521

527522

528-
# --------------------------------------------------------------------------- #
529-
# Tokenization
530-
# --------------------------------------------------------------------------- #
531-
532-
533-
def _wrapped_token(array: Any) -> Any:
534-
"""A token for the wrapped array.
535-
536-
In order of preference: the array's own `__dask_tokenize__`;
537-
`dask.base.tokenize` when dask is importable (imported lazily — this package
538-
never requires it); otherwise a local fallback that digests the contents of
539-
a small array.
540-
541-
Tokens can differ depending on whether Dask is available and on the source
542-
hook. This fallback does not provide a portable content identifier.
543-
544-
Above `_TOKEN_DIGEST_LIMIT`, or when conversion is unavailable, the local
545-
fallback adds a fresh UUID on each call, so repeated calls normally differ.
546-
The returned tuple is still equal to itself. Below the limit, conversion
547-
can read the source, and the hash is of its NumPy buffer bytes. Object-array
548-
buffer bytes contain object references, not a recursive content snapshot.
549-
"""
550-
hook = getattr(array, "__dask_tokenize__", None)
551-
if hook is not None:
552-
try:
553-
return hook()
554-
# A failing source hook falls through to the remaining tokenization paths.
555-
except Exception: # pragma: no cover - a hook that refuses to run
556-
pass
557-
try:
558-
# dask is an optional peer, never a dependency of this package, so it is
559-
# imported here and its absence is ordinary.
560-
from dask.base import tokenize # pyright: ignore[reportMissingImports]
561-
except ImportError:
562-
pass
563-
else:
564-
return tokenize(array)
565-
566-
shape = tuple(int(s) for s in getattr(array, "shape", ()))
567-
dtype = getattr(array, "dtype", None)
568-
structural = (type(array).__qualname__, shape, str(dtype))
569-
# A fresh identifier per call when contents cannot be identified. It
570-
# is the shape and dtype that would otherwise be mistaken for an identity,
571-
# so they are kept alongside it for a reader looking at a graph.
572-
unidentified = (*structural, "unidentified", uuid.uuid4().hex)
573-
574-
# Decide whether to digest the contents from the *declared* size. Measuring
575-
# it by converting first would read the whole array — a multi-gigabyte store
576-
# pulled into memory by a token call, which is the opposite of the point.
577-
itemsize = getattr(dtype, "itemsize", None)
578-
if not isinstance(itemsize, int) or itemsize * math.prod(shape) > _TOKEN_DIGEST_LIMIT:
579-
return unidentified
580-
try:
581-
contents = np.ascontiguousarray(array)
582-
# Failed NumPy conversion leaves this source unidentified.
583-
except Exception:
584-
return unidentified
585-
return (*structural, hashlib.sha256(contents.tobytes()).hexdigest())
586-
587-
588523
# --------------------------------------------------------------------------- #
589524
# The wrapper
590525
# --------------------------------------------------------------------------- #
@@ -1249,15 +1184,18 @@ def __dask_tokenize__(self) -> Any:
12491184
tokens produce equal tokens; arbitrary semantically equivalent mappings
12501185
are not guaranteed to serialize identically.
12511186
1252-
Source tokenization can read or hash data and need not be deterministic
1253-
on every fallback path; see `_wrapped_token`. The reader and partitioning
1254-
are omitted under the contract that they preserve values. Cache users
1255-
must also account for source mutation and the source's token semantics.
1187+
Dask tokenizes the wrapped source using its normal dispatch and
1188+
determinism policy. This may read or hash source values. Dask is
1189+
imported only when this method is called and is otherwise optional.
1190+
The reader and partitioning are omitted because they must preserve
1191+
values. Mutating a source does not update keys in existing Dask graphs.
12561192
"""
1193+
from dask.base import tokenize # pyright: ignore[reportMissingImports]
1194+
12571195
canonical = json.dumps(self._transform.to_json(), sort_keys=True)
12581196
return (
12591197
type(self).__qualname__,
1260-
_wrapped_token(self._array),
1198+
tokenize(self._array),
12611199
hashlib.sha256(canonical.encode()).hexdigest(),
12621200
)
12631201

‎packages/zarr-indexing/tests/test_lazy_array.py‎

Lines changed: 64 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -1776,6 +1776,7 @@ def test_nonfirst_partition_transform_directly_addresses_its_array() -> None:
17761776

17771777

17781778
def test_partition_token_encodes_its_public_global_transform() -> None:
1779+
pytest.importorskip("dask.base")
17791780
source = np.arange(8)
17801781
base = LazyArray.from_numpy(source)
17811782
partition_view = list(base.with_parts((4,)).parts())[1].view
@@ -1907,6 +1908,7 @@ def test_with_parts_validates_strictly(parts: Any, match: str) -> None:
19071908

19081909
def test_dask_token_is_deterministic_and_discriminating() -> None:
19091910
"""Same data and same view token alike; a different selection differs."""
1911+
pytest.importorskip("dask.base")
19101912
data = reference()
19111913
base = LazyArray(data)
19121914
assert base.__dask_tokenize__() == LazyArray(reference()).__dask_tokenize__()
@@ -1923,6 +1925,7 @@ def test_dask_token_is_deterministic_and_discriminating() -> None:
19231925

19241926

19251927
def test_reader_and_partitioning_do_not_change_dask_identity() -> None:
1928+
pytest.importorskip("dask.base")
19261929
base = LazyArray(reference())
19271930
token = base.__dask_tokenize__()
19281931
assert base.with_reader(numpy_reader).__dask_tokenize__() == token
@@ -2000,7 +2003,6 @@ def test_pickle_round_trip() -> None:
20002003
view = LazyArray(reference()).with_parts((2, 2, 2)).lazy[1:6, ::2].lazy.oindex[[3, 0, 0], :, :]
20012004
restored = pickle.loads(pickle.dumps(view))
20022005
assert restored.shape == view.shape
2003-
assert restored.__dask_tokenize__() == view.__dask_tokenize__()
20042006
np.testing.assert_array_equal(np.asarray(restored.result()), np.asarray(view.result()))
20052007

20062008

@@ -2010,7 +2012,7 @@ def test_pickle_round_trip() -> None:
20102012

20112013

20122014
def test_dask_from_array_roundtrip() -> None:
2013-
"""A `LazyArray` is a drop-in dask source — no translation ceremony."""
2015+
"""Dask can tokenize and read a wrapper over a Zarr source."""
20142016
da = pytest.importorskip("dask.array")
20152017
source = make_source("zarr")
20162018

@@ -2435,25 +2437,6 @@ def test_a_masked_source_keeps_its_mask_when_the_view_is_empty(parts: Any) -> No
24352437
assert np.asarray(got).shape == (3, 0), parts
24362438

24372439

2438-
def test_a_large_array_without_dask_refuses_to_claim_equality(
2439-
monkeypatch: pytest.MonkeyPatch,
2440-
) -> None:
2441-
"""Without Dask, arrays above the digest limit receive distinct fallback tokens."""
2442-
import sys
2443-
2444-
monkeypatch.setitem(sys.modules, "dask.base", None)
2445-
big = np.zeros(1 << 19, dtype=np.int64)
2446-
other = big.copy()
2447-
other[0] = 1
2448-
2449-
assert LazyArray(big).__dask_tokenize__() != LazyArray(other).__dask_tokenize__()
2450-
assert LazyArray(big).__dask_tokenize__() != LazyArray(big).__dask_tokenize__()
2451-
2452-
# Below the limit the contents are digested, so equal data still tokens alike.
2453-
small = np.zeros(8, dtype=np.int64)
2454-
assert LazyArray(small).__dask_tokenize__() == LazyArray(small.copy()).__dask_tokenize__()
2455-
2456-
24572440
# ---------------------------------------------------------------------------
24582441
# Completeness and partition spellings
24592442
# ---------------------------------------------------------------------------
@@ -2576,3 +2559,63 @@ def test_fancy_composition_over_an_empty_axis() -> None:
25762559
scalar = composed.lazy.vindex[..., np.array(1)]
25772560
assert scalar.shape == (2, 0)
25782561
assert np.asarray(scalar.result()).shape == (2, 0)
2562+
2563+
2564+
@pytest.mark.parametrize("kind", ["numpy", "object", "masked", "registered", "hook"])
2565+
def test_source_token_uses_dask_policy(kind: str) -> None:
2566+
dask_base = pytest.importorskip("dask.base")
2567+
2568+
class RegisteredArray(ForeignArray):
2569+
pass
2570+
2571+
class VersionedArray(ForeignArray):
2572+
def __dask_tokenize__(self) -> Any:
2573+
return ("versioned-source", 1)
2574+
2575+
dask_base.normalize_token.register(RegisteredArray, lambda source: ("registered-source", 1))
2576+
data = np.arange(4)
2577+
sources = {
2578+
"numpy": data,
2579+
"object": data.astype(object),
2580+
"masked": np.ma.masked_greater(data, 2),
2581+
"registered": RegisteredArray(data, None),
2582+
"hook": VersionedArray(data, None),
2583+
}
2584+
source = sources[kind]
2585+
assert LazyArray(source).__dask_tokenize__()[1] == dask_base.tokenize(source)
2586+
2587+
2588+
def test_source_token_preserves_dask_determinism_requirement() -> None:
2589+
dask_base = pytest.importorskip("dask.base")
2590+
dask_tokenize = pytest.importorskip("dask.tokenize")
2591+
2592+
class UnserializableArray(ForeignArray):
2593+
def __reduce_ex__(self, protocol: int) -> Any:
2594+
raise TypeError("cannot serialize source")
2595+
2596+
source = UnserializableArray(np.arange(4), None)
2597+
for value in (source, LazyArray(source)):
2598+
with pytest.raises(dask_tokenize.TokenizationError):
2599+
dask_base.tokenize(value, ensure_deterministic=True)
2600+
2601+
2602+
def test_source_token_preserves_hook_failure() -> None:
2603+
pytest.importorskip("dask.base")
2604+
2605+
class RefusingArray(ForeignArray):
2606+
def __dask_tokenize__(self) -> Any:
2607+
raise RuntimeError("source version unavailable")
2608+
2609+
with pytest.raises(RuntimeError, match="source version unavailable"):
2610+
LazyArray(RefusingArray(np.arange(4), None)).__dask_tokenize__()
2611+
2612+
2613+
def test_dask_is_only_required_for_tokenization(monkeypatch: pytest.MonkeyPatch) -> None:
2614+
import sys
2615+
2616+
monkeypatch.setitem(sys.modules, "dask.base", None)
2617+
data = np.arange(4)
2618+
view = LazyArray(data).lazy[1:]
2619+
np.testing.assert_array_equal(view.result(), data[1:])
2620+
with pytest.raises(ModuleNotFoundError):
2621+
view.__dask_tokenize__()

0 commit comments

Comments
 (0)