diff --git a/Justfile b/Justfile index 6e56644207..b3f187e231 100644 --- a/Justfile +++ b/Justfile @@ -57,7 +57,7 @@ coverage-serve *args: # Run slow Hypothesis tests and write coverage.xml hypothesis *args: - hatch run {{ quote(hatch_env) }}:coverage run --source=src -m pytest -nauto --run-slow-hypothesis tests/test_properties.py tests/test_store/test_stateful* "$@" + hatch run {{ quote(hatch_env) }}:coverage run --source=src -m pytest -nauto --run-slow-hypothesis tests/test_properties.py tests/test_store/test_stateful* tests/test_array_stateful.py "$@" hatch run {{ quote(hatch_env) }}:coverage xml # Validate executable documentation code blocks diff --git a/changes/4334.bugfix.md b/changes/4334.bugfix.md new file mode 100644 index 0000000000..ea7ec480ce --- /dev/null +++ b/changes/4334.bugfix.md @@ -0,0 +1,7 @@ +A chunk edge length is now always at least 1, while an array extent may be 0. Every spelling of one chunk spanning an axis (`chunks=-1`, `chunks=False`, `chunks="auto"`, `shards=-1`, `shards=False`) gives chunk size 1 on a zero-length axis, and a shard spanning an axis is a multiple of the inner chunk size, so `shards=-1` no longer fails when the axis length is not a multiple of it. Rectilinear chunk grids can be created on a zero-length dimension with a non-empty list of positive chunk sizes, which are kept for later growth. + +`FixedDimension(size=0, ...)` now raises a `ValueError`. `ArrayV2Metadata(chunks=(0,))` is still accepted and written as given; an array built from such metadata (with `create_hierarchy`, for example) reads the chunk size as a stored chunk size of 0 is read (below). `create_hierarchy` now builds every array and group before it deletes or stores anything, so a node that cannot be built fails with the store untouched. The regular chunk grid metadata class reads a `bool` chunk edge length as the `int` it equals, and rejects a string or a mapping as a chunk shape as a whole. `RectilinearChunkGridMetadata`, which is experimental, now reads a `bool` edge length as the `int` it equals and rejects a NumPy integer or a float edge length such as `4.0` with a `TypeError`; it used to keep them as given, so that it could not store a NumPy integer, could not read back a `bool`, and stored a float as a JSON float, which the Zarr specification does not allow. + +Stored metadata with a regular chunk size of 0 or JSON `false` (Zarr format 2 `chunks`, Zarr format 3 `regular` `chunk_shape`) is now read as one chunk spanning the axis (a multiple of the inner chunk size for sharded arrays, which previously failed to open). zarr-python wrote such sizes for arrays created with a zero-length axis until 3.4, and, from 2.18.7 to 3.2.1, for an explicit chunk size of 0 or `False` on an axis of any length, as in `chunks=(0,)` (those arrays could store no data). On a zero-length axis this opens silently, as before; appending to such an axis then stores chunks of size 1. On an axis of positive length it opens with a `ZarrUserWarning`, which says that the axis holds only the fill value and how to store valid metadata: `array.update_attributes({})`, then `zarr.consolidate_metadata` if the metadata is consolidated. JSON `true`, as zarr-python 3.0 and 3.2 wrote for a chunk size of `True`, is read as 1, silently, in the chunk grid (regular, and the explicit edges and run-length encoded sizes of a rectilinear one) and in the inner chunk shape of every sharding codec, nested or not. A stored rectilinear chunk grid whose edge lengths are integral JSON floats (`[[4.0, 2]]`), as zarr-python 3.2 wrote for float edges, is read with those edges as integers, silently; a float anywhere no release wrote one (a regular chunk shape, the inner chunk shape of a sharding codec, a run-length repeat count, an edge below 1) is rejected. + +Writing data to an array read from such metadata (other than an empty selection) first stores the upgrade of the metadata the store then holds, unless another writer has stored valid metadata since, so that other readers find the chunks written; in a `ZipStore` this adds a second entry for the metadata document, as every metadata update does. If the metadata the store then holds lays out chunks differently (another writer resized the array keeping the chunk size of 0, say), the write raises a `ValueError` asking to reopen the array, and stores nothing. Apart from the array's own metadata writes (`update_attributes`, `resize`), which store the upgrade as before, this is the only operation that stores it: a group's consolidated metadata keeps such an array's metadata as it was stored (`zarr.consolidate_metadata` copies it as the array's own document holds it) until the array stores its upgrade through the group's handle, and changing a group reads and writes no metadata of its members. diff --git a/docs/user-guide/arrays.md b/docs/user-guide/arrays.md index 140bc4bd09..e8df79a102 100644 --- a/docs/user-guide/arrays.md +++ b/docs/user-guide/arrays.md @@ -707,6 +707,12 @@ z.append(np.arange(10, dtype='float64')) print(f"After append: shape={z.shape}, chunk_sizes={z.write_chunk_sizes}") ``` +A rectilinear array can also be created with a zero-length dimension: because no +non-empty list of positive chunk sizes can sum to 0, the chunk sizes given for such a +dimension are stored as-is and describe the chunks the dimension will grow into +on `append` or `resize` — the same state as resizing an existing rectilinear +dimension down to 0. + ### Compressors and filters Rectilinear arrays work with all codecs — compressors, filters, and checksums. diff --git a/src/zarr/core/_json.py b/src/zarr/core/_json.py index efe8152a4f..66f6236f44 100644 --- a/src/zarr/core/_json.py +++ b/src/zarr/core/_json.py @@ -41,6 +41,12 @@ def buffer_to_json(buffer: Buffer) -> JSON: return cast("JSON", json.loads(buffer.to_bytes())) +def json_equal(a: JSON, b: JSON) -> bool: + """Whether two JSON values have the same JSON encoding. Python compares `True` and + `1`, or `1.0` and `1`, as equal; JSON does not.""" + return json.dumps(a) == json.dumps(b) + + def buffer_to_json_object(buffer: Buffer) -> dict[str, JSON]: """Parse the contents of a `Buffer` as a JSON object (a `dict`). diff --git a/src/zarr/core/array.py b/src/zarr/core/array.py index 5a8d6bf57e..e402382a33 100644 --- a/src/zarr/core/array.py +++ b/src/zarr/core/array.py @@ -118,7 +118,13 @@ ArrayV2MetadataDict, ArrayV3Metadata, ) -from zarr.core.metadata.io import save_metadata +from zarr.core.metadata.io import ( + ARRAY_DOCUMENTS, + parse_stored_array, + read_documents, + save_metadata, + upsert_metadata, +) from zarr.core.metadata.v2 import ( CompressorLikev2, get_object_codec_id, @@ -201,13 +207,22 @@ def _chunk_sizes_from_shape( return tuple(result) -def parse_array_metadata(data: Any) -> ArrayMetadata: +def parse_array_metadata(data: Any, path: str | None = None) -> ArrayMetadata: + """Array metadata from a metadata object or a metadata document, naming the array at + `path` in warnings about how an invalid document was read. + + `ArrayV2Metadata` accepts a chunk size of 0, as it always has, though only an + invalid document holds one: such metadata is read as the documents it would store + are (see `zarr.core.metadata.upgrades`), so an array can be built from it. No data + was read or written under that chunk size, so the reading is silent.""" + if isinstance(data, ArrayV2Metadata) and 0 in data.chunks: + return parse_stored_array(data.to_buffer_dict(default_buffer_prototype()), 2) if isinstance(data, ArrayMetadata): return data - elif isinstance(data, dict): + if isinstance(data, dict): zarr_format = data.get("zarr_format") if zarr_format == 3: - meta_out = ArrayV3Metadata.from_dict(data) + meta_out = ArrayV3Metadata.from_dict(data, path=path) if len(meta_out.storage_transformers) > 0: msg = ( f"Array metadata contains storage transformers: {meta_out.storage_transformers}." @@ -216,7 +231,7 @@ def parse_array_metadata(data: Any) -> ArrayMetadata: raise ValueError(msg) return meta_out elif zarr_format == 2: - return ArrayV2Metadata.from_dict(data) + return ArrayV2Metadata.from_dict(data, path=path) else: raise ValueError(f"Invalid zarr_format: {zarr_format}. Expected 2 or 3") raise TypeError # pragma: no cover @@ -404,7 +419,7 @@ def __init__( store_path: StorePath, config: ArrayConfigLike | None = None, ) -> None: - metadata_parsed = parse_array_metadata(metadata) + metadata_parsed = parse_array_metadata(metadata, str(store_path)) config_parsed = parse_array_config(config) object.__setattr__(self, "metadata", metadata_parsed) @@ -765,7 +780,7 @@ def from_dict( ValueError If the dictionary data is invalid or incompatible with either Zarr format 2 or 3 array creation. """ - metadata = parse_array_metadata(data) + metadata = parse_array_metadata(data, str(store_path)) return cls(metadata=metadata, store_path=store_path) @classmethod @@ -1610,10 +1625,50 @@ async def get_coordinate_selection( return out_array async def _save_metadata(self, metadata: ArrayMetadata, ensure_parents: bool = False) -> None: - """ - Asynchronously save the array metadata. - """ + """Store `metadata` as this array's own documents, then clear the + `_stored_document` mark (see `_stored_document_replaced`).""" await save_metadata(self.store_path, metadata, ensure_parents=ensure_parents) + self._stored_document_replaced() + + def _stored_document_replaced(self) -> None: + """Record that the store no longer holds a document of this array that needs an + upgrade: it holds the upgrade, a valid document, or none. The metadata this handle + holds, which a consolidated group handle may share, then stops standing for the + document it was read from (see `mark_upgraded`), so no later write through either + handle stores that document again.""" + object.__setattr__(self.metadata, "_stored_document", None) + + async def _store_upgraded_document(self) -> None: + """Store the upgrade of this array's current stored document, if it needs one, + before chunks are written under this handle's metadata. + + Only for metadata read from a document that had to be upgraded (see + `zarr.core.metadata.upgrades`). The document is read again, because the store may + hold a newer one than this handle's metadata. If that one lays out chunks + differently (the array was resized since by software that kept the invalid chunk + size), this handle would write chunks no reader finds, so it raises and stores + nothing. If it needs no upgrade (the array was re-saved since, possibly by another + implementation), it is left as written; if there is none, there is nothing to + upgrade. Storing the same upgrade twice is harmless, so concurrent callers need + no coordination. + """ + if self.metadata._stored_document is None: + return + zarr_format = self.metadata.zarr_format + documents = await read_documents(self.store_path, ARRAY_DOCUMENTS[zarr_format]) + try: + current = parse_stored_array(documents, zarr_format) + except ArrayNotFoundError: + pass + else: + if _chunk_layout(current) != _chunk_layout(self.metadata): + raise ValueError( + f"The metadata stored for the array at {str(self.store_path)!r} has " + "changed since this array was opened: reopen the array to write to it." + ) + if current._stored_document is not None: + await upsert_metadata(self.store_path, current, documents) + self._stored_document_replaced() async def _set_selection( self, @@ -1623,6 +1678,10 @@ async def _set_selection( prototype: BufferPrototype, fields: Fields | None = None, ) -> None: + if product(indexer.shape) > 0: + # Chunks are about to be stored under the upgraded metadata, so store it + # first: every reader of the store then agrees with them. + await self._store_upgraded_document() return await _set_selection( self.store_path, self.metadata, @@ -1674,16 +1733,10 @@ async def setitem( - This method is asynchronous and should be awaited. - Supports basic indexing, where the selection is contiguous and does not involve advanced indexing. """ - return await _setitem( - self.store_path, - self.metadata, - self.codec_pipeline, - self.config, - self._chunk_grid, - selection, - value, - prototype=prototype, - ) + if prototype is None: + prototype = default_buffer_prototype() + indexer = BasicIndexer(selection, shape=self.metadata.shape, chunk_grid=self._chunk_grid) + return await self._set_selection(indexer, value, prototype=prototype) @property def oindex(self) -> AsyncOIndex[T_ArrayMetadata]: @@ -4846,6 +4899,16 @@ async def create_array( ) +def _chunk_layout( + metadata: ArrayMetadata, +) -> tuple[tuple[int, ...] | ChunkGridMetadata, tuple[int, ...] | None]: + """How an array's chunks are laid out: its chunk grid and, if it is sharded, the + inner chunk shape.""" + grid = metadata.chunks if isinstance(metadata, ArrayV2Metadata) else metadata.chunk_grid + sharding = _sharding_codec(metadata) + return grid, None if sharding is None else sharding.chunk_shape + + def _sharding_codec(metadata: ArrayMetadata) -> ShardingCodec | None: """The array's sharding codec, or None if the array is not sharded. @@ -5827,58 +5890,6 @@ async def _set_selection( ) -async def _setitem( - store_path: StorePath, - metadata: ArrayMetadata, - codec_pipeline: CodecPipeline, - config: ArrayConfig, - chunk_grid: ChunkGrid, - selection: BasicSelection, - value: npt.ArrayLike, - prototype: BufferPrototype | None = None, -) -> None: - """ - Set values in the array using basic indexing. - - Parameters - ---------- - store_path : StorePath - The store path of the array. - metadata : ArrayMetadata - The array metadata. - codec_pipeline : CodecPipeline - The codec pipeline for encoding/decoding. - config : ArrayConfig - The array configuration. - chunk_grid : ChunkGrid - The chunk grid. - selection : BasicSelection - The selection defining the region of the array to set. - value : npt.ArrayLike - The values to be written into the selected region of the array. - prototype : BufferPrototype or None, optional - A prototype buffer that defines the structure and properties of the array chunks being modified. - If None, the default buffer prototype is used. - """ - if prototype is None: - prototype = default_buffer_prototype() - indexer = BasicIndexer( - selection, - shape=metadata.shape, - chunk_grid=chunk_grid, - ) - return await _set_selection( - store_path, - metadata, - codec_pipeline, - config, - chunk_grid, - indexer, - value, - prototype=prototype, - ) - - async def _resize( array: AsyncArray[ArrayV2Metadata] | AsyncArray[ArrayV3Metadata], new_shape: ShapeLike, @@ -5928,7 +5939,7 @@ async def _delete_key(key: str) -> None: ) # Write new metadata - await save_metadata(array.store_path, new_metadata) + await array._save_metadata(new_metadata) # Update metadata and chunk_grid (in place) object.__setattr__(array, "metadata", new_metadata) @@ -5993,15 +6004,7 @@ async def _append( slice(None) if i != axis else slice(old_shape[i], new_shape[i]) for i in range(len(array.shape)) ) - await _setitem( - array.store_path, - array.metadata, - array.codec_pipeline, - array.config, - array._chunk_grid, - append_selection, - data, - ) + await array.setitem(append_selection, data) return new_shape @@ -6028,7 +6031,7 @@ async def _update_attributes( array.metadata.attributes.update(new_attributes) # Write new metadata - await save_metadata(array.store_path, array.metadata) + await array._save_metadata(array.metadata) return array diff --git a/src/zarr/core/chunk_grids.py b/src/zarr/core/chunk_grids.py index cc27366027..399943f72d 100644 --- a/src/zarr/core/chunk_grids.py +++ b/src/zarr/core/chunk_grids.py @@ -49,20 +49,17 @@ class FixedDimension: """Uniform chunk size. Boundary chunks contain less data but are encoded at full size by the codec pipeline.""" - size: int # chunk edge length (>= 0) - extent: int # array dimension length + size: int # chunk edge length (>= 1) + extent: int # array dimension length (>= 0) nchunks: int = field(init=False, repr=False) ngridcells: int = field(init=False, repr=False) def __post_init__(self) -> None: - if self.size < 0: - raise ValueError(f"FixedDimension size must be >= 0, got {self.size}") + if self.size < 1: + raise ValueError(f"FixedDimension size must be >= 1, got {self.size}") if self.extent < 0: raise ValueError(f"FixedDimension extent must be >= 0, got {self.extent}") - if self.size == 0: - n = 0 - else: - n = ceildiv(self.extent, self.size) + n = ceildiv(self.extent, self.size) object.__setattr__(self, "nchunks", n) object.__setattr__(self, "ngridcells", n) @@ -71,8 +68,6 @@ def index_to_chunk(self, idx: int) -> int: raise IndexError(f"Negative index {idx} is not allowed") if idx >= self.extent: raise IndexError(f"Index {idx} is out of bounds for extent {self.extent}") - if self.size == 0: - return 0 return idx // self.size def chunk_offset(self, chunk_ix: int) -> int: @@ -97,8 +92,6 @@ def data_size(self, chunk_ix: int) -> int: Does not validate *chunk_ix* — callers must ensure it is in ``[0, nchunks)``. Use ``ChunkGrid.__getitem__`` for safe access. """ - if self.size == 0: - return 0 return max(0, min(self.size, self.extent - chunk_ix * self.size)) @property @@ -112,8 +105,6 @@ def _unique_edge_lengths(self) -> Iterable[int]: return (self.size,) def indices_to_chunks(self, indices: npt.NDArray[np.intp]) -> npt.NDArray[np.intp]: - if self.size == 0: - return np.zeros_like(indices) return indices // self.size def with_extent(self, new_extent: int) -> FixedDimension: @@ -660,6 +651,17 @@ class ChunkLayout(NamedTuple): inner: ChunkLayout | None = None +def full_span_chunk_size(span: int, unit: int = 1) -> int: + """The edge length of one chunk spanning an axis of length `span`. + + This is the smallest positive multiple of `unit` that covers `span`, so a + zero-length axis gets a chunk of size `unit` and zero chunks. `unit` is the + size the chunk must be a multiple of: the inner chunk size for a shard, 1 + otherwise. + """ + return unit * max(1, ceildiv(span, unit)) + + def _guess_regular_chunks( shape: tuple[int, ...] | int, typesize: int, @@ -696,12 +698,12 @@ def _guess_regular_chunks( if isinstance(shape, int): shape = (shape,) + # Start from one chunk spanning each axis, then halve axes until the chunk is small enough. + chunks = np.array([full_span_chunk_size(s) for s in shape], dtype="=f8") if typesize == 0: - return shape + return tuple(int(x) for x in chunks) ndims = len(shape) - # require chunks to have non-zero length for all dimensions - chunks = np.maximum(np.array(shape, dtype="=f8"), 1) # Determine the optimal chunk size in bytes using a PyTables expression. # This is kept as a float. @@ -736,7 +738,7 @@ def _guess_regular_chunks( return tuple(int(x) for x in chunks) -def normalize_chunks_1d(chunks: int | Iterable[object], span: int) -> DimensionGrid: +def normalize_chunks_1d(chunks: int | Iterable[object], span: int, unit: int = 1) -> DimensionGrid: """ Normalize a one-dimensional chunk specification into a dimension grid: `FixedDimension` for scalar chunk sizes, `VaryingDimension` for explicit @@ -744,13 +746,15 @@ def normalize_chunks_1d(chunks: int | Iterable[object], span: int) -> DimensionG the span, and the uniform form is O(1) in the number of chunks — a dimension with `2**62` chunks must not materialize one entry per chunk. - `-1` means "one chunk covering the entire span." + `-1` means "one chunk covering the entire span", sized by + `full_span_chunk_size(span, unit)`. Explicit chunk size lists must sum to the span exactly and always produce `VaryingDimension`, even when the sizes happen to be uniform: the input syntax declares the grid kind, so a per-chunk list is preserved as a rectilinear dimension rather than silently collapsed to a regular one, which would change how the dimension grows on resize. For scalar sizes - the last chunk may overhang the span. + the last chunk may overhang the span. On a zero-length span any non-empty + list of positive sizes is kept: the chunks the axis grows into. """ # `numbers.Integral` rather than `int` so that numpy integer scalars (which are not # `int` subclasses) take the uniform-chunk path instead of being treated as a sequence. @@ -761,9 +765,7 @@ def normalize_chunks_1d(chunks: int | Iterable[object], span: int) -> DimensionG if chunk_size < -1 or chunk_size == 0: raise ValueError(f"Chunk size must be positive or -1, got {chunk_size}") if chunk_size == -1: - # A zero-length span still gets chunk size 1 (chunk sizes must be positive), - # matching the auto-chunking clamp in _guess_regular_chunks. - return FixedDimension(size=max(span, 1), extent=span) + return FixedDimension(size=full_span_chunk_size(span, unit), extent=span) return FixedDimension(size=chunk_size, extent=span) else: try: @@ -788,7 +790,7 @@ def normalize_chunks_1d(chunks: int | Iterable[object], span: int) -> DimensionG ints: list[int] = [int(c) for c in chunk_list] # type: ignore[call-overload] if any(c <= 0 for c in ints): raise ValueError(f"All chunk sizes must be positive, got {ints}") - if sum(ints) != span: + if span > 0 and sum(ints) != span: raise ValueError(f"Chunk sizes {ints} do not sum to span {span}") return VaryingDimension(ints, extent=span) @@ -796,6 +798,7 @@ def normalize_chunks_1d(chunks: int | Iterable[object], span: int) -> DimensionG def normalize_chunks_nd( chunks: Any, shape: tuple[int, ...], + unit: tuple[int, ...] | None = None, ) -> ChunkGrid: """ Normalize a chunk specification into a `ChunkGrid`. @@ -815,6 +818,10 @@ def normalize_chunks_nd( `ChunkGrid` directly. `chunks=None` and `chunks=True` are rejected here — the caller is responsible for choosing between explicit sizes and auto-chunking. + + `unit` gives, per axis, the size a chunk must be a multiple of (the inner + chunk shape, when normalizing a shard shape); it only affects the chunks + that `-1` and `False` derive from the span. """ from zarr.core.metadata.v3 import RectilinearChunkGridMetadata, RegularChunkGridMetadata @@ -828,8 +835,7 @@ def normalize_chunks_nd( f'{chunks!r} is not a valid chunk input. Use chunks=None or chunks="auto" from the top-level API for auto-chunking, or pass an int / tuple of ints.' ) - # handle no chunking: one chunk covering every axis. Routed through the -1 sentinel so - # the zero-length-axis clamp lives in one place (normalize_chunks_1d). + # handle no chunking: one chunk covering every axis. if chunks is False: chunks = -1 @@ -843,8 +849,13 @@ def normalize_chunks_nd( f"chunks has {len(chunks)} dimensions but shape has {len(shape)} dimensions" ) + if unit is None: + unit = (1,) * len(shape) return ChunkGrid( - dimensions=tuple(normalize_chunks_1d(c, span=s) for c, s in zip(chunks, shape, strict=True)) + dimensions=tuple( + normalize_chunks_1d(c, span=s, unit=u) + for c, s, u in zip(chunks, shape, unit, strict=True) + ) ) @@ -884,8 +895,8 @@ def _guess_num_chunks_per_axis_shard( In other words the shard would be a (2,2,2) grid of (2,2,2) chunks i.e., prod(chunk_shape) * (returned_val ** len(chunk_shape)) * item_size = 256 bytes. - Degenerate chunk shapes — a 0-dimensional shape, or one containing a zero-length - axis — return 1, as the search loop's stopping conditions can never be met. + Degenerate inputs — a 0-dimensional chunk shape, or a zero-byte chunk — return 1, + as the search loop's stopping conditions can never be met. Parameters ---------- @@ -906,8 +917,8 @@ def _guess_num_chunks_per_axis_shard( if max_bytes < bytes_per_chunk: return 1 num_axes = len(chunk_shape) - # For a 0-dimensional chunk shape or one with a zero-length axis, both loop - # conditions below are constant, so the loop would never terminate. + # For a 0-dimensional chunk shape or a zero-byte chunk, both loop conditions + # below are constant, so the loop would never terminate. if num_axes == 0 or bytes_per_chunk == 0: return 1 chunks_per_shard = 1 @@ -990,5 +1001,5 @@ def resolve_outer_and_inner_chunks( else: shard_flat = cast("tuple[int, ...]", shard_shape) - outer = normalize_chunks_nd(shard_flat, array_shape) + outer = normalize_chunks_nd(shard_flat, array_shape, unit=chunk_shape_flat) return ChunkLayout(outer_chunks=outer, inner=ChunkLayout(outer_chunks=chunks)) diff --git a/src/zarr/core/common.py b/src/zarr/core/common.py index ba01f6c19f..7ec4e2d04b 100644 --- a/src/zarr/core/common.py +++ b/src/zarr/core/common.py @@ -276,8 +276,47 @@ def _default_zarr_format() -> ZarrFormat: return cast("ZarrFormat", int(zarr_config.get("default_zarr_format", 3))) -def expand_rle(data: Sequence[int | list[int]]) -> list[int]: - """Expand a mixed array of bare integers and RLE pairs. +def _subject(name: str, axis: int | None) -> str: + """`name` as the subject of an error message, prefixed by the dimension `axis`.""" + return name[0].upper() + name[1:] if axis is None else f"Dimension {axis}: {name}" + + +def _parse_positive_int(value: object, name: str, axis: int | None) -> int: + """`value` as an `int` of at least 1. A `bool` is read as the `int` it equals; any + other type, a NumPy integer or a float (even an integral one: stored documents with + integral floats are read by `zarr.core.metadata.upgrades`), is rejected.""" + subject = _subject(name, axis) + if not isinstance(value, int): + raise TypeError(f"{subject} must be an int, got {value!r}") + if value < 1: + raise ValueError(f"{subject} must be >= 1, got {value!r}") + return int(value) + + +def parse_chunk_edge(size: object, axis: int | None = None) -> int: + """Check that `size` is a chunk edge length: an `int` of at least 1 (a `bool` is + read as the `int` it equals). + + This is the one rule for chunk edge lengths in metadata: bare chunk sizes, explicit + edges and run-length encoded sizes. `axis`, when given, is named in the error. + """ + return _parse_positive_int(size, "chunk edge length", axis) + + +def parse_chunk_shape(data: object) -> tuple[int, ...]: + """Check a regular chunk shape: an iterable, other than a string or a mapping, of one + chunk edge length per axis (see `parse_chunk_edge`).""" + match data: + case str() | Mapping(): + pass + case Iterable(): + return tuple(parse_chunk_edge(size, axis) for axis, size in enumerate(data)) + raise TypeError(f"A chunk shape must be an iterable of chunk edge lengths, got {data!r}") + + +def expand_rle(data: Sequence[object], axis: int | None = None) -> list[int]: + """Expand a mixed array of bare integers and RLE pairs, the edges of dimension + `axis` (named in errors, when given). Per the rectilinear chunk grid spec, each element can be: - a bare integer (an explicit edge length) @@ -285,20 +324,15 @@ def expand_rle(data: Sequence[int | list[int]]) -> list[int]: """ result: list[int] = [] for item in data: - if isinstance(item, (int, float)) and not isinstance(item, bool): - val = int(item) - if val < 1: - raise ValueError(f"Chunk edge length must be >= 1, got {val}") - result.append(val) - elif isinstance(item, list) and len(item) == 2: - size, count = int(item[0]), int(item[1]) - if size < 1: - raise ValueError(f"Chunk edge length must be >= 1, got {size}") - if count < 1: - raise ValueError(f"RLE repeat count must be >= 1, got {count}") - result.extend([size] * count) + if isinstance(item, list): + if len(item) != 2: + subject = _subject("RLE entries", axis) + raise ValueError(f"{subject} must be an integer or [size, count], got {item}") + size, count = item + repeat = _parse_positive_int(count, "RLE repeat count", axis) + result.extend([parse_chunk_edge(size, axis)] * repeat) else: - raise ValueError(f"RLE entries must be an integer or [size, count], got {item}") + result.append(parse_chunk_edge(item, axis)) return result diff --git a/src/zarr/core/group.py b/src/zarr/core/group.py index d734e6b7cd..d76106b1cd 100644 --- a/src/zarr/core/group.py +++ b/src/zarr/core/group.py @@ -147,11 +147,16 @@ class ConsolidatedMetadata: must_understand: Literal[False] = False def to_dict(self) -> dict[str, JSON]: + """The consolidated metadata document. An array read from a stored document that + had to be upgraded is written as that document was stored, so every reader of + the consolidated metadata reads it as upgraded again (see `mark_upgraded`).""" return { "kind": self.kind, "must_understand": self.must_understand, "metadata": { k: v.to_dict() + if isinstance(v, GroupMetadata) or v._stored_document is None + else dict(v._stored_document) for k, v in sorted( self.flattened_metadata.items(), key=lambda item: ( @@ -163,7 +168,9 @@ def to_dict(self) -> dict[str, JSON]: } @classmethod - def from_dict(cls, data: dict[str, JSON]) -> ConsolidatedMetadata: + def from_dict(cls, data: dict[str, JSON], *, path: str | None = None) -> ConsolidatedMetadata: + """Read consolidated metadata, naming each member by its path under the group at + `path` (or relative to that group, when `path` is not given) in warnings.""" data = dict(data) kind = data.get("kind") @@ -177,6 +184,7 @@ def from_dict(cls, data: dict[str, JSON]) -> ConsolidatedMetadata: metadata: dict[str, ArrayV2Metadata | ArrayV3Metadata | GroupMetadata] = {} if raw_metadata: for k, v in raw_metadata.items(): + member = k if path is None else _join_paths([path, k]) if not isinstance(v, dict): raise TypeError( f"Invalid value for metadata items. key='{k}', type='{type(v).__name__}'" @@ -188,16 +196,16 @@ def from_dict(cls, data: dict[str, JSON]) -> ConsolidatedMetadata: if zarr_format == 3: node_type = parse_node_type(v.get("node_type", None)) if node_type == "group": - metadata[k] = GroupMetadata.from_dict(v) + metadata[k] = GroupMetadata.from_dict(v, path=member) elif node_type == "array": - metadata[k] = ArrayV3Metadata.from_dict(v) + metadata[k] = ArrayV3Metadata.from_dict(v, path=member) else: assert_never(node_type) elif zarr_format == 2: if "shape" in v: - metadata[k] = ArrayV2Metadata.from_dict(v) + metadata[k] = ArrayV2Metadata.from_dict(v, path=member) else: - metadata[k] = GroupMetadata.from_dict(v) + metadata[k] = GroupMetadata.from_dict(v, path=member) else: assert_never(zarr_format) @@ -415,7 +423,9 @@ def __init__( object.__setattr__(self, "consolidated_metadata", consolidated_metadata) @classmethod - def from_dict(cls, data: dict[str, Any]) -> GroupMetadata: + def from_dict(cls, data: dict[str, Any], *, path: str | None = None) -> GroupMetadata: + """Read a stored group document; `path` names the group in warnings about the + consolidated metadata it holds.""" data = dict(data) node_type = data.pop("node_type", None) if node_type not in ("group", None): @@ -424,7 +434,9 @@ def from_dict(cls, data: dict[str, Any]) -> GroupMetadata: ) consolidated_metadata = data.pop("consolidated_metadata", None) if consolidated_metadata: - data["consolidated_metadata"] = ConsolidatedMetadata.from_dict(consolidated_metadata) + data["consolidated_metadata"] = ConsolidatedMetadata.from_dict( + consolidated_metadata, path=path + ) zarr_format = data.get("zarr_format") if zarr_format == 2 or zarr_format is None: @@ -695,7 +707,7 @@ def from_dict( msg = f"Node type in metadata ({node_type}) is not 'group'" raise GroupNotFoundError(msg) return cls( - metadata=GroupMetadata.from_dict(data), + metadata=GroupMetadata.from_dict(data, path=str(store_path)), store_path=store_path, ) @@ -3078,6 +3090,7 @@ async def create_hierarchy( # ensure that all nodes have the same zarr_format, and add implicit groups as needed nodes_parsed = _parse_hierarchy_dict(data=nodes_normed_keys) redundant_implicit_groups = [] + to_delete_keys: list[str] = [] # empty hierarchies should be a no-op if len(nodes_parsed) > 0: @@ -3110,13 +3123,7 @@ async def create_hierarchy( if overwrite: # we will remove any nodes that collide with arrays and non-implicit groups defined in # nodes - - # track the keys of nodes we need to delete - to_delete_keys = [] - to_delete_keys.extend( - [k for k, v in nodes_parsed.items() if k not in implicit_group_keys] - ) - await asyncio.gather(*(store.delete_dir(key) for key in to_delete_keys)) + to_delete_keys = [k for k in nodes_parsed if k not in implicit_group_keys] else: # This type is long. coros: ( @@ -3176,7 +3183,11 @@ async def create_hierarchy( else: nodes_explicit[k] = v - async for key, node in create_nodes(store=store, nodes=nodes_explicit): + # Build every node before deleting or storing anything: metadata that no array or + # group can be built from then fails with the store untouched. + built = _build_nodes(store, nodes_explicit) + await asyncio.gather(*(store.delete_dir(key) for key in to_delete_keys)) + async for key, node in _store_nodes(store, nodes_explicit, built): yield key, node @@ -3206,7 +3217,26 @@ async def create_nodes( AsyncGroup | AsyncArray The created nodes in the order they are created. """ + async for key, node in _store_nodes(store, nodes, _build_nodes(store, nodes)): + yield key, node + +def _build_nodes( + store: Store, nodes: Mapping[str, GroupMetadata | ArrayV2Metadata | ArrayV3Metadata] +) -> dict[str, AsyncGroup | AnyAsyncArray]: + """The array or group each of `nodes` describes, at its path in `store`.""" + return { + path: _build_node(store=store, path=path, metadata=meta) for path, meta in nodes.items() + } + + +async def _store_nodes( + store: Store, + nodes: Mapping[str, GroupMetadata | ArrayV2Metadata | ArrayV3Metadata], + built: Mapping[str, AsyncGroup | AnyAsyncArray], +) -> AsyncIterator[tuple[str, AsyncGroup | AnyAsyncArray]]: + """Store the metadata of `nodes` and yield the nodes `_build_nodes` built from them + (see `create_nodes`).""" # Note: the only way to alter this value is via the config. If that's undesirable for some reason, # then we should consider adding a keyword argument to this function semaphore = asyncio.Semaphore(config.get("async.concurrency")) @@ -3234,7 +3264,7 @@ async def create_nodes( node_name = created_key[: created_key.rfind("/")] meta_out = nodes[node_name] if meta_out.zarr_format == 3: - yield node_name, _build_node(store=store, path=node_name, metadata=meta_out) + yield node_name, built[node_name] else: # For zarr v2 # we only want to yield when both the metadata and attributes are created @@ -3249,7 +3279,7 @@ async def create_nodes( meta_done = _join_paths([node_name, ZARRAY_JSON]) in created_object_keys if meta_done and attrs_done: - yield node_name, _build_node(store=store, path=node_name, metadata=meta_out) + yield node_name, built[node_name] continue @@ -3504,7 +3534,9 @@ async def _read_metadata_v3(store: Store, path: str) -> ArrayV3Metadata | GroupM ) if zarr_json_bytes is None: raise FileNotFoundError(path) - return _build_metadata_v3(buffer_to_json_object(zarr_json_bytes)) + return _build_metadata_v3( + buffer_to_json_object(zarr_json_bytes), path=str(StorePath(store, path)) + ) async def _read_metadata_v2(store: Store, path: str) -> ArrayV2Metadata | GroupMetadata: @@ -3539,7 +3571,7 @@ async def _read_metadata_v2(store: Store, path: str) -> ArrayV2Metadata | GroupM else: zmeta = buffer_to_json_object(zgroup_bytes) - return _build_metadata_v2(zmeta, zattrs) + return _build_metadata_v2(zmeta, zattrs, path=str(StorePath(store, path))) async def _read_group_metadata_v2(store: Store, path: str) -> GroupMetadata: @@ -3570,7 +3602,9 @@ async def _read_group_metadata( return await _read_group_metadata_v3(store=store, path=path) -def _build_metadata_v3(zarr_json: dict[str, JSON]) -> ArrayV3Metadata | GroupMetadata: +def _build_metadata_v3( + zarr_json: dict[str, JSON], *, path: str | None = None +) -> ArrayV3Metadata | GroupMetadata: """ Convert a dict representation of Zarr V3 metadata into the corresponding metadata class. """ @@ -3579,9 +3613,9 @@ def _build_metadata_v3(zarr_json: dict[str, JSON]) -> ArrayV3Metadata | GroupMet raise MetadataValidationError(msg) match zarr_json: case {"node_type": "array"}: - return ArrayV3Metadata.from_dict(zarr_json) + return ArrayV3Metadata.from_dict(zarr_json, path=path) case {"node_type": "group"}: - return GroupMetadata.from_dict(zarr_json) + return GroupMetadata.from_dict(zarr_json, path=path) case _: # pragma: no cover raise ValueError( "invalid value for `node_type` key in metadata document" @@ -3589,16 +3623,16 @@ def _build_metadata_v3(zarr_json: dict[str, JSON]) -> ArrayV3Metadata | GroupMet def _build_metadata_v2( - zarr_json: dict[str, JSON], attrs_json: dict[str, JSON] + zarr_json: dict[str, JSON], attrs_json: dict[str, JSON], *, path: str | None = None ) -> ArrayV2Metadata | GroupMetadata: """ Convert a dict representation of Zarr V2 metadata into the corresponding metadata class. """ match zarr_json: case {"shape": _}: - return ArrayV2Metadata.from_dict(zarr_json | {"attributes": attrs_json}) + return ArrayV2Metadata.from_dict(zarr_json | {"attributes": attrs_json}, path=path) case _: # pragma: no cover - return GroupMetadata.from_dict(zarr_json | {"attributes": attrs_json}) + return GroupMetadata.from_dict(zarr_json | {"attributes": attrs_json}, path=path) @overload diff --git a/src/zarr/core/metadata/io.py b/src/zarr/core/metadata/io.py index 7b63f5493b..c5c1c24da3 100644 --- a/src/zarr/core/metadata/io.py +++ b/src/zarr/core/metadata/io.py @@ -1,19 +1,140 @@ from __future__ import annotations import asyncio -from typing import TYPE_CHECKING +from enum import Enum +from itertools import zip_longest +from typing import TYPE_CHECKING, Final, NamedTuple from zarr.abc.store import set_or_delete +from zarr.core._json import buffer_to_json_object, json_equal from zarr.core.buffer.core import default_buffer_prototype -from zarr.errors import ContainsArrayError +from zarr.core.buffer.cpu import buffer_prototype as cpu_buffer_prototype +from zarr.core.common import ZARR_JSON, ZARRAY_JSON, ZATTRS_JSON +from zarr.core.metadata.upgrades import mark_upgraded, upgrade_array_document +from zarr.errors import ArrayNotFoundError, ContainsArrayError from zarr.storage._common import StorePath, ensure_no_existing_node if TYPE_CHECKING: - from zarr.core.common import ZarrFormat + from collections.abc import Iterable, Iterator, Mapping + + from zarr.core.buffer import Buffer + from zarr.core.common import JSON, ZarrFormat from zarr.core.group import GroupMetadata from zarr.core.metadata import ArrayMetadata +class _Absent(Enum): + ABSENT = "absent" + + +ABSENT: Final = _Absent.ABSENT +"""Where a `DocumentChange` has no value: the document, member or element is not there.""" + + +class DocumentChange(NamedTuple): + """A JSON value that differs between a node's stored metadata documents and the + documents its metadata would store.""" + + path: tuple[str | int, ...] + """Where the value is: the document's key (e.g. `.zarray`), then the object members + and array indices within it.""" + stored: JSON | _Absent + new: JSON | _Absent + + +def diff_documents( + stored: Mapping[str, JSON], new: Mapping[str, JSON] +) -> tuple[DocumentChange, ...]: + """What differs between a node's stored metadata documents and the documents its + metadata would store, each keyed by its store key (a document missing from `stored` + is not stored). Empty if they are identical. + + Objects and arrays are compared member by member, so a change is the smallest value + that differs; any other two values are identical only if their JSON encodings are, so + `true` differs from `1` and `1.0` from `1`. + """ + return tuple(_diff((), stored, new)) + + +def _diff( + path: tuple[str | int, ...], stored: JSON | _Absent, new: JSON | _Absent +) -> Iterator[DocumentChange]: + match stored, new: + case dict(), dict(): + for key in dict.fromkeys([*stored, *new]): + yield from _diff((*path, key), stored.get(key, ABSENT), new.get(key, ABSENT)) + case list(), list(): + for index, pair in enumerate(zip_longest(stored, new, fillvalue=ABSENT)): + yield from _diff((*path, index), *pair) + case _ if stored is ABSENT or new is ABSENT or not json_equal(stored, new): + yield DocumentChange(path, stored, new) + + +async def store_documents(store_path: StorePath, documents: Mapping[str, Buffer]) -> None: + """Store metadata documents encoded by `to_buffer_dict` under `store_path`.""" + await asyncio.gather( + *(set_or_delete(store_path / key, value) for key, value in documents.items()) + ) + + +async def read_documents(store_path: StorePath, keys: Iterable[str]) -> dict[str, Buffer]: + """The documents stored under `store_path` at `keys`, read concurrently, by key; a + key with no document is left out.""" + keys = tuple(keys) + buffers = await asyncio.gather( + *((store_path / key).get(prototype=cpu_buffer_prototype) for key in keys) + ) + return {key: buf for key, buf in zip(keys, buffers, strict=True) if buf is not None} + + +async def upsert_metadata( + store_path: StorePath, metadata: ArrayMetadata | GroupMetadata, stored: Mapping[str, Buffer] +) -> tuple[DocumentChange, ...]: + """Store the documents of `metadata` under `store_path` that differ from `stored`, + the documents `read_documents` read there, and return how they differed (see + `diff_documents`): empty if nothing was stored. + + The documents are encoded before any is stored, so metadata that cannot be stored + fails with the store untouched. + """ + documents = metadata.to_buffer_dict(default_buffer_prototype()) + changes = diff_documents( + {key: buffer_to_json_object(buf) for key, buf in stored.items() if key in documents}, + {key: buffer_to_json_object(buf) for key, buf in documents.items()}, + ) + changed = {change.path[0] for change in changes} + await store_documents(store_path, {k: v for k, v in documents.items() if k in changed}) + return changes + + +ARRAY_DOCUMENTS: Final[Mapping[ZarrFormat, tuple[str, ...]]] = { + 2: (ZARRAY_JSON, ZATTRS_JSON), + 3: (ZARR_JSON,), +} +"""The store keys of the metadata documents of an array of each Zarr format.""" + + +def parse_stored_array(documents: Mapping[str, Buffer], zarr_format: ZarrFormat) -> ArrayMetadata: + """The metadata of an array from its documents (by store key, see `ARRAY_DOCUMENTS`), + read with the upgrades but without their warnings (whoever asks has warned, or reads + metadata built in code), and marked (see `mark_upgraded`) if they had to be + upgraded. Raises `ArrayNotFoundError` if there is no array document among them.""" + from zarr.core.array import ( + _array_metadata_dict_v2, + _array_metadata_dict_v3, + parse_array_metadata, + ) + + if zarr_format == 2 and ZARRAY_JSON in documents: + stored = _array_metadata_dict_v2(documents[ZARRAY_JSON], documents.get(ZATTRS_JSON)) + elif zarr_format == 3 and ZARR_JSON in documents: + stored = _array_metadata_dict_v3(documents[ZARR_JSON]) + else: + raise ArrayNotFoundError(f"No Zarr format {zarr_format} array metadata document.") + upgraded, readings = upgrade_array_document(stored, zarr_format) + return mark_upgraded(parse_array_metadata(dict(upgraded)), stored, [None] * len(readings), None) + + def _build_parents(store_path: StorePath, zarr_format: ZarrFormat) -> dict[str, GroupMetadata]: from zarr.core.group import GroupMetadata diff --git a/src/zarr/core/metadata/upgrades.py b/src/zarr/core/metadata/upgrades.py new file mode 100644 index 0000000000..728a0e53f5 --- /dev/null +++ b/src/zarr/core/metadata/upgrades.py @@ -0,0 +1,272 @@ +"""Upgrades that read invalid stored array metadata documents written by older software. + +This is the only place invalid metadata is read leniently. An upgrade maps a stored +array metadata document (parsed JSON) to a valid one. `ArrayV2Metadata.from_dict` and +`ArrayV3Metadata.from_dict` apply the upgrades for their Zarr format, so every path +that parses a stored document, including consolidated metadata, goes through them. + +A reading warns only where the user must act on it; a reading that gives what zarr +read from the same document before these upgrades existed is silent, so a document +that opened without a warning still does. The warnings are given once per document, +after the upgraded document has passed the metadata constructor, so an invalid +document raises its own error, not a warning about how it was read. Silent or not, +metadata read from an upgraded document is marked (see `mark_upgraded`), so the array +stores the upgrade before it writes chunks under it. + +To read another kind of invalid document, add an upgrade to `ARRAY_UPGRADES`. +""" + +from __future__ import annotations + +import copy +import json +import warnings +from collections.abc import Callable, Iterable, Mapping, Sequence +from itertools import chain, repeat +from typing import TYPE_CHECKING, Final, TypeGuard, cast + +from zarr.core._json import json_equal +from zarr.core.chunk_grids import full_span_chunk_size +from zarr.errors import ZarrUserWarning + +if TYPE_CHECKING: + from zarr.core.common import JSON, ZarrFormat + +type ArrayDocument = Mapping[str, JSON] +type Upgrade = Callable[[ArrayDocument], tuple[ArrayDocument, str | None] | None] +"""Returns `None` if the document needs no upgrade, else the upgraded document and, +if the user must act on how it was read, a warning saying so (else `None`).""" + +RESAVE_HINT: Final = ( + "To store valid metadata, open the array writable and call `array.update_attributes({})`; " + "if a group holds consolidated metadata for the array, then also call " + "`zarr.consolidate_metadata` on that group." +) + + +def mark_upgraded[M]( + metadata: M, stored: ArrayDocument, readings: Sequence[str | None], path: str | None +) -> M: + """Record that `metadata` was read from the document `stored`, which needed the + upgrades whose `readings` `upgrade_array_document` returned, if any: keep a copy of + `stored` as `_stored_document` on `metadata`, so the array stores the upgrade before + it writes chunks under it and consolidated metadata stores it as it was stored, and + warn once with the readings that are warnings, naming the array at `path` when the + caller knows it.""" + if readings: + object.__setattr__(metadata, "_stored_document", copy.deepcopy(stored)) + if messages := [reading for reading in readings if reading is not None]: + subject = "" if path is None else f"Array {path!r}: " + # The synchronous API parses metadata on zarr's IO thread, whose stack holds no + # user code, so the warning points at the `from_dict` that read the document. + warnings.warn(f"{subject}{' '.join(messages)} {RESAVE_HINT}", ZarrUserWarning, stacklevel=2) + return metadata + + +def _is_int_list(value: object) -> TypeGuard[list[int]]: + """Whether `value` is a JSON array of integers (JSON `false` and `true` count).""" + return isinstance(value, list) and all(isinstance(v, int) for v in value) + + +def _read_chunk_size(size: JSON, span: int | None, unit: int) -> tuple[int, str | None] | None: + """Read one entry of a stored regular chunk shape as a chunk edge length. + + Returns the edge length and, where the user must act on how it was read, how it was + read; `None` if the entry cannot be read, which leaves it for the metadata + constructors to check. A JSON int >= 1 is kept, JSON `true` is read as 1, and 0 or + JSON `false` is read as one chunk spanning the axis of length `span`, a multiple of + `unit` (the inner chunk size of a shard): on an axis of positive length no chunk + can have been stored under it, so the array holds only its fill value. `span` is + `None` where no stored 0 is known, as in the inner chunk shape of a sharding codec: + 0 is then left as stored. + """ + match size: + case True: + return 1, None + case int() if size >= 1: + return size, None + case int() if size == 0 and span is not None: + edge = full_span_chunk_size(span, unit) + if span == 0: + return edge, None + return edge, ( + f"one chunk spanning the dimension ({edge}), and as no chunk can be stored " + "under a chunk size of 0, the array holds only its fill value" + ) + return None + + +def _read_chunk_shape( + stored: JSON, spans: Sequence[int | None], units: Iterable[int] = () +) -> tuple[list[int], str | None] | None: + """Read a stored regular chunk shape, entry by entry (see `_read_chunk_size`), for + axes of lengths `spans` whose chunks are multiples of `units` (1 where not given). + + Returns the chunk shape and, where the user must act on how an entry was read, a + sentence saying how the chunk shape was read; `None` if it cannot be read. + """ + if not (isinstance(stored, list) and len(stored) == len(spans)): + return None + edges: list[int] = [] + readings: list[str] = [] + axes = zip(stored, spans, chain(units, repeat(1)), strict=False) + for axis, (size, span, unit) in enumerate(axes): + read = _read_chunk_size(size, span, unit) + if read is None: + return None + edge, how = read + edges.append(edge) + if how is not None: + readings.append(f"{json.dumps(size)} in dimension {axis} as {how}") + if not readings: + return edges, None + return edges, ( + f"The stored chunk shape {json.dumps(stored)} is invalid: chunk sizes must be " + f"integers of at least 1. It is read as {edges}, reading {'; '.join(readings)}." + ) + + +def _invalid_chunk_sizes_v2(doc: ArrayDocument) -> tuple[ArrayDocument, str | None] | None: + shape = doc.get("shape") + if not _is_int_list(shape): + return None + stored = doc.get("chunks") + match _read_chunk_shape(stored, shape): + case chunks, reading if not json_equal(chunks, stored): + return {**doc, "chunks": chunks}, reading + return None + + +def _read_codec(codec: JSON) -> JSON: + """Read a stored codec: the inner chunk shape of a sharding codec is read as a chunk + shape with no known axis lengths (see `_read_chunk_shape`), and so are those of the + sharding codecs nested in its codecs.""" + match codec: + case {"name": "sharding_indexed", "configuration": Mapping() as configuration}: + upgraded = dict(configuration) + stored = configuration.get("chunk_shape") + if isinstance(stored, list) and ( + read := _read_chunk_shape(stored, [None] * len(stored)) + ): + upgraded["chunk_shape"] = read[0] + if isinstance(codecs := configuration.get("codecs"), list): + upgraded["codecs"] = [_read_codec(inner) for inner in codecs] + # The mapping pattern does not narrow `codec` for mypy. + return {**cast("Mapping[str, JSON]", codec), "configuration": upgraded} + return codec + + +def _invalid_inner_chunk_sizes_v3(doc: ArrayDocument) -> tuple[ArrayDocument, str | None] | None: + stored = doc.get("codecs") + if not isinstance(stored, list): + return None + codecs = [_read_codec(codec) for codec in stored] + if json_equal(codecs, stored): + return None + return {**doc, "codecs": codecs}, None + + +def _inner_chunk_shape(doc: ArrayDocument) -> list[int] | None: + """The inner chunk shape of the sharding codec of a Zarr format 3 array document: + `[]` if it has none, `None` if its inner chunk sizes are not all integers of at + least 1 (the unit of its outer chunk shape is then unknown).""" + match doc.get("codecs"): + case list() as codecs: + for codec in codecs: + match codec: + case {"name": "sharding_indexed", "configuration": {"chunk_shape": inner}}: + return inner if _is_int_list(inner) and min(inner, default=1) >= 1 else None + return [] + + +def _invalid_chunk_sizes_v3(doc: ArrayDocument) -> tuple[ArrayDocument, str | None] | None: + grid = doc.get("chunk_grid") + shape = doc.get("shape") + if not (isinstance(grid, Mapping) and grid.get("name") == "regular" and _is_int_list(shape)): + return None + if (units := _inner_chunk_shape(doc)) is None: + return None + configuration = grid.get("configuration") + if not isinstance(configuration, Mapping): + return None + stored = configuration.get("chunk_shape") + match _read_chunk_shape(stored, shape, units): + case chunk_shape, reading if not json_equal(chunk_shape, stored): + upgraded = {**configuration, "chunk_shape": chunk_shape} + return {**doc, "chunk_grid": {**grid, "configuration": upgraded}}, reading + return None + + +def _read_edge_length(edge: JSON) -> JSON: + """Read one stored chunk edge length of a rectilinear chunk grid: JSON `true` is + read as 1 and an integral JSON float of at least 1 (`4.0`) as the `int` it equals. + Anything else is kept, for the metadata constructors to check.""" + match edge: + case True: + return 1 + case float() if edge.is_integer() and edge >= 1: + return int(edge) + return edge + + +def _read_rectilinear_axis(spec: JSON) -> JSON: + """Read the stored chunk edge lengths of one axis of a rectilinear chunk grid (see + `_read_edge_length`): its explicit edges and the sizes of its run-length encoded + `[size, count]` pairs. A bare chunk size and a repeat count are kept as stored: no + writer stored them as `true` or as floats.""" + match spec: + case list(): + return [_read_rectilinear_entry(entry) for entry in spec] + return spec + + +def _read_rectilinear_entry(entry: JSON) -> JSON: + match entry: + case [size, count]: + return [_read_edge_length(size), count] + return _read_edge_length(entry) + + +def _invalid_edge_lengths_v3(doc: ArrayDocument) -> tuple[ArrayDocument, str | None] | None: + grid = doc.get("chunk_grid") + if not (isinstance(grid, Mapping) and grid.get("name") == "rectilinear"): + return None + configuration = grid.get("configuration") + if not isinstance(configuration, Mapping): + return None + stored = configuration.get("chunk_shapes") + if not isinstance(stored, list): + return None + read = [_read_rectilinear_axis(axis) for axis in stored] + if json_equal(read, stored): + return None + upgraded = {**configuration, "chunk_shapes": read} + return {**doc, "chunk_grid": {**grid, "configuration": upgraded}}, None + + +ARRAY_UPGRADES: Final[Mapping[ZarrFormat, tuple[Upgrade, ...]]] = { + 2: (_invalid_chunk_sizes_v2,), + # The inner chunk shape is read first: it gives the unit of the outer chunk shape. + # The rectilinear edge lengths are read last, after any upgrade that yields a + # rectilinear chunk grid. + 3: (_invalid_inner_chunk_sizes_v3, _invalid_chunk_sizes_v3, _invalid_edge_lengths_v3), +} +"""The upgrades of an array document of each Zarr format, applied in order.""" + + +def upgrade_array_document( + doc: ArrayDocument, zarr_format: ZarrFormat +) -> tuple[ArrayDocument, list[str | None]]: + """Apply the upgrades for `zarr_format` to a stored array metadata document. + + Returns the upgraded document and the reading of each upgrade that changed it (a + warning, or `None` for a silent one), for `mark_upgraded` once the document has + been validated. + """ + readings: list[str | None] = [] + for upgrade in ARRAY_UPGRADES[zarr_format]: + upgraded = upgrade(doc) + if upgraded is not None: + doc, reading = upgraded + readings.append(reading) + return doc, readings diff --git a/src/zarr/core/metadata/v2.py b/src/zarr/core/metadata/v2.py index 5822a228b9..7f6a33e93b 100644 --- a/src/zarr/core/metadata/v2.py +++ b/src/zarr/core/metadata/v2.py @@ -4,7 +4,7 @@ import warnings from collections.abc import Iterable, Sequence from functools import cached_property -from typing import TYPE_CHECKING, Any, Literal, TypedDict, cast +from typing import TYPE_CHECKING, Any, ClassVar, Literal, TypedDict, cast from zarr.abc.metadata import Metadata from zarr.abc.numcodec import Numcodec, _is_numcodec @@ -25,6 +25,7 @@ TBaseScalar, ZDType, ) + from zarr.core.metadata.upgrades import ArrayDocument from dataclasses import dataclass, field, fields, replace @@ -38,11 +39,13 @@ ZARRAY_JSON, ZATTRS_JSON, MemoryOrder, + ShapeLike, parse_shapelike, ) from zarr.core.config import config, parse_indexing_order from zarr.core.json_parse import parse_field from zarr.core.metadata.common import parse_attributes +from zarr.core.metadata.upgrades import mark_upgraded, upgrade_array_document class ArrayV2MetadataDict(TypedDict): @@ -70,6 +73,9 @@ class ArrayV2Metadata(Metadata): compressor: Numcodec | None attributes: dict[str, JSON] = field(default_factory=dict) zarr_format: Literal[2] = field(init=False, default=2) + _stored_document: ClassVar[ArrayDocument | None] = None + """The stored document `from_dict` read this metadata from, if it had to upgrade it + (set on the instance by `mark_upgraded`): the store may still hold it.""" def __init__( self, @@ -88,7 +94,7 @@ def __init__( Metadata for a Zarr format 2 array. """ shape_parsed = parse_shapelike(shape) - chunks_parsed = parse_shapelike(chunks) + chunks_parsed = parse_chunks(chunks, shape_parsed) compressor_parsed = parse_compressor(compressor) order_parsed = parse_indexing_order(order) dimension_separator_parsed = parse_separator(dimension_separator) @@ -110,9 +116,6 @@ def __init__( object.__setattr__(self, "fill_value", fill_value_parsed) object.__setattr__(self, "attributes", attributes_parsed) - # ensure that the metadata document is consistent - _ = parse_metadata(self) - @property def ndim(self) -> int: return len(self.shape) @@ -149,9 +152,13 @@ def to_buffer_dict(self, prototype: BufferPrototype) -> dict[str, Buffer]: } @classmethod - def from_dict(cls, data: dict[str, Any]) -> ArrayV2Metadata: - # Make a copy to protect the original from modification. - _data = data.copy() + def from_dict(cls, data: dict[str, Any], *, path: str | None = None) -> ArrayV2Metadata: + """Read a stored `.zarray` document (with its attributes). An invalid document + that `zarr.core.metadata.upgrades` can read is read as upgraded; a reading the user + must act on warns, naming the array at `path`.""" + upgraded, readings = upgrade_array_document(data, 2) + # a new dict, because we are modifying it + _data: dict[str, Any] = dict(upgraded) # Check that the zarr_format attribute is correct. _ = parse_zarr_format(_data.pop("zarr_format")) @@ -159,7 +166,7 @@ def from_dict(cls, data: dict[str, Any]) -> ArrayV2Metadata: # which could be in filters or as a compressor. # we will reference a hard-coded collection of object codec ids for this search. - _filters, _compressor = (data.get("filters"), data.get("compressor")) + _filters, _compressor = (_data.get("filters"), _data.get("compressor")) if _filters is not None: _filters = cast("tuple[dict[str, JSON], ...]", _filters) object_codec_id = get_object_codec_id(tuple(_filters) + (_compressor,)) @@ -168,7 +175,7 @@ def from_dict(cls, data: dict[str, Any]) -> ArrayV2Metadata: # we add a layer of indirection here around the dtype attribute of the array metadata # because we also need to know the object codec id, if any, to resolve the data type dtype_spec: DTypeSpec_V2 = { - "name": data["dtype"], + "name": _data["dtype"], "object_codec_id": object_codec_id, } dtype = get_data_type_from_json(dtype_spec, zarr_format=2) @@ -202,7 +209,7 @@ def from_dict(cls, data: dict[str, Any]) -> ArrayV2Metadata: _data = {k: v for k, v in _data.items() if k in expected} - return cls(**_data) + return mark_upgraded(cls(**_data), data, readings, path) def to_dict(self) -> dict[str, JSON]: zarray_dict = super().to_dict() @@ -323,14 +330,16 @@ def parse_compressor(data: object) -> Numcodec | None: raise ValueError(msg) -def parse_metadata(data: ArrayV2Metadata) -> ArrayV2Metadata: - if (l_chunks := len(data.chunks)) != (l_shape := len(data.shape)): - msg = ( +def parse_chunks(chunks: ShapeLike, shape: tuple[int, ...]) -> tuple[int, ...]: + """Check a chunk shape: one non-negative integer per array axis (see + `parse_shapelike`). Stored chunk sizes of 0 are read by `zarr.core.metadata.upgrades`.""" + chunks_parsed = parse_shapelike(chunks) + if len(chunks_parsed) != len(shape): + raise ValueError( f"The `shape` and `chunks` attributes must have the same length. " - f"`chunks` has length {l_chunks}, but `shape` has length {l_shape}." + f"`chunks` has length {len(chunks_parsed)}, but `shape` has length {len(shape)}." ) - raise ValueError(msg) - return data + return chunks_parsed def get_object_codec_id(maybe_object_codecs: Sequence[JSON]) -> str | None: diff --git a/src/zarr/core/metadata/v3.py b/src/zarr/core/metadata/v3.py index 11f3eb593d..b654c36978 100644 --- a/src/zarr/core/metadata/v3.py +++ b/src/zarr/core/metadata/v3.py @@ -3,7 +3,7 @@ import json from collections.abc import Iterable, Mapping, Sequence from dataclasses import dataclass, field, replace -from typing import TYPE_CHECKING, Any, Final, Literal, NotRequired, TypeGuard, cast +from typing import TYPE_CHECKING, Any, ClassVar, Final, Literal, NotRequired, TypeGuard, cast from typing_extensions import TypedDict @@ -26,6 +26,8 @@ NamedRequiredConfig, compress_rle, expand_rle, + parse_chunk_edge, + parse_chunk_shape, parse_named_configuration, parse_shapelike, validate_rectilinear_edges, @@ -36,6 +38,10 @@ from zarr.core.dtype.common import check_dtype_spec_v3 from zarr.core.json_parse import parse_field from zarr.core.metadata.common import parse_attributes +from zarr.core.metadata.upgrades import ( + mark_upgraded, + upgrade_array_document, +) from zarr.errors import MetadataValidationError, NodeTypeValidationError from zarr.registry import get_codec_class @@ -46,6 +52,7 @@ from zarr.core.buffer import Buffer, BufferPrototype from zarr.core.chunk_grids import ChunkGrid from zarr.core.dtype.wrapper import TBaseDType, TBaseScalar + from zarr.core.metadata.upgrades import ArrayDocument def parse_zarr_format(data: object) -> Literal[3]: @@ -212,17 +219,6 @@ class RectilinearChunkGridMetadataConfig(TypedDict): ] -def _parse_chunk_shape(chunk_shape: Iterable[int]) -> tuple[int, ...]: - """Validate and normalize a regular chunk shape. - - Delegates to ``_validate_chunk_shapes`` — a regular chunk shape is just - a sequence of bare ints (one per dimension), each of which must be >= 1. - """ - result = _validate_chunk_shapes(tuple(chunk_shape)) - # Regular grids only have bare ints — cast is safe after validation - return cast(tuple[int, ...], result) - - def _validate_chunk_shapes( chunk_shapes: Sequence[int | Sequence[int]], ) -> tuple[int | tuple[int, ...], ...]: @@ -233,23 +229,14 @@ def _validate_chunk_shapes( """ result: list[int | tuple[int, ...]] = [] for dim_idx, dim_spec in enumerate(chunk_shapes): - if isinstance(dim_spec, int): - if dim_spec < 1: - raise ValueError( - f"Dimension {dim_idx}: integer chunk edge length must be >= 1, got {dim_spec}" - ) - result.append(dim_spec) - else: - edges = tuple(dim_spec) - if not edges: - raise ValueError(f"Dimension {dim_idx} has no chunk edges.") - bad = [i for i, e in enumerate(edges) if e < 1] - if bad: - raise ValueError( - f"Dimension {dim_idx} has invalid edge lengths at indices {bad}: " - f"{[edges[i] for i in bad]}" - ) - result.append(edges) + match dim_spec: + case list() | tuple(): + edges = tuple(parse_chunk_edge(edge, dim_idx) for edge in dim_spec) + if not edges: + raise ValueError(f"Dimension {dim_idx} has no chunk edges.") + result.append(edges) + case _: + result.append(parse_chunk_edge(dim_spec, dim_idx)) return tuple(result) @@ -264,7 +251,7 @@ class RegularChunkGridMetadata(Metadata): chunk_shape: tuple[int, ...] def __post_init__(self) -> None: - chunk_shape_parsed = _parse_chunk_shape(self.chunk_shape) + chunk_shape_parsed = parse_chunk_shape(self.chunk_shape) object.__setattr__(self, "chunk_shape", chunk_shape_parsed) @property @@ -281,7 +268,11 @@ def to_dict(self) -> RegularChunkGridMetadataJSON: # type: ignore[override] def from_dict(cls, data: RegularChunkGridMetadataJSON) -> Self: # type: ignore[override] parse_named_configuration(data, "regular") # validate name configuration = data["configuration"] - return cls(chunk_shape=_parse_chunk_shape(configuration["chunk_shape"])) + return cls(chunk_shape=parse_chunk_shape(configuration["chunk_shape"])) + + +class RectilinearChunksDisabledError(ValueError): + """Rectilinear chunk grids are used while the `array.rectilinear_chunks` flag is off.""" @dataclass(frozen=True, kw_only=True) @@ -304,7 +295,7 @@ class RectilinearChunkGridMetadata(Metadata): def __post_init__(self) -> None: if not config.get("array.rectilinear_chunks"): - raise ValueError( + raise RectilinearChunksDisabledError( "Rectilinear chunk grids are experimental and disabled by default. " "Enable them with: zarr.config.set({'array.rectilinear_chunks': True}) " "or set the environment variable ZARR_ARRAY__RECTILINEAR_CHUNKS=True" @@ -366,18 +357,10 @@ def from_dict(cls, data: RectilinearChunkGridMetadataJSON) -> Self: # type: ign configuration = data["configuration"] validate_rectilinear_kind(configuration.get("kind")) raw_shapes = configuration["chunk_shapes"] - parsed: list[int | tuple[int, ...]] = [] - for dim_spec in raw_shapes: - if isinstance(dim_spec, int): - if dim_spec < 1: - raise ValueError(f"Integer chunk edge length must be >= 1, got {dim_spec}") - parsed.append(dim_spec) - elif isinstance(dim_spec, list): - parsed.append(tuple(expand_rle(dim_spec))) - else: - raise TypeError( - f"Invalid chunk_shapes entry: expected int or list, got {type(dim_spec)}" - ) + parsed = [ + tuple(expand_rle(dim_spec, axis)) if isinstance(dim_spec, list) else dim_spec + for axis, dim_spec in enumerate(raw_shapes) + ] return cls(chunk_shapes=tuple(parsed)) @@ -490,6 +473,9 @@ class ArrayV3Metadata(Metadata): node_type: Literal["array"] = field(default="array", init=False) storage_transformers: tuple[dict[str, JSON], ...] extra_fields: dict[str, AllowedExtraField] + _stored_document: ClassVar[ArrayDocument | None] = None + """The stored document `from_dict` read this metadata from, if it had to upgrade it + (set on the instance by `mark_upgraded`): the store may still hold it.""" def __init__( self, @@ -630,9 +616,13 @@ def to_buffer_dict(self, prototype: BufferPrototype) -> dict[str, Buffer]: return {ZARR_JSON: json_to_buffer(self.to_dict(), prototype=prototype, indent=indent)} @classmethod - def from_dict(cls, data: dict[str, JSON]) -> Self: - # make a copy because we are modifying the dict - _data = data.copy() + def from_dict(cls, data: dict[str, JSON], *, path: str | None = None) -> Self: + """Read a stored `zarr.json` array document. An invalid document that + `zarr.core.metadata.upgrades` can read is read as upgraded; a reading the user + must act on warns, naming the array at `path`.""" + upgraded, readings = upgrade_array_document(data, 3) + # a new dict, because we are modifying it + _data = dict(upgraded) # check that the zarr_format attribute is correct _ = parse_zarr_format(_data.pop("zarr_format")) @@ -672,7 +662,7 @@ def from_dict(cls, data: dict[str, JSON]) -> Self: # TODO: replace this with a real type check! _data_typed = cast(ArrayMetadataJSON_V3, _data) - return cls( + metadata = cls( shape=_data_typed["shape"], chunk_grid=_data_typed["chunk_grid"], # type: ignore[arg-type] chunk_key_encoding=_data_typed["chunk_key_encoding"], # type: ignore[arg-type] @@ -686,6 +676,7 @@ def from_dict(cls, data: dict[str, JSON]) -> Self: extra_fields=allowed_extra_fields, storage_transformers=_data_typed.get("storage_transformers", ()), # type: ignore[arg-type] ) + return mark_upgraded(metadata, data, readings, path) def to_dict(self) -> dict[str, JSON]: out_dict = super().to_dict() diff --git a/tests/test_array.py b/tests/test_array.py index 3ae502a4e8..b153f24c0f 100644 --- a/tests/test_array.py +++ b/tests/test_array.py @@ -732,6 +732,18 @@ def test_resize_1d(store: MemoryStore, zarr_format: ZarrFormat) -> None: assert new_shape == result.shape +@pytest.mark.parametrize("chunks", [(1,), (2,), (4,)]) +def test_resize_sharded_keeps_cells_beyond_shape(chunks: tuple[int, ...]) -> None: + """A shard kept by a shrinking resize keeps its cells beyond the new shape, and a + later write to the shard leaves them alone, so they come back when the array grows.""" + arr = zarr.create_array({}, shape=(4,), chunks=chunks, shards=(4,), dtype="int16", fill_value=0) + arr[:] = [1, 2, 3, 4] + arr.resize((2,)) + arr[:] = [9, 9] + arr.resize((4,)) + np.testing.assert_array_equal(arr[:], [9, 9, 3, 4]) + + @pytest.mark.parametrize("store", ["memory"], indirect=True) def test_resize_2d(store: MemoryStore, zarr_format: ZarrFormat) -> None: z = zarr.create( diff --git a/tests/test_array_stateful.py b/tests/test_array_stateful.py new file mode 100644 index 0000000000..5319dd6f89 --- /dev/null +++ b/tests/test_array_stateful.py @@ -0,0 +1,313 @@ +"""A stateful test of one array's life: create, append, resize, write, reopen. + +The model is a NumPy array of what the store holds, not a second zarr array, so a bug +in zarr's chunk grid logic cannot hide by being made on both sides. Zero-length axes +are drawn on purpose, both at creation and by resizing and appending, and so are +stored chunk sizes of 0. + +The model tracks cells beyond the array's shape too, because chunks do: a shrinking +`resize` deletes exactly the chunks outside the new grid, a chunk it keeps keeps its +cells beyond the new shape (and they come back if the array grows again), and a write +that covers every in-bounds cell of an unsharded chunk rewrites the whole chunk, +resetting its cells beyond the shape to the fill value. A sharded array rewrites a +shard through its inner chunks, each judged against the shard rather than the array +shape, so a write never resets cells beyond the shape. +""" + +from __future__ import annotations + +import itertools +import json +import warnings +from typing import Any, Literal + +import hypothesis.extra.numpy as npst +import hypothesis.strategies as st +import numpy as np +import pytest +from hypothesis import event, note +from hypothesis.stateful import ( + RuleBasedStateMachine, + initialize, + invariant, + precondition, + rule, +) + +import zarr +from zarr.core.buffer import cpu, default_buffer_prototype +from zarr.core.chunk_grids import ChunkGrid +from zarr.core.sync import sync +from zarr.errors import ZarrUserWarning +from zarr.storage import MemoryStore + +pytestmark = [ + pytest.mark.slow_hypothesis, + pytest.mark.filterwarnings("ignore::zarr.core.dtype.common.UnstableSpecificationWarning"), +] + +DTYPE = np.dtype("int16") +METADATA_KEYS = (".zarray", ".zattrs", "zarr.json") +MAX_SIDE = 6 + + +def _rectilinear_dim(extent: int) -> st.SearchStrategy[int | list[int]]: + """A bare step, or an edge list covering `extent` (any edges for extent 0). + + A small local copy of what `zarr.testing.strategies` draws for rectilinear + declarations, so this test does not depend on that module's experimental API. + """ + steps = st.integers(min_value=1, max_value=MAX_SIDE) + if extent == 0: + return steps | st.lists(steps, min_size=1, max_size=3) + if extent == 1: + return steps | st.just([1]) + cuts = st.lists(st.integers(min_value=1, max_value=extent - 1), unique=True, max_size=3) + edges = cuts.map( + lambda c: [b - a for a, b in zip([0, *sorted(c)], [*sorted(c), extent], strict=True)] + ) + return steps | edges + + +async def _list(store: MemoryStore, prefix: str) -> list[str]: + return [key async for key in store.list_prefix(prefix)] + + +class ArrayLifecycle(RuleBasedStateMachine): + def __init__(self) -> None: + super().__init__() + self._rectilinear = zarr.config.set({"array.rectilinear_chunks": True}) + self._rectilinear.__enter__() + self.store = MemoryStore() + self.path = "a" + self.shape: tuple[int, ...] = (0,) + self.fill = 0 + # What the store holds, indexed like the array and extending past its shape. + self.stored: np.ndarray[Any, np.dtype[np.int16]] = np.zeros((0,), dtype=DTYPE) + # Axes stored with chunk size 0; the store warns until its metadata is re-saved. + self.legacy_axes: list[int] = [] + + # -------------------------------------------------------------- creation + @initialize(data=st.data()) + def create(self, data: st.DataObject) -> None: + zarr_format: Literal[2, 3] = data.draw(st.sampled_from([3, 2]), label="zarr_format") + shape = data.draw( + npst.array_shapes(min_dims=1, max_dims=3, min_side=0, max_side=MAX_SIDE), + label="shape", + ) + self.fill = data.draw(st.integers(-3, 3), label="fill_value") + # sampled_from favours early entries; the less common spellings go first. + spellings = ["ints", "-1", "False", "auto"] + if zarr_format == 3: + spellings[:0] = ["sharded", "rectilinear"] + spelling = data.draw(st.sampled_from(spellings), label="chunk spelling") + event(f"chunks: {spelling}") + + chunks: Any + shards: Any = None + if spelling in ("-1", "False"): + chunks = {"-1": -1, "False": False}[spelling] + elif spelling == "auto": + chunks = "auto" + elif spelling == "rectilinear": + chunks = [data.draw(_rectilinear_dim(s)) for s in shape] + if not any(isinstance(c, list) for c in chunks): + chunks[0] = [chunks[0]] if shape[0] == 0 else [shape[0]] + else: + chunks = tuple(data.draw(st.integers(1, 3)) for _ in shape) + if spelling == "sharded": + shards = tuple(c * data.draw(st.integers(1, 3)) for c in chunks) + note(f"create {shape=} {chunks=} {shards=} {zarr_format=} fill={self.fill}") + zarr.create_array( + self.store, + name=self.path, + shape=shape, + chunks=chunks, + shards=shards, + dtype=DTYPE, + fill_value=self.fill, + zarr_format=zarr_format, + ) + self.shape = shape + self.stored = np.full(shape, self.fill, dtype=DTYPE) + + if spelling != "rectilinear" and data.draw(st.booleans(), label="legacy zero"): + # A stored chunk size of 0, as zarr-python wrote for arrays created with a + # zero-length axis; older releases could then grow the axis without storing + # a chunk, so any extent is possible, but zero-length axes come first. + # Sharded arrays store it in the outer grid. + axes = st.lists(st.integers(0, len(shape) - 1), min_size=1, unique=True) + if zero_axes := [axis for axis, extent in enumerate(shape) if extent == 0]: + axes = st.lists(st.sampled_from(zero_axes), min_size=1, unique=True) | axes + self.legacy_axes = data.draw(axes, label="axes stored with chunk size 0") + stored_zero = data.draw(st.sampled_from([0, False]), label="stored zero") + self._rewrite_stored_chunks(zarr_format, self.legacy_axes, stored_zero) + event("legacy zero chunk size") + else: + # Re-saving valid metadata, as the warning tells users to, changes nothing. + arr = self._open() + arr.update_attributes({}) + assert self._open().metadata == arr.metadata + + def _rewrite_stored_chunks( + self, zarr_format: Literal[2, 3], axes: list[int], value: Any + ) -> None: + key = f"{self.path}/{'.zarray' if zarr_format == 2 else 'zarr.json'}" + buf = sync(self.store.get(key, prototype=default_buffer_prototype())) + assert buf is not None + doc = json.loads(buf.to_bytes()) + sizes = ( + doc["chunks"] if zarr_format == 2 else doc["chunk_grid"]["configuration"]["chunk_shape"] + ) + for axis in axes: + sizes[axis] = value + sync(self.store.set(key, cpu.Buffer.from_bytes(json.dumps(doc).encode()))) + + def _open(self) -> zarr.Array[Any]: + with warnings.catch_warnings(record=True) as record: + warnings.simplefilter("always", ZarrUserWarning) + arr = zarr.open_array(self.store, path=self.path, mode="r+") + # A stored chunk size of 0 is read silently on an empty axis; on a non-empty one + # it warns that the axis holds only the fill value. Either way it is upgraded. + warned = any(issubclass(w.category, ZarrUserWarning) for w in record) + must_warn = any(self.shape[axis] > 0 for axis in self.legacy_axes) + assert warned is must_warn, [str(w.message) for w in record] + assert (arr.metadata._stored_document is not None) is bool(self.legacy_axes) + return arr + + # ----------------------------------------------------------------- model + def _cover(self, shape: tuple[int, ...]) -> None: + """Grow `stored` with the fill value so that it covers `shape`.""" + pad = [(0, max(0, n - s)) for n, s in zip(shape, self.stored.shape, strict=True)] + if any(after for _, after in pad): + self.stored = np.pad(self.stored, pad, constant_values=self.fill) + + def _model_resize(self, grid: ChunkGrid, new_shape: tuple[int, ...]) -> None: + """Delete exactly the chunks of `grid` outside the grid for `new_shape`.""" + kept = [] + for dim, old, new in zip(grid.dimensions, self.shape, new_shape, strict=True): + if new >= old: + kept.append(slice(None)) + elif new == 0: + kept.append(slice(0, 0)) + else: + last = dim.index_to_chunk(new - 1) + kept.append(slice(0, dim.chunk_offset(last) + dim.chunk_size(last))) + stored = np.full_like(self.stored, self.fill) + stored[tuple(kept)] = self.stored[tuple(kept)] + self.stored = stored + self.shape = new_shape + self._cover(new_shape) + + def _model_write(self, arr: zarr.Array[Any], region: tuple[slice, ...], values: Any) -> None: + """Write `values`; an unsharded chunk whose in-bounds cells are all written is + rewritten whole.""" + grid = ChunkGrid.from_metadata(arr.metadata) + if arr.shards is None and all(r.stop > r.start for r in region): + # Per axis, each chunk the region touches: (start, stop, all in-bounds cells written). + per_axis: list[list[tuple[int, int, bool]]] = [] + for dim, r, extent in zip(grid.dimensions, region, self.shape, strict=True): + spans = [] + for c in range(dim.index_to_chunk(r.start), dim.index_to_chunk(r.stop - 1) + 1): + lo = dim.chunk_offset(c) + hi = lo + dim.chunk_size(c) + spans.append((lo, hi, r.start <= lo and r.stop >= min(hi, extent))) + per_axis.append(spans) + self._cover(tuple(max(hi for _, hi, _ in axis) for axis in per_axis)) + for combo in itertools.product(*per_axis): + if all(complete for *_, complete in combo): + self.stored[tuple(slice(lo, hi) for lo, hi, _ in combo)] = self.fill + self.stored[region] = values + + # ----------------------------------------------------------------- rules + @rule(data=st.data()) + def append(self, data: st.DataObject) -> None: + arr = self._open() + axes = st.integers(0, len(self.shape) - 1) + if self.legacy_axes: + # What a user of an older release did next: grow an axis stored with chunk size 0. + axes = st.sampled_from(self.legacy_axes) | axes + axis = data.draw(axes, label="axis") + block_shape = list(self.shape) + block_shape[axis] = data.draw(st.integers(0, 4), label="rows") + block = data.draw(npst.arrays(DTYPE, tuple(block_shape)), label="block") + note(f"append {block.shape} along {axis} to {self.shape}") + if self.shape[axis] == 0 and block.shape[axis]: + event("append to a zero-length axis") + if axis in self.legacy_axes and block.shape[axis]: + event("grow an axis stored with chunk size 0") + old_extent = self.shape[axis] + arr.append(block, axis=axis) + self._model_resize(ChunkGrid.from_metadata(arr.metadata), arr.shape) + region = tuple( + slice(old_extent, s) if i == axis else slice(0, s) for i, s in enumerate(arr.shape) + ) + self._model_write(arr, region, block) + # Growing the array rewrites its metadata, which stores any correction. + self.legacy_axes = [] + + @rule(data=st.data()) + def resize(self, data: st.DataObject) -> None: + arr = self._open() + extents = [st.integers(0, MAX_SIDE) for _ in self.shape] + for axis in self.legacy_axes: + # What a user of an older release did next: grow an axis stored with chunk size 0. + extents[axis] = st.integers(self.shape[axis] + 1, MAX_SIDE + 1) | extents[axis] + new_shape = data.draw(st.tuples(*extents), label="new shape") + note(f"resize {self.shape} -> {new_shape}") + if any(new_shape[axis] > self.shape[axis] for axis in self.legacy_axes): + event("grow an axis stored with chunk size 0") + grid = ChunkGrid.from_metadata(arr.metadata) + arr.resize(new_shape) + self._model_resize(grid, new_shape) + self.legacy_axes = [] + + @rule(data=st.data()) + def write(self, data: st.DataObject) -> None: + arr = self._open() + region = tuple( + slice(*sorted(data.draw(st.tuples(st.integers(0, s), st.integers(0, s))))) + for s in self.shape + ) + shape = tuple(r.stop - r.start for r in region) + values = data.draw(npst.arrays(DTYPE, shape), label="values") + note(f"write {region}") + arr[region] = values + self._model_write(arr, region, values) + if all(shape): + # Writing chunks first stores the metadata they are written under; a write + # of nothing stores nothing. + self.legacy_axes = [] + + @precondition(lambda self: self.legacy_axes) + @rule() + def resave_metadata(self) -> None: + """What the warning for an invalid stored chunk size tells users to do: store + the metadata as read.""" + arr = self._open() + read = arr.metadata + arr.update_attributes({}) + self.legacy_axes = [] + assert self._open().metadata == read + + def teardown(self) -> None: + self._rectilinear.__exit__(None, None, None) + + # ------------------------------------------------------------ invariants + @invariant() + def no_chunk_under_an_invalid_document(self) -> None: + """While the stored document is still one that is upgraded on read, which other + readers may reject or read differently, no chunk is stored under it.""" + if self.legacy_axes: + keys = sync(_list(self.store, f"{self.path}/")) + assert set(keys) <= {f"{self.path}/{name}" for name in METADATA_KEYS}, keys + + @invariant() + def matches_model(self) -> None: + arr = self._open() + assert arr.shape == self.shape + expected = self.stored[tuple(slice(0, s) for s in self.shape)] + np.testing.assert_array_equal(np.asarray(arr[...]), expected) + + +TestArrayLifecycle = ArrayLifecycle.TestCase diff --git a/tests/test_chunk_grids.py b/tests/test_chunk_grids.py index 40133700a8..207f121d11 100644 --- a/tests/test_chunk_grids.py +++ b/tests/test_chunk_grids.py @@ -13,6 +13,7 @@ VaryingDimension, _guess_num_chunks_per_axis_shard, _guess_regular_chunks, + full_span_chunk_size, normalize_chunks_1d, normalize_chunks_nd, resolve_outer_and_inner_chunks, @@ -209,6 +210,10 @@ def test_chunk_layout_nested() -> None: ExpectFail( input=([10, 20], 100), exception=ValueError, id="wrong-sum", msg="do not sum to span" ), + # Only a zero-length span keeps edges that do not sum to it. + ExpectFail( + input=([3, 5], 1), exception=ValueError, id="over-sum", msg="do not sum to span 1" + ), # Nested/RLE form for a single dim is rejected with offending indices. ExpectFail( input=([[3, 3], 1], 7), @@ -318,6 +323,11 @@ def test_normalize_chunks_nd_errors(case: ExpectFail[tuple[Any, tuple[int, ...]] output=VaryingDimension([10, 20, 70], extent=100), id="explicit-irregular", ), + # on a zero-length span any non-empty list of positive edges is kept: the + # chunks the axis grows into. + Expect( + input=([3, 5], 0), output=VaryingDimension([3, 5], extent=0), id="explicit-zero-span" + ), ], ids=lambda c: c.id, ) @@ -367,35 +377,39 @@ def test_create_0d_array_auto_shards_with_target_shard_size() -> None: assert arr.shards == () -@pytest.mark.parametrize("chunks", [-1, False], ids=["minus-one", "false"]) -@pytest.mark.parametrize("shape", [(0,), (0, 4), (4, 0)], ids=["1d", "2d-lead", "2d-trail"]) +# -- Zero-length dimensions -- + + +@pytest.mark.parametrize( + ("span", "unit", "expected"), + [(0, 1, 1), (5, 1, 5), (0, 4, 4), (8, 4, 8), (10, 4, 12)], +) +def test_full_span_chunk_size(span: int, unit: int, expected: int) -> None: + """One chunk spanning an axis is the smallest positive multiple of `unit` covering it.""" + assert full_span_chunk_size(span, unit) == expected + + +@pytest.mark.parametrize("chunks", [-1, False, "auto"]) +@pytest.mark.parametrize( + "shape", + [(0,), (0, 4), (4, 0), (0, 0), ()], + ids=["1d", "2d-lead", "2d-trail", "2d-both", "0d"], +) @pytest.mark.parametrize( ("zarr_format", "shards", "target_shard_size_bytes"), - [ - (2, None, None), - (3, None, None), - (3, "auto", None), - (3, "auto", 128 * 1024 * 1024), - ], + [(2, None, None), (3, None, None), (3, "auto", None), (3, "auto", 128 * 1024 * 1024)], ids=["v2", "v3", "v3-auto-shards", "v3-auto-shards-budget"], ) -def test_create_zero_length_array_full_span_chunks( - chunks: int | bool, +def test_create_zero_length_array( + chunks: Any, shape: tuple[int, ...], zarr_format: Literal[2, 3], shards: Literal["auto"] | None, target_shard_size_bytes: int | None, ) -> None: - """`chunks=-1` and `chunks=False` on a zero-length axis must resolve to chunk size 1. - - Both spellings mean "one chunk covering the whole axis". They used to resolve to chunk - size 0 on zero-length axes, which broke every downstream path differently: a ValueError - from the Zarr format 3 chunk grid metadata, a ZeroDivisionError with shards="auto", an - infinite loop with a shard size budget (https://github.com/zarr-developers/zarr-python/issues/4304), - and invalid `chunks: [0]` metadata for Zarr format 2 that silently corrupted reads after - a resize. - """ - expected_chunks = tuple(max(s, 1) for s in shape) + """Every spelling of one chunk spanning the axis gives chunk size 1 on a + zero-length axis, and the array can grow along that axis and shrink back.""" + expected = tuple(max(s, 1) for s in shape) warns = ( pytest.warns(ZarrUserWarning, match="Automatic shard shape inference is experimental") if shards == "auto" @@ -410,23 +424,65 @@ def test_create_zero_length_array_full_span_chunks( shards=shards, zarr_format=zarr_format, ) - assert arr.chunks == expected_chunks - assert arr.shards == (expected_chunks if shards == "auto" else None) - - # The stored chunk grid must be the clamped shape, whichever format wrote it. + assert arr.chunks == expected + assert arr.shards == (None if shards is None else expected) meta = cast(dict[str, Any], arr.metadata.to_dict()) - if zarr_format == 2: - assert meta["chunks"] == expected_chunks - else: - assert meta["chunk_grid"]["configuration"]["chunk_shape"] == expected_chunks + stored = ( + meta["chunks"] if zarr_format == 2 else meta["chunk_grid"]["configuration"]["chunk_shape"] + ) + assert tuple(stored) == expected - # The array must remain usable: grow the empty axis and round-trip data through it. + if not shape: + arr[...] = 7 + assert arr[...] == 7 + return axis = shape.index(0) - grown = tuple(2 if s == 0 else s for s in shape) - arr.append(np.full(grown, 7, dtype="int64"), axis=axis) - assert arr.shape == grown - np.testing.assert_array_equal(arr[...], np.full(grown, 7, dtype="int64")) - resized = tuple(3 if s == 0 else s for s in shape) - arr.resize(resized) - assert arr.shape == resized - assert int(np.asarray(arr[...]).sum()) == 7 * np.prod(grown) + grown = tuple(2 if i == axis else s for i, s in enumerate(shape)) + data = np.full(grown, 7, dtype="int64") + arr.append(data, axis=axis) + np.testing.assert_array_equal(arr[...], data) + arr.resize(shape) + assert np.asarray(arr[...]).shape == shape + + +@pytest.mark.parametrize( + ("shape", "chunks", "shards", "expected"), + [ + ((0, 20), (5, 5), -1, (5, 20)), + ((0,), (4,), False, (4,)), + ((10,), (4,), -1, (12,)), + ((8, 0), (4, 3), (-1, 6), (8, 6)), + ], +) +def test_create_full_span_shards( + shape: tuple[int, ...], chunks: tuple[int, ...], shards: Any, expected: tuple[int, ...] +) -> None: + """A shard spanning the axis is a multiple of the inner chunk, even on a zero-length axis.""" + arr = zarr.create_array(store={}, shape=shape, chunks=chunks, shards=shards, dtype="int8") + assert arr.shards == expected + assert arr.chunks == chunks + + +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_create_zero_chunk_rejected(zarr_format: Literal[2, 3]) -> None: + """An explicit chunk size of 0 is rejected up front, even for a zero-length axis.""" + with pytest.raises(ValueError, match="Chunk size must be positive or -1, got 0"): + zarr.create_array(store={}, shape=(0,), chunks=(0,), dtype="int64", zarr_format=zarr_format) + + +def test_rectilinear_zero_extent_matches_resize() -> None: + """Creating a rectilinear axis at length 0 equals resizing one down to 0.""" + with zarr.config.set({"array.rectilinear_chunks": True}): + created = zarr.create_array(store={}, shape=(0,), chunks=[[2, 2]], dtype="int64") + resized = zarr.create_array(store={}, shape=(4,), chunks=[[2, 2]], dtype="int64") + resized.resize((0,)) + created_meta = cast(dict[str, Any], created.metadata.to_dict()) + resized_meta = cast(dict[str, Any], resized.metadata.to_dict()) + assert created_meta["chunk_grid"] == resized_meta["chunk_grid"] + assert created_meta["shape"] == resized_meta["shape"] == (0,) + + created.append(np.arange(3, dtype="int64")) + resized.append(np.arange(3, dtype="int64")) + np.testing.assert_array_equal(created[...], np.arange(3)) + np.testing.assert_array_equal(resized[...], np.arange(3)) + assert created.write_chunk_sizes == resized.write_chunk_sizes == ((2, 1),) diff --git a/tests/test_codecs/test_sharding.py b/tests/test_codecs/test_sharding.py index e2619a6ce2..831285ae9f 100644 --- a/tests/test_codecs/test_sharding.py +++ b/tests/test_codecs/test_sharding.py @@ -760,17 +760,24 @@ def test_structured_dtype_fill_value() -> None: assert np.array_equal(arr[:], expected) -def test_pickle() -> None: +@pytest.mark.parametrize( + "codec", + [ + ShardingCodec(chunk_shape=(8, 8)), + ShardingCodec(chunk_shape=(8, 8), subchunk_write_order="lexicographic"), + ShardingCodec(chunk_shape=(0,)), + ], + ids=["default", "lexicographic", "zero-chunk-size"], +) +def test_pickle(codec: ShardingCodec) -> None: """ShardingCodec round-trips through pickle, including the non-serialized ``subchunk_write_order`` (which ``to_dict`` omits and which must not silently - revert to the ``morton`` default).""" - codec = ShardingCodec(chunk_shape=(8, 8)) - assert pickle.loads(pickle.dumps(codec)) == codec - - ordered = ShardingCodec(chunk_shape=(8, 8), subchunk_write_order="lexicographic") - restored = pickle.loads(pickle.dumps(ordered)) - assert restored == ordered - assert restored.subchunk_write_order == "lexicographic" + revert to the ``morton`` default), and an inner chunk size of 0, which the + constructor accepts.""" + restored = pickle.loads(pickle.dumps(codec)) + assert restored == codec + assert restored.chunk_shape == codec.chunk_shape + assert restored.subchunk_write_order == codec.subchunk_write_order @pytest.mark.parametrize("store", ["local", "memory"], indirect=["store"]) diff --git a/tests/test_group.py b/tests/test_group.py index 31fbd138cd..b6da55fa2a 100644 --- a/tests/test_group.py +++ b/tests/test_group.py @@ -1919,6 +1919,30 @@ async def test_create_hierarchy( assert expected_meta == {k: v.metadata for k, v in created.items()} +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_create_hierarchy_unbuildable_node_leaves_store_untouched( + monkeypatch: pytest.MonkeyPatch, zarr_format: ZarrFormat +) -> None: + """`create_hierarchy` builds every node before it deletes or stores anything, so a + node that cannot be built fails with the store untouched, even when overwriting.""" + store = MemoryStore() + group = zarr.create_group(store, zarr_format=zarr_format) + group.create_array("a", shape=(2,), chunks=(1,), dtype="int8")[:] = [1, 2] + before = dict(store._store_dict) + + def unbuildable(**kwargs: object) -> None: + raise RuntimeError("cannot build this node") + + monkeypatch.setattr(zarr.core.group, "_build_node", unbuildable) + with pytest.raises(RuntimeError, match="cannot build this node"): + dict( + zarr.create_hierarchy( + store=store, nodes={"a": GroupMetadata(zarr_format=zarr_format)}, overwrite=True + ) + ) + assert store._store_dict == before + + @pytest.mark.parametrize("store", ["memory"], indirect=True) @pytest.mark.parametrize("extant_node", ["array", "group"]) @pytest.mark.parametrize("impl", ["async", "sync"]) diff --git a/tests/test_metadata/test_io.py b/tests/test_metadata/test_io.py new file mode 100644 index 0000000000..7fb4a06914 --- /dev/null +++ b/tests/test_metadata/test_io.py @@ -0,0 +1,214 @@ +"""Tests for comparing stored metadata documents with the documents metadata would store, +and for storing only the documents that differ.""" + +from __future__ import annotations + +import json +from typing import TYPE_CHECKING, Any, Literal + +import pytest + +import zarr +from zarr.core.buffer import cpu +from zarr.core.metadata import ArrayV2Metadata, ArrayV3Metadata +from zarr.core.metadata.io import ( + ABSENT, + ARRAY_DOCUMENTS, + DocumentChange, + diff_documents, + read_documents, + upsert_metadata, +) +from zarr.core.sync import sync +from zarr.storage import MemoryStore, StorePath + +if TYPE_CHECKING: + from zarr.abc.store import Store + from zarr.core.buffer import Buffer + from zarr.core.common import JSON + +V3_DOC: dict[str, JSON] = { + "zarr_format": 3, + "node_type": "array", + "shape": [3], + "chunk_grid": {"name": "regular", "configuration": {"chunk_shape": [3]}}, + "fill_value": 0, +} + + +def _with(doc: dict[str, JSON], **changes: JSON) -> dict[str, JSON]: + return {**doc, **changes} + + +@pytest.mark.parametrize( + ("stored", "new", "expected"), + [ + ({"zarr.json": V3_DOC}, {"zarr.json": V3_DOC}, ()), + ( + {"zarr.json": V3_DOC}, + {"zarr.json": _with(V3_DOC, fill_value=1)}, + ((("zarr.json", "fill_value"), 0, 1),), + ), + ( + { + "zarr.json": _with( + V3_DOC, chunk_grid={"name": "regular", "configuration": {"chunk_shape": [0]}} + ) + }, + {"zarr.json": V3_DOC}, + ((("zarr.json", "chunk_grid", "configuration", "chunk_shape", 0), 0, 3),), + ), + ( + {"zarr.json": V3_DOC}, + {"zarr.json": _with(V3_DOC, shape=[3, 4])}, + ((("zarr.json", "shape", 1), ABSENT, 4),), + ), + ( + {"zarr.json": _with(V3_DOC, attributes={"a": 1})}, + {"zarr.json": _with(V3_DOC, dimension_names=["x"])}, + ( + (("zarr.json", "attributes"), {"a": 1}, ABSENT), + (("zarr.json", "dimension_names"), ABSENT, ["x"]), + ), + ), + ( + {"zarr.json": _with(V3_DOC, fill_value=True)}, + {"zarr.json": _with(V3_DOC, fill_value=1)}, + ((("zarr.json", "fill_value"), True, 1),), + ), + ( + {"zarr.json": _with(V3_DOC, fill_value=1.0)}, + {"zarr.json": _with(V3_DOC, fill_value=1)}, + ((("zarr.json", "fill_value"), 1.0, 1),), + ), + ( + {"zarr.json": _with(V3_DOC, fill_value=float("nan"))}, + {"zarr.json": _with(V3_DOC, fill_value=float("nan"))}, + (), + ), + ( + {"zarr.json": _with(V3_DOC, shape={"0": 3})}, + {"zarr.json": V3_DOC}, + ((("zarr.json", "shape"), {"0": 3}, [3]),), + ), + ( + {".zarray": {"shape": [3], "chunks": [0]}}, + {".zarray": {"shape": [3], "chunks": [3]}, ".zattrs": {}}, + (((".zarray", "chunks", 0), 0, 3), ((".zattrs",), ABSENT, {})), + ), + ], + ids=[ + "identical", + "changed-scalar", + "changed-nested-list", + "longer-list", + "removed-and-added-key", + "bool-is-not-int", + "float-is-not-int", + "nan-is-nan", + "object-is-not-array", + "v2-two-documents", + ], +) +def test_diff_documents( + stored: dict[str, JSON], new: dict[str, JSON], expected: tuple[Any, ...] +) -> None: + """Documents are compared value by value, each change named by its JSON path under + the document's key; values other than objects and arrays are identical only if + their JSON encodings are. `diff_documents` is total over JSON: it raises no error.""" + assert diff_documents(stored, new) == tuple(DocumentChange(*change) for change in expected) + + +class _CountingStore(MemoryStore): + """A memory store that counts the values set in it.""" + + sets = 0 + + async def set(self, key: str, value: Buffer, byte_range: tuple[int, int] | None = None) -> None: + self.sets += 1 + await super().set(key, value, byte_range) + + +def _documents(store: Store) -> dict[str, Any]: + assert isinstance(store, MemoryStore) + return {key: json.loads(value.to_bytes()) for key, value in store._store_dict.items()} + + +def _legacy(zarr_format: Literal[2, 3]) -> tuple[StorePath, ArrayV2Metadata | ArrayV3Metadata]: + """An array stored with chunk shape `[0]`, and the metadata its upgrade reads.""" + store = _CountingStore() + array = zarr.create_array( + store, shape=(3,), chunks=(3,), dtype="int16", zarr_format=zarr_format + ) + key = ".zarray" if zarr_format == 2 else "zarr.json" + doc = _documents(store)[key] + if zarr_format == 2: + doc["chunks"] = [0] + else: + doc["chunk_grid"]["configuration"]["chunk_shape"] = [0] + sync(store.set(key, cpu.Buffer.from_bytes(json.dumps(doc).encode()))) + return StorePath(store), array.metadata + + +def _upsert( + store_path: StorePath, metadata: ArrayV2Metadata | ArrayV3Metadata +) -> tuple[DocumentChange, ...]: + """Upsert `metadata` against the documents stored at `store_path`.""" + stored = sync(read_documents(store_path, ARRAY_DOCUMENTS[metadata.zarr_format])) + return sync(upsert_metadata(store_path, metadata, stored)) + + +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_upsert_metadata_stores_documents_that_differ(zarr_format: Literal[2, 3]) -> None: + """The documents that differ from the stored ones are stored, and the changes are + returned.""" + store_path, metadata = _legacy(zarr_format) + assert isinstance(store_path.store, _CountingStore) + store_path.store.sets = 0 + key = ".zarray" if zarr_format == 2 else "zarr.json" + chunk_path = ( + ("chunks", 0) if zarr_format == 2 else ("chunk_grid", "configuration", "chunk_shape", 0) + ) + + changes = _upsert(store_path, metadata) + + assert changes == (DocumentChange((key, *chunk_path), 0, 3),) + assert store_path.store.sets == 1 + stored = _documents(store_path.store) + assert stored[key] == json.loads(metadata.to_buffer_dict(cpu.buffer_prototype)[key].to_bytes()) + + +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_upsert_metadata_identical_stores_nothing(zarr_format: Literal[2, 3]) -> None: + """Metadata identical to what is stored stores nothing.""" + store = _CountingStore() + array = zarr.create_array( + store, shape=(3,), chunks=(3,), dtype="int16", zarr_format=zarr_format + ) + store.sets = 0 + + assert _upsert(StorePath(store), array.metadata) == () + assert store.sets == 0 + + +def test_upsert_metadata_unstorable_leaves_store_untouched(monkeypatch: pytest.MonkeyPatch) -> None: + """Metadata that cannot be encoded fails before the store is written.""" + store_path, metadata = _legacy(3) + before = _documents(store_path.store) + + def refuse(*args: object) -> None: + raise ValueError("cannot be stored") + + monkeypatch.setattr(ArrayV3Metadata, "to_buffer_dict", refuse) + with pytest.raises(ValueError, match="cannot be stored"): + _upsert(store_path, metadata) + assert _documents(store_path.store) == before + + +def test_upsert_metadata_stored_document_not_an_object() -> None: + """A stored document that is not a JSON object is not overwritten.""" + store_path, metadata = _legacy(3) + sync(store_path.store.set("zarr.json", cpu.Buffer.from_bytes(b"[]"))) + with pytest.raises(TypeError, match="Expected a JSON object, got list"): + _upsert(store_path, metadata) + assert _documents(store_path.store) == {"zarr.json": []} diff --git a/tests/test_metadata/test_upgrades.py b/tests/test_metadata/test_upgrades.py new file mode 100644 index 0000000000..cf33aa2f15 --- /dev/null +++ b/tests/test_metadata/test_upgrades.py @@ -0,0 +1,1173 @@ +"""Tests for the upgrades that read invalid stored array metadata documents.""" + +from __future__ import annotations + +import asyncio +import dataclasses +import json +import re +import warnings +from typing import TYPE_CHECKING, Any, Literal, cast + +import numpy as np +import pytest + +import zarr +from zarr.codecs import ShardingCodec +from zarr.codecs.numcodecs import Quantize +from zarr.core.array import AsyncArray +from zarr.core.group import ConsolidatedMetadata +from zarr.core.metadata import ArrayV2Metadata, ArrayV3Metadata +from zarr.core.metadata.upgrades import ( + RESAVE_HINT, + upgrade_array_document, +) +from zarr.core.metadata.v3 import RectilinearChunkGridMetadata, RegularChunkGridMetadata +from zarr.core.sync import sync +from zarr.dtype import Int16 +from zarr.errors import ZarrUserWarning +from zarr.storage import LocalStore, MemoryStore, StorePath +from zarr.storage._common import make_store_path + +if TYPE_CHECKING: + from collections.abc import Callable + from pathlib import Path + + from zarr.core.common import JSON, ZarrFormat + from zarr.types import AnyArray + + +def _v2_doc(shape: list[int], chunks: list[Any]) -> dict[str, JSON]: + return { + "zarr_format": 2, + "shape": shape, + "chunks": chunks, + "dtype": " dict[str, JSON]: + bytes_codec: dict[str, JSON] = {"name": "bytes", "configuration": {"endian": "little"}} + codecs: list[JSON] = [bytes_codec] + if inner is not None: + codecs = [ + { + "name": "sharding_indexed", + "configuration": { + "chunk_shape": inner, + "codecs": [bytes_codec], + "index_codecs": [bytes_codec, {"name": "crc32c"}], + "index_location": "end", + }, + } + ] + return { + "zarr_format": 3, + "node_type": "array", + "shape": shape, + "data_type": "int16", + "chunk_grid": {"name": "regular", "configuration": {"chunk_shape": chunk_shape}}, + "chunk_key_encoding": {"name": "default", "configuration": {"separator": "/"}}, + "fill_value": 0, + "codecs": codecs, + } + + +def _stored_chunks(doc: dict[str, Any]) -> Any: + return ( + doc["chunks"] + if doc["zarr_format"] == 2 + else doc["chunk_grid"]["configuration"]["chunk_shape"] + ) + + +def _chunk_shapes(metadata: ArrayV2Metadata | ArrayV3Metadata) -> tuple[Any, ...]: + """The chunk shape of `metadata`, then the inner chunk shape of each sharding codec, + from the outermost in.""" + if isinstance(metadata, ArrayV2Metadata): + return (metadata.chunks,) + assert isinstance(metadata.chunk_grid, RegularChunkGridMetadata) + shapes: list[Any] = [metadata.chunk_grid.chunk_shape] + codecs: tuple[Any, ...] = metadata.codecs + while sharding := next((c for c in codecs if isinstance(c, ShardingCodec)), None): + shapes.append(sharding.chunk_shape) + codecs = sharding.codecs + return tuple(shapes) + + +def _nested_sharded_doc(inner: list[Any], nested: list[Any]) -> dict[str, JSON]: + """A Zarr format 3 document of shape `[8]` in one shard of chunk shape `inner`, whose + codecs shard each chunk again, in chunks of shape `nested`.""" + doc = _v3_doc([8], [8], inner=inner) + outer = cast("dict[str, Any]", cast("list[JSON]", doc["codecs"])[0]) + configuration = outer["configuration"] + configuration["codecs"] = [{**outer, "configuration": {**configuration, "chunk_shape": nested}}] + return doc + + +@pytest.mark.parametrize( + ("doc", "expected", "upgraded", "warning"), + [ + (_v2_doc([10, 10], [4, 5]), ((4, 5),), False, None), + (_v3_doc([0, 0], [1, 1]), ((1, 1),), False, None), + (_v3_doc([10], [4], inner=[2]), ((4,), (2,)), False, None), + (_v2_doc([0, 4], [0, 4]), ((1, 4),), True, None), + (_v3_doc([0], [False]), ((1,),), True, None), + (_v2_doc([5], [True]), ((1,),), True, None), + (_v3_doc([5, 4], [True, 4]), ((1, 4),), True, None), + ( + _v2_doc([3], [0]), + ((3,),), + True, + ( + r"^The stored chunk shape \[0\] is invalid: .* read as \[3\], reading 0 in " + r"dimension 0 as one chunk spanning the dimension \(3\), and .* holds only " + r"its fill value\.$" + ), + ), + ( + _v3_doc([4, 3], [4, 0]), + ((4, 3),), + True, + r"reading 0 in dimension 1 as .* holds only its fill value\.$", + ), + ( + _v2_doc([0, 3], [0, 0]), + ((1, 3),), + True, + r"read as \[1, 3\], reading 0 in dimension 1 as .* holds only its fill value\.$", + ), + (_v3_doc([0], [0], inner=[4]), ((4,), (4,)), True, None), + ( + _v3_doc([10], [0], inner=[4]), + ((12,), (4,)), + True, + r"spanning the dimension \(12\), and .* holds only its fill value", + ), + (_v3_doc([0, 3], [0, 3], inner=[2, 3]), ((2, 3), (2, 3)), True, None), + (_v3_doc([5], [True], inner=[True]), ((1,), (1,)), True, None), + (_nested_sharded_doc([4], [2]), ((8,), (4,), (2,)), False, None), + (_nested_sharded_doc([4], [True]), ((8,), (4,), (1,)), True, None), + ], + ids=[ + "v2-valid", + "v3-valid-empty-axes", + "v3-valid-sharded", + "v2-zero-empty-axis", + "v3-false-empty-axis", + "v2-true", + "v3-true", + "v2-zero-grown-axis", + "v3-zero-grown-axis", + "v2-zero-empty-and-grown-axes", + "v3-sharded-zero-empty-axis", + "v3-sharded-zero-grown-axis", + "v3-sharded-zero-2d", + "v3-sharded-true-inner-and-outer", + "v3-nested-sharded-valid", + "v3-nested-sharded-true", + ], +) +def test_upgrade_array_document( + doc: dict[str, JSON], expected: tuple[Any, ...], upgraded: bool, warning: str | None +) -> None: + """Valid documents pass unchanged. A stored chunk size of 0 or `false` is read as one + chunk spanning the axis (a multiple of the inner chunk when sharded) and `true` as 1, + in the chunk shape and in the inner chunk shape of every sharding codec, nested or + not. `from_dict` marks the metadata of an upgraded document; it warns once, naming + the array, only where a chunk size of 0 was stored for a non-empty axis (which then + holds only its fill value), saying how that part was read and how to re-save. The + other readings give what zarr read before, so they are silent.""" + upgraded_doc, readings = upgrade_array_document(doc, cast("ZarrFormat", doc["zarr_format"])) + assert { + k: v for k, v in upgraded_doc.items() if k not in ("chunks", "chunk_grid", "codecs") + } == {k: v for k, v in doc.items() if k not in ("chunks", "chunk_grid", "codecs")} + assert bool(readings) is upgraded + if not upgraded: + assert upgraded_doc is doc + metadata_cls = ArrayV2Metadata if doc["zarr_format"] == 2 else ArrayV3Metadata + with warnings.catch_warnings(record=True) as record: + warnings.simplefilter("always") + metadata = metadata_cls.from_dict(dict(doc), path="group/array") + assert _chunk_shapes(metadata) == expected + assert metadata._stored_document == (doc if upgraded else None) + messages = [str(w.message) for w in record] + if warning is None: + assert messages == [] + else: + [message] = messages + assert message.startswith("Array 'group/array': ") + assert message.endswith(RESAVE_HINT) + assert re.search( + warning, message.removeprefix("Array 'group/array': ").removesuffix(f" {RESAVE_HINT}") + ) + + +@pytest.mark.parametrize( + ("doc", "error"), + [ + (_v2_doc([4], [0]) | {"order": "Z"}, "Failed to parse input for 'order'"), + (_v3_doc([4], [0], inner=[2, 2]), "need to have the same number of dimensions"), + ], + ids=["v2", "v3-sharded"], +) +def test_invalid_upgraded_document_raises_without_warning(doc: dict[str, JSON], error: str) -> None: + """A document the upgrades read that the metadata constructor then rejects raises + that error, without first warning how it was read.""" + metadata_cls = ArrayV2Metadata if doc["zarr_format"] == 2 else ArrayV3Metadata + assert upgrade_array_document(doc, cast("ZarrFormat", doc["zarr_format"]))[1] + with warnings.catch_warnings(): + warnings.simplefilter("error", ZarrUserWarning) + with pytest.raises(ValueError, match=error): + metadata_cls.from_dict(doc) + + +def _read_strictly(doc: dict[str, JSON]) -> ArrayV2Metadata | ArrayV3Metadata: + """Read `doc`, failing on any warning that it was upgraded.""" + metadata_cls = ArrayV2Metadata if doc["zarr_format"] == 2 else ArrayV3Metadata + with warnings.catch_warnings(): + warnings.simplefilter("error", ZarrUserWarning) + return metadata_cls.from_dict(doc) + + +def _open_strictly(path: Path, mode: Literal["r", "a", "r+"] = "r") -> AnyArray: + """Open the array at `path`, failing on any warning that its document was upgraded.""" + with warnings.catch_warnings(): + warnings.simplefilter("error", ZarrUserWarning) + array = zarr.open_array(store=path, mode=mode) + assert isinstance(array, zarr.Array) + return array + + +def test_stored_negative_chunk_size_rejected() -> None: + """No known writer stored a negative chunk size: it is rejected, not upgraded.""" + with pytest.raises(ValueError, match="^Expected all values to be non-negative"): + _read_strictly(_v2_doc([4], [-1])) + + +def test_stored_chunk_shape_ndim_mismatch_rejected() -> None: + """A chunk shape with the wrong number of dimensions is not upgraded, so its 0 is + rejected.""" + with pytest.raises(ValueError, match="^Dimension 0: chunk edge length must be >= 1, got 0$"): + _read_strictly(_v3_doc([4, 4], [0])) + + +@pytest.mark.parametrize("inner", [[0], [False]]) +def test_stored_zero_chunk_size_of_shard_with_invalid_inner_chunk_shape_rejected( + inner: list[Any], +) -> None: + """A stored chunk size of 0 of a sharded array is read in multiples of the inner + chunk size; if that is not an integer of at least 1, the 0 is not upgraded, so it + is rejected.""" + with pytest.raises(ValueError, match="^Dimension 0: chunk edge length must be >= 1, got 0$"): + _read_strictly(_v3_doc([4], [0], inner=inner)) + + +def _rectilinear_doc(shape: list[int], chunk_shapes: list[Any]) -> dict[str, JSON]: + return _v3_doc(shape, [1] * len(shape)) | { + "chunk_grid": { + "name": "rectilinear", + "configuration": {"kind": "inline", "chunk_shapes": chunk_shapes}, + } + } + + +@pytest.mark.parametrize( + ("chunk_shapes", "expected", "upgraded"), + [ + ([[[4, 2]], [[5, 2]]], ((4, 4), (5, 5)), False), + ([[[4.0, 2]], [[5, 2]]], ((4, 4), (5, 5)), True), + ([[3.0, 5.0], [[5, 2]]], ((3, 5), (5, 5)), True), + ([[[4.0, 2], 2.0], [[5, 2]]], ((4, 4, 2), (5, 5)), True), + ([[[4.0, 2]], [10]], ((4, 4), (10,)), True), + ([[[4.0, 2]], [5.0, 5]], ((4, 4), (5, 5)), True), + ([[True, 4], [[5, 2]]], ((1, 4), (5, 5)), True), + ([[[True, 2], 3], [[5, 2]]], ((1, 1, 3), (5, 5)), True), + ], + ids=[ + "valid", + "float-rle-size", + "float-edges", + "float-rle-size-and-edge", + "float-sharded-outer", + "float-2d", + "true-edge", + "true-rle-size", + ], +) +def test_read_invalid_edges_in_rectilinear_grid( + chunk_shapes: list[Any], expected: tuple[tuple[int, ...], ...], upgraded: bool +) -> None: + """A stored rectilinear chunk grid whose explicit edges or run-length encoded sizes + are integral floats or JSON `true`, as zarr-python wrote them when given float or + `True` edges, is read with those edges as the `int`s they equal. `from_dict` marks + the metadata as upgraded, silently: zarr read these edges so before.""" + shape = [sum(edges) for edges in expected] + doc = _rectilinear_doc(shape, chunk_shapes) + with zarr.config.set({"array.rectilinear_chunks": True}): + metadata = _read_strictly(doc) + assert metadata.chunk_grid == RectilinearChunkGridMetadata(chunk_shapes=expected) + assert metadata._stored_document == (doc if upgraded else None) + + +@pytest.mark.parametrize( + ("doc", "error"), + [ + (_v3_doc([20], [10.0]), "Dimension 0: chunk edge length must be an int, got 10.0"), + (_v3_doc([4], [4.5]), "Dimension 0: chunk edge length must be an int, got 4.5"), + ( + _v3_doc([8], [4], inner=[2.0]), + "Expected an iterable of integers. Got [2.0] instead.", + ), + (_v2_doc([20], [10.0]), "Expected an iterable of integers. Got [10.0] instead."), + ( + _rectilinear_doc([8], [[[4, 2.0]]]), + "Dimension 0: RLE repeat count must be an int, got 2.0", + ), + ( + _rectilinear_doc([8], [4.0]), + "Dimension 0: chunk edge length must be an int, got 4.0", + ), + ( + _rectilinear_doc([8], [[0.0, 8]]), + "Dimension 0: chunk edge length must be an int, got 0.0", + ), + ], + ids=[ + "regular", + "regular-fractional", + "sharding-inner", + "v2", + "rle-count", + "rectilinear-bare", + "rectilinear-zero", + ], +) +def test_stored_float_chunk_size_rejected(doc: dict[str, JSON], error: str) -> None: + """A float chunk size is read only where zarr-python stored one, as an edge of at + least 1 of a rectilinear chunk grid. Anywhere else it is rejected, as zarr 3.4.0 + rejected a stored regular chunk size of `10.0`.""" + with zarr.config.set({"array.rectilinear_chunks": True}), pytest.raises(TypeError) as info: + _read_strictly(doc) + assert info.match(re.escape(error)) + + +@pytest.mark.parametrize( + ("chunks", "stored", "resaved"), + [ + ([[4, 4], [5, 5]], [[[4.0, 2]], [[5, 2]]], "[[[4, 2]], [[5, 2]]]"), + ([[1, 3, 4], [5, 5]], [[True, 3, 4], [[5, 2]]], "[[1, 3, 4], [[5, 2]]]"), + ], + ids=["float", "true"], +) +def test_invalid_edges_round_trip( + tmp_path: Path, chunks: list[list[int]], stored: list[Any], resaved: str +) -> None: + """A store whose rectilinear chunk grid holds the float or `true` edges zarr-python + wrote opens silently, reads its data, and stores its edges as `int`s before the + first write.""" + path = tmp_path / "rectilinear.zarr" + data = np.arange(80, dtype="int16").reshape(8, 10) + with zarr.config.set({"array.rectilinear_chunks": True}): + zarr.create_array(path, shape=data.shape, chunks=chunks, dtype="int16")[...] = data + _rewrite_doc( + path, 3, lambda doc: doc["chunk_grid"]["configuration"].update(chunk_shapes=stored) + ) + arr = _open_strictly(path, mode="a") + np.testing.assert_array_equal(arr[...], data) + arr[0, 0] = -1 + written = json.loads((path / "zarr.json").read_text())["chunk_grid"]["configuration"] + assert json.dumps(written["chunk_shapes"]) == resaved + data[0, 0] = -1 + np.testing.assert_array_equal(_open_strictly(path)[...], data) + + +def test_v2_constructor_rejects_chunks_of_wrong_length() -> None: + with pytest.raises(ValueError, match="`chunks` has length 1, but `shape` has length 2"): + ArrayV2Metadata(shape=(4, 4), chunks=(2,), dtype=Int16(), fill_value=0, order="C") + + +def _v2_metadata(chunks: Any) -> ArrayV2Metadata: + return ArrayV2Metadata(shape=(4,), chunks=chunks, dtype=Int16(), fill_value=0, order="C") + + +def _rectilinear(chunk_shapes: tuple[Any, ...]) -> RectilinearChunkGridMetadata: + with zarr.config.set({"array.rectilinear_chunks": True}): + return RectilinearChunkGridMetadata(chunk_shapes=chunk_shapes) + + +def _rectilinear_from_dict(chunk_shapes: list[Any]) -> RectilinearChunkGridMetadata: + with zarr.config.set({"array.rectilinear_chunks": True}): + return RectilinearChunkGridMetadata.from_dict( + { + "name": "rectilinear", + "configuration": {"kind": "inline", "chunk_shapes": chunk_shapes}, + } + ) + + +def _second_edge(grid: RectilinearChunkGridMetadata) -> int: + edges = grid.chunk_shapes[0] + assert isinstance(edges, tuple) + return edges[1] + + +CHUNK_EDGE_SITES: dict[str, Callable[[Any], object]] = { + "regular": lambda size: RegularChunkGridMetadata(chunk_shape=(size,)).chunk_shape[0], + "rectilinear-bare": lambda size: _rectilinear((size,)).chunk_shapes[0], + "rectilinear-edge": lambda size: _second_edge(_rectilinear(((4, size),))), + "rectilinear-bare-json": lambda size: _rectilinear_from_dict([size]).chunk_shapes[0], + "rectilinear-edge-json": lambda size: _second_edge(_rectilinear_from_dict([[4, size]])), + "rectilinear-rle-json": lambda size: _second_edge(_rectilinear_from_dict([[[size, 2]]])), +} +"""Each place chunk grid metadata reads a chunk edge length, returning the edge it read.""" + + +@pytest.mark.parametrize("site", CHUNK_EDGE_SITES) +@pytest.mark.parametrize( + ("size", "expected"), + [(4, 4), (True, 1)], + ids=["int", "bool"], +) +def test_metadata_reads_integer_chunk_edge(site: str, size: object, expected: int) -> None: + """Chunk grid metadata reads an `int` or a `bool` as the `int` chunk edge length it + equals.""" + edge = CHUNK_EDGE_SITES[site](size) + assert type(edge) is int + assert edge == expected + + +@pytest.mark.parametrize("site", CHUNK_EDGE_SITES) +@pytest.mark.parametrize( + "size", + [4.0, np.float64(4.0), 4.5, float("inf"), "4", None, np.int64(4)], + ids=["float", "numpy-float", "fractional", "inf", "str", "none", "numpy-int"], +) +def test_metadata_rejects_non_integer_chunk_edge(site: str, size: object) -> None: + """A chunk edge length in metadata built in code is an `int`: a float is rejected, + even an integral one (stored documents with integral floats are read by the + upgrades), and so is a NumPy integer, as zarr 3.4.0 rejected one.""" + with pytest.raises( + TypeError, + match=re.escape(f"Dimension 0: chunk edge length must be an int, got {size!r}"), + ): + CHUNK_EDGE_SITES[site](size) + + +@pytest.mark.parametrize("site", CHUNK_EDGE_SITES) +@pytest.mark.parametrize("size", [0, False, -1]) +def test_metadata_rejects_chunk_edge_below_one(site: str, size: int) -> None: + """Chunk grid metadata built in code is strict: a chunk edge length below 1 is + rejected, without a warning.""" + with warnings.catch_warnings(): + warnings.simplefilter("error", ZarrUserWarning) + with pytest.raises( + ValueError, match=f"Dimension 0: chunk edge length must be >= 1, got {size}" + ): + CHUNK_EDGE_SITES[site](size) + + +@pytest.mark.parametrize("chunk_shape", [4, np.int64(4), None, "44", {"4": 4}]) +def test_regular_chunk_grid_rejects_chunk_shape_not_a_sequence(chunk_shape: Any) -> None: + """A regular chunk shape is an iterable of chunk edge lengths, but not a string or a + mapping, which is rejected as a whole, not entry by entry.""" + with pytest.raises( + TypeError, + match=re.escape( + f"A chunk shape must be an iterable of chunk edge lengths, got {chunk_shape!r}" + ), + ): + RegularChunkGridMetadata(chunk_shape=chunk_shape) + + +def _sharding_chunk_shape(chunks: Any) -> tuple[tuple[int, ...], object]: + codec = ShardingCodec(chunk_shape=chunks) + configuration = cast("dict[str, JSON]", codec.to_dict()["configuration"]) + return codec.chunk_shape, configuration["chunk_shape"] + + +CHUNK_SHAPE_SITES: dict[str, Callable[[Any], tuple[tuple[int, ...], object]]] = { + "v2": lambda chunks: ((md := _v2_metadata(chunks)).chunks, md.to_dict()["chunks"]), + "sharding-inner": _sharding_chunk_shape, +} +"""`ArrayV2Metadata` and `ShardingCodec` read a chunk shape as an array shape, returning +the chunk shape and the value `to_dict` writes for it.""" + + +@pytest.mark.parametrize("site", CHUNK_SHAPE_SITES) +@pytest.mark.parametrize( + ("chunks", "expected"), + [ + ((4,), (4,)), + ([4], (4,)), + (4, (4,)), + (np.int64(4), (4,)), + ((np.int64(4),), (4,)), + (np.array([4]), (4,)), + ((True,), (1,)), + ((0,), (0,)), + ((False,), (0,)), + (range(4, 5), (4,)), + ], +) +def test_chunk_shape_read_as_array_shape( + site: str, chunks: object, expected: tuple[int, ...] +) -> None: + """`ArrayV2Metadata` and `ShardingCodec` read their chunk shape as `parse_shapelike` + reads an array shape: an integer or an iterable of non-negative integers, including + NumPy integers and bools. A chunk size of 0 is written back as given; reading a + stored 0 is `zarr.core.metadata.upgrades`' business.""" + parsed, written = CHUNK_SHAPE_SITES[site](chunks) + assert parsed == expected + assert all(type(size) is int for size in parsed) + assert written == expected + + +@pytest.mark.parametrize("site", CHUNK_SHAPE_SITES) +def test_chunk_shape_read_as_array_shape_rejects_negative(site: str) -> None: + with pytest.raises(ValueError, match="Expected all values to be non-negative"): + CHUNK_SHAPE_SITES[site]((-1,)) + + +@pytest.mark.parametrize("site", CHUNK_SHAPE_SITES) +@pytest.mark.parametrize("chunks", [(4.0,), "4", None]) +def test_chunk_shape_read_as_array_shape_rejects_non_integer(site: str, chunks: object) -> None: + with pytest.raises(TypeError, match="Expected an"): + CHUNK_SHAPE_SITES[site](chunks) + + +def _rewrite_doc(path: Path, zarr_format: Literal[2, 3], edit: Any) -> None: + doc_path = path / (".zarray" if zarr_format == 2 else "zarr.json") + doc = json.loads(doc_path.read_text()) + edit(doc) + doc_path.write_text(json.dumps(doc)) + + +@pytest.mark.parametrize( + ("zarr_format", "shape", "stored", "inner", "expected", "warns"), + [ + (2, (0, 4), [0, 4], None, (1, 4), False), + (3, (5,), [True], None, (1,), False), + (3, (10,), [0], (4,), (12,), True), + ], + ids=["v2-empty-2d", "v3-true", "v3-sharded-grown"], +) +def test_legacy_chunk_size_round_trip( + tmp_path: Path, + zarr_format: Literal[2, 3], + shape: tuple[int, ...], + stored: list[Any], + inner: tuple[int, ...] | None, + expected: tuple[int, ...], + warns: bool, +) -> None: + """A store whose metadata holds a chunk size written by older software opens (with a + warning where a non-empty axis was stored with chunk size 0), reads and appends under + the upgraded grid, and re-saves valid metadata.""" + path = tmp_path / "legacy.zarr" + arr = zarr.create_array( + store=path, + shape=shape, + chunks=inner or expected, + shards=expected if inner else None, + dtype="int16", + fill_value=0, + zarr_format=zarr_format, + ) + data = np.arange(np.prod(shape), dtype="int16").reshape(shape) + arr[...] = data + + def store_legacy(doc: dict[str, Any]) -> None: + if zarr_format == 2: + doc["chunks"] = stored + else: + doc["chunk_grid"]["configuration"]["chunk_shape"] = stored + + _rewrite_doc(path, zarr_format, store_legacy) + + with warnings.catch_warnings(record=True) as record: + warnings.simplefilter("always", ZarrUserWarning) + arr = zarr.open_array(store=path, mode="a") + assert [ + re.match(r"^Array '.*legacy\.zarr': .* is read as", str(w.message)) is not None + for w in record + ] == [True] * warns + assert (arr.shards or arr.chunks) == expected + np.testing.assert_array_equal(arr[...], data) + + block = np.full((2, *shape[1:]), 7, dtype="int16") + arr.append(block) + arr.update_attributes({}) + with warnings.catch_warnings(): + warnings.simplefilter("error", ZarrUserWarning) + reopened = zarr.open_array(store=path) + np.testing.assert_array_equal(reopened[...], np.concatenate([data, block])) + + +@pytest.mark.parametrize( + ("zarr_format", "shape", "inner", "expected"), + [(2, (3,), None, (3,)), (3, (3,), None, (3,)), (3, (10,), (4,), (12,))], + ids=["v2", "v3", "v3-sharded"], +) +@pytest.mark.parametrize("api", ["sync", "async", "async-concurrent"]) +def test_write_stores_upgraded_metadata_first( + tmp_path: Path, + zarr_format: Literal[2, 3], + shape: tuple[int, ...], + inner: tuple[int, ...] | None, + expected: tuple[int, ...], + api: str, +) -> None: + """Writing chunks to an array read from an upgraded document first stores the + upgraded metadata, so readers that do not upgrade (or read it differently) see + the chunks the write stored.""" + path = tmp_path / "legacy.zarr" + zarr.create_array( + store=path, + shape=shape, + chunks=inner or expected, + shards=expected if inner else None, + dtype="int16", + fill_value=0, + zarr_format=zarr_format, + ) + _rewrite_doc( + path, + zarr_format, + lambda doc: ( + doc.update(chunks=[0]) + if zarr_format == 2 + else doc["chunk_grid"]["configuration"].update(chunk_shape=[0]) + ), + ) + data = np.arange(1, shape[0] + 1, dtype="int16") + with pytest.warns(ZarrUserWarning, match="is read as"): + arr = zarr.open_array(store=path, mode="r+") + if api == "sync": + arr[:] = data + elif api == "async": + sync(arr.async_array.setitem(slice(None), data)) + else: + + async def write_twice() -> None: + # Both writes find the metadata not yet stored, and both store it. + await asyncio.gather(*(arr.async_array.setitem(slice(None), data) for _ in "ab")) + + sync(write_twice()) + + assert arr.metadata._stored_document is None + with warnings.catch_warnings(): + warnings.simplefilter("error", ZarrUserWarning) + reopened = zarr.open_array(store=path, mode="r") + assert (reopened.shards or reopened.chunks) == expected + np.testing.assert_array_equal(reopened[...], data) + + +def _legacy_array(path: Path, zarr_format: Literal[2, 3]) -> None: + """Store an array of shape (3,) whose stored chunk shape is `[0]`.""" + zarr.create_array(store=path, shape=(3,), chunks=(3,), dtype="int16", zarr_format=zarr_format) + _rewrite_doc( + path, + zarr_format, + lambda doc: ( + doc.update(chunks=[0]) + if zarr_format == 2 + else doc["chunk_grid"]["configuration"].update(chunk_shape=[0]) + ), + ) + + +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_stale_handle_write_keeps_newer_metadata( + tmp_path: Path, zarr_format: Literal[2, 3] +) -> None: + """A handle read from an upgraded document stores the upgrade of what the store + holds when it first writes chunks; if another handle stored valid metadata since, + it stores no metadata and writes only its chunks.""" + path = tmp_path / "legacy.zarr" + _legacy_array(path, zarr_format) + with pytest.warns(ZarrUserWarning, match="is read as"): + stale = zarr.open_array(store=path, mode="r+") + with pytest.warns(ZarrUserWarning, match="is read as"): + other = zarr.open_array(store=path, mode="r+") + other.append(np.arange(1, 7, dtype="int16")) + other.attrs["x"] = 1 + + stale[0] = 9 + + assert stale.metadata._stored_document is None + reopened = _open_strictly(path) + assert reopened.shape == (9,) + assert reopened.attrs.asdict() == {"x": 1} + np.testing.assert_array_equal(reopened[...], [9, 0, 0, 1, 2, 3, 4, 5, 6]) + + +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_stale_handle_write_keeps_valid_document_as_written( + tmp_path: Path, zarr_format: Literal[2, 3] +) -> None: + """If the document the store holds when a handle read from an upgraded document + first writes chunks needs no upgrade, it is left as written, even where zarr would + encode the same metadata differently (as another implementation may have written it).""" + path = tmp_path / "legacy.zarr" + _legacy_array(path, zarr_format) + with pytest.warns(ZarrUserWarning, match="is read as"): + stale = zarr.open_array(store=path, mode="r+") + # Valid metadata for the same array, as another writer might store it: chunk size + # 3, without the optional members zarr writes, as compact JSON. + doc_path = path / (".zarray" if zarr_format == 2 else "zarr.json") + doc = json.loads(doc_path.read_text()) + for optional in ("dimension_separator", "attributes", "storage_transformers"): + doc.pop(optional, None) + _stored_chunks(doc)[0] = 3 + doc_path.write_text(json.dumps(doc, separators=(",", ":"))) + written = doc_path.read_bytes() + + stale[0] = 9 + + assert doc_path.read_bytes() == written + np.testing.assert_array_equal(_open_strictly(path)[...], [9, 0, 0]) + + +def _store_zero(doc: dict[str, Any]) -> None: + _stored_chunks(doc)[0] = 0 + + +def _resize_to_10(doc: dict[str, Any]) -> None: + doc["shape"] = [10] + + +def _halve_inner_chunk_shape(doc: dict[str, Any]) -> None: + doc["codecs"][0]["configuration"]["chunk_shape"] = [2] + + +@pytest.mark.parametrize( + ("zarr_format", "sharded", "change"), + [(2, False, _resize_to_10), (3, False, _resize_to_10), (3, True, _halve_inner_chunk_shape)], + ids=["v2-resized", "v3-resized", "v3-sharded-inner-chunk-shape"], +) +def test_stale_handle_write_after_chunk_grid_change_raises( + tmp_path: Path, + zarr_format: Literal[2, 3], + sharded: bool, + change: Callable[[dict[str, Any]], None], +) -> None: + """If the document the store holds when a handle read from an upgraded document + first writes chunks lays out chunks differently from the handle's metadata (the + array was resized by software that kept the stored chunk size of 0, which now reads + as a larger chunk, or its inner chunk shape changed), the handle's chunks would not + be found under it: the write raises and stores nothing.""" + path = tmp_path / "legacy.zarr" + if sharded: + zarr.create_array(store=path, shape=(3,), chunks=(4,), shards=(4,), dtype="int16") + _rewrite_doc(path, 3, _store_zero) + else: + _legacy_array(path, zarr_format) + with pytest.warns(ZarrUserWarning, match="is read as"): + stale = zarr.open_array(store=path, mode="r+") + _rewrite_doc(path, zarr_format, change) + documents = {p.name: p.read_bytes() for p in path.iterdir()} + + with pytest.raises(ValueError, match="has changed since this array was opened: reopen"): + stale[0:3] = [7, 8, 9] + + assert stale.metadata._stored_document is not None + assert {p.name: p.read_bytes() for p in path.iterdir()} == documents + + +def _set_attribute(array: AnyArray) -> None: + array.attrs["x"] = 1 + + +def _grow(array: AnyArray) -> None: + array.resize((9,)) + + +@pytest.mark.parametrize("operation", [_set_attribute, _grow], ids=["attrs", "resize"]) +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_failed_metadata_save_keeps_stored_document( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + zarr_format: Literal[2, 3], + operation: Callable[[AnyArray], None], +) -> None: + """If storing an array's metadata fails, the array still stands for the document + the store holds, which still needs its upgrade.""" + path = tmp_path / "legacy.zarr" + _legacy_array(path, zarr_format) + doc_name = ".zarray" if zarr_format == 2 else "zarr.json" + with pytest.warns(ZarrUserWarning, match="is read as"): + arr = zarr.open_array(store=path, mode="r+") + stored = (path / doc_name).read_bytes() + original_set = LocalStore.set + + async def failing_set(self: LocalStore, key: str, *args: Any, **kwargs: Any) -> None: + if key == doc_name: + raise OSError(f"cannot store {key}") + await original_set(self, key, *args, **kwargs) + + monkeypatch.setattr(LocalStore, "set", failing_set) + + with pytest.raises(OSError, match=f"cannot store {re.escape(doc_name)}"): + operation(arr) + + assert arr.metadata._stored_document is not None + assert (path / doc_name).read_bytes() == stored + + +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_write_without_stored_document(zarr_format: Literal[2, 3]) -> None: + """An array read from an upgraded document that no store holds (as + `AsyncArray.from_dict` builds one) writes its chunks as any array does: there is no + stored document to upgrade.""" + store = MemoryStore() + doc = _v2_doc([3], [True]) if zarr_format == 2 else _v3_doc([3], [True]) + array = zarr.Array(AsyncArray.from_dict(StorePath(store), doc)) + upgraded = array.metadata._stored_document is not None + + array[:] = [1, 2, 3] + + assert (upgraded, array.metadata._stored_document) == (True, None) + np.testing.assert_array_equal(array[:], [1, 2, 3]) + assert not [key for key in store._store_dict if key.endswith((".zarray", "zarr.json"))] + + +@pytest.mark.parametrize(("shape", "expected"), [((0,), (1,)), ((3,), (3,))]) +def test_array_from_metadata_with_chunk_size_zero(shape: tuple[int], expected: tuple[int]) -> None: + """`ArrayV2Metadata` accepts a chunk size of 0, as a stored document may hold it. An + array built from such metadata reads it as the upgrades read that document, silently + (no data was read or written under it): `create_hierarchy` stores the metadata as + given and yields such an array, which stores the upgrade before its first write.""" + metadata = ArrayV2Metadata(shape=shape, chunks=(0,), dtype=Int16(), fill_value=0, order="C") + store = MemoryStore() + with warnings.catch_warnings(): + warnings.simplefilter("error", ZarrUserWarning) + nodes = dict(zarr.create_hierarchy(store=store, nodes={"a": metadata})) + array = nodes["a"] + assert isinstance(array, zarr.Array) + assert array.chunks == expected + assert json.loads(store._store_dict["a/.zarray"].to_bytes())["chunks"] == [0] + + array[...] = 1 + + # Writing an empty selection stores no chunks, so it stores no metadata either. + resaved = list(expected) if array.size else [0] + assert json.loads(store._store_dict["a/.zarray"].to_bytes())["chunks"] == resaved + np.testing.assert_array_equal(zarr.open_array(store, path="a")[...], np.ones(shape)) + + +def test_array_from_metadata_with_numpy_scalar_codec_configuration() -> None: + """An array is built from metadata whose codec configuration holds NumPy scalars + (which are not JSON values), as zarr always built one.""" + array = zarr.create_array( + MemoryStore(), shape=(4,), chunks=(2,), dtype="f8", filters=[Quantize(digits=3, dtype="f8")] + ) + assert isinstance(array.metadata, ArrayV3Metadata) + # The codec keeps its configuration as given, though it is typed as JSON. + digits = cast("JSON", np.int64(3)) + codecs = (Quantize(digits=digits, dtype="f8"), *array.metadata.codecs[1:]) + metadata = dataclasses.replace(array.metadata, codecs=codecs) + + assert AsyncArray(metadata, StorePath(MemoryStore())).metadata is metadata + + +def _rewrite_consolidated( + path: Path, zarr_format: Literal[2, 3], name: str, edit: Callable[[dict[str, Any]], None] +) -> None: + """Edit the document of the member `name` in the consolidated metadata at `path`.""" + if zarr_format == 2: + document = json.loads((path / ".zmetadata").read_text()) + edit(document["metadata"][f"{name}/.zarray"]) + (path / ".zmetadata").write_text(json.dumps(document)) + else: + _rewrite_doc(path, 3, lambda doc: edit(doc["consolidated_metadata"]["metadata"][name])) + + +def _consolidated_member(path: Path, zarr_format: Literal[2, 3], name: str) -> Any: + if zarr_format == 2: + return json.loads((path / ".zmetadata").read_text())["metadata"][f"{name}/.zarray"] + return json.loads((path / "zarr.json").read_text())["consolidated_metadata"]["metadata"][name] + + +def _flagged_consolidated_group(path: Path, zarr_format: Literal[2, 3], member: str) -> None: + """A group whose consolidated copy of the array `member` holds the stored chunk size + 0, and whose array `b` is valid.""" + group = zarr.open_group(path, mode="w", zarr_format=zarr_format) + parent, _, name = member.rpartition("/") + (group.require_group(parent) if parent else group).create_array( + name, shape=(3,), chunks=(3,), dtype="int16" + ) + group.create_array("b", shape=(1,), chunks=(1,), dtype="int16") + zarr.consolidate_metadata(path) + _rewrite_consolidated(path, zarr_format, member, _store_zero) + + +def _record_store_access(monkeypatch: pytest.MonkeyPatch) -> list[tuple[str, str]]: + """Record every `get` and `set` of a `LocalStore`, as `(method, key)`.""" + accesses: list[tuple[str, str]] = [] + for method in ("get", "set"): + original = getattr(LocalStore, method) + + async def recording( + self: LocalStore, + key: str, + *args: Any, + _method: str = method, + _original: Any = original, + **kwargs: Any, + ) -> Any: + accesses.append((_method, key)) + return await _original(self, key, *args, **kwargs) + + monkeypatch.setattr(LocalStore, method, recording) + return accesses + + +@pytest.mark.filterwarnings("ignore:Consolidated metadata is currently not part:UserWarning") +@pytest.mark.parametrize("zarr_format", [2, 3]) +@pytest.mark.parametrize("member", ["a", "g/a"]) +@pytest.mark.parametrize("operation", ["attrs", "update_attributes_async", "delete-member"]) +def test_group_write_stores_upgraded_consolidated_member_as_stored( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + zarr_format: Literal[2, 3], + member: str, + operation: str, +) -> None: + """A group write stores only the group's own documents, reading none: it reads and + writes no member document, and stores the consolidated copy of a member read from a + document that had to be upgraded exactly as it was stored, so every reader of the + consolidated metadata reads that member as upgraded again.""" + path = tmp_path / "group.zarr" + _flagged_consolidated_group(path, zarr_format, member) + stored = json.dumps(_consolidated_member(path, zarr_format, member)) + with pytest.warns(ZarrUserWarning, match="is read as"): + group = zarr.open_group(path, mode="r+", use_consolidated=True) + accesses = _record_store_access(monkeypatch) + + if operation == "attrs": + group.attrs["x"] = 1 + elif operation == "update_attributes_async": + sync(group.update_attributes_async({"x": 1})) + else: + del group["b"] + + own = {"zarr.json"} if zarr_format == 3 else {".zgroup", ".zattrs", ".zmetadata"} + assert {method for method, _ in accesses} == {"set"} + assert {key for _, key in accesses} == own + assert json.dumps(_consolidated_member(path, zarr_format, member)) == stored + with pytest.warns(ZarrUserWarning, match="is read as"): + zarr.open_group(path, mode="r", use_consolidated=True) + + +@pytest.mark.filterwarnings("ignore:Consolidated metadata is currently not part:UserWarning") +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_concurrent_deletions_leave_no_upgraded_member( + tmp_path: Path, zarr_format: Literal[2, 3] +) -> None: + """Members deleted concurrently through one consolidated group handle stay deleted, + also those read from documents that had to be upgraded: no group write stores a + member document.""" + path = tmp_path / "group.zarr" + zarr.open_group(path, mode="w", zarr_format=zarr_format) + for name in ("a", "b"): + _legacy_array(path / name, zarr_format) + with pytest.warns(ZarrUserWarning, match="is read as"): + zarr.consolidate_metadata(path) + with pytest.warns(ZarrUserWarning, match="is read as"): + group = zarr.open_group(path, mode="r+", use_consolidated=True) + + async def delete_both() -> None: + await asyncio.gather(*(group._async_group.delitem(name) for name in ("a", "b"))) + + sync(delete_both()) + + assert not (path / "a").exists() + assert not (path / "b").exists() + + +def _consolidated_legacy_member(path: Path, zarr_format: Literal[2, 3]) -> Path: + """A group whose array `a` is stored with chunk shape `[0]`, consolidated; the path + of the array's own document.""" + zarr.open_group(path, mode="w", zarr_format=zarr_format) + _legacy_array(path / "a", zarr_format) + with pytest.warns(ZarrUserWarning, match="is read as"): + zarr.consolidate_metadata(path) + return path / "a" / (".zarray" if zarr_format == 2 else "zarr.json") + + +@pytest.mark.filterwarnings("ignore:Consolidated metadata is currently not part:UserWarning") +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_consolidated_upgraded_member_stored_by_its_first_write( + tmp_path: Path, zarr_format: Literal[2, 3] +) -> None: + """Consolidating a group copies the document of a member that had to be upgraded as + it is stored, and leaves that document as it is. The first chunk write through the + consolidated metadata stores the member's upgrade before its chunks, the one write + that stores it; consolidating again then copies the upgrade.""" + path = tmp_path / "group.zarr" + document = _consolidated_legacy_member(path, zarr_format) + legacy = document.read_bytes() + assert json.dumps(_consolidated_member(path, zarr_format, "a")).encode() == legacy + + with pytest.warns(ZarrUserWarning, match="is read as"): + array = zarr.open_group(path, mode="r+", use_consolidated=True)["a"] + assert isinstance(array, zarr.Array) + array[:] = [7, 8, 9] + + np.testing.assert_array_equal(_open_strictly(path / "a")[...], [7, 8, 9]) + zarr.consolidate_metadata(path) + assert _stored_chunks(_consolidated_member(path, zarr_format, "a")) == [3] + + +@pytest.mark.filterwarnings("ignore:Consolidated metadata is currently not part:UserWarning") +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_consolidated_upgraded_member_write_after_chunk_grid_change_raises( + tmp_path: Path, zarr_format: Literal[2, 3] +) -> None: + """A member read from its consolidated copy, whose own document has since been + resized by software that kept the stored chunk size of 0, would write chunks no + reader finds: its first write raises and stores nothing.""" + path = tmp_path / "group.zarr" + _consolidated_legacy_member(path, zarr_format) + with pytest.warns(ZarrUserWarning, match="is read as"): + array = zarr.open_group(path, mode="r+", use_consolidated=True)["a"] + assert isinstance(array, zarr.Array) + _rewrite_doc(path / "a", zarr_format, _resize_to_10) + documents = {p: p.read_bytes() for p in path.rglob("*") if p.is_file()} + + with pytest.raises(ValueError, match="has changed since this array was opened: reopen"): + array[0:3] = [7, 8, 9] + + assert {p: p.read_bytes() for p in path.rglob("*") if p.is_file()} == documents + + +@pytest.mark.filterwarnings("ignore:Consolidated metadata is currently not part:UserWarning") +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_consolidated_upgraded_member_attributes_kept_by_group_write( + tmp_path: Path, zarr_format: Literal[2, 3] +) -> None: + """Setting attributes of a member read from consolidated metadata stores the member's + upgraded document with them. A later write of the group, whose consolidated metadata + shares the member's metadata, then stores that upgrade: the new attributes, not the + legacy document as it was stored before.""" + path = tmp_path / "group.zarr" + _consolidated_legacy_member(path, zarr_format) + with pytest.warns(ZarrUserWarning, match="is read as"): + group = zarr.open_group(path, mode="r+", use_consolidated=True) + + group["a"].attrs["x"] = 1 + group.attrs["y"] = 2 + + with warnings.catch_warnings(record=True) as record: + warnings.simplefilter("always") + attributes = dict(zarr.open_group(path, mode="r", use_consolidated=True)["a"].attrs) + assert attributes == {"x": 1} + assert _stored_chunks(_consolidated_member(path, zarr_format, "a")) == [3] + assert not [w for w in record if "is read as" in str(w.message)] + + +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_upgraded_metadata_keeps_the_document_it_read(zarr_format: Literal[2, 3]) -> None: + """Metadata read from a document that had to be upgraded keeps that document as it + was read, whatever becomes of the caller's dict (or the objects in it, which the + metadata's attributes may hold): consolidated metadata stores it as read.""" + doc = _v2_doc([0], [0]) if zarr_format == 2 else _v3_doc([0], [0]) + doc["attributes"] = {"k": [1]} + read = json.loads(json.dumps(doc)) + metadata_cls = ArrayV2Metadata if zarr_format == 2 else ArrayV3Metadata + metadata = metadata_cls.from_dict(doc) + + cast("list[int]", metadata.attributes["k"]).append(2) + doc["shape"] = [4] + + assert ConsolidatedMetadata(metadata={"a": metadata}).to_dict()["metadata"] == {"a": read} + + +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_empty_write_stores_no_metadata(tmp_path: Path, zarr_format: Literal[2, 3]) -> None: + """A write of an empty selection stores no chunks, so it stores no metadata either.""" + path = tmp_path / "legacy.zarr" + _legacy_array(path, zarr_format) + documents = {p.name: p.read_bytes() for p in path.iterdir()} + with pytest.warns(ZarrUserWarning, match="is read as"): + array = zarr.open_array(store=path, mode="r+") + + array[0:0] = np.empty(0, dtype="int16") + + assert array.metadata._stored_document is not None + assert {p.name: p.read_bytes() for p in path.iterdir()} == documents + + +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_async_array_from_dict_names_array(tmp_path: Path, zarr_format: Literal[2, 3]) -> None: + """`AsyncArray.from_dict` names the array at its store path in the upgrade warning.""" + store_path = sync(make_store_path(tmp_path / "legacy.zarr")) + doc = _v2_doc([4], [0]) if zarr_format == 2 else _v3_doc([4], [0]) + with pytest.warns(ZarrUserWarning, match=f"^Array {re.escape(repr(str(store_path)))}: "): + AsyncArray.from_dict(store_path, doc) + + +@pytest.mark.filterwarnings("ignore:Consolidated metadata is currently not part:UserWarning") +@pytest.mark.parametrize("zarr_format", [2, 3]) +def test_legacy_chunk_size_consolidated(tmp_path: Path, zarr_format: Literal[2, 3]) -> None: + """Consolidated metadata goes through the same upgrade as the arrays' own documents, + with one warning naming each array by its path; re-saving the arrays and + consolidating again leaves a group that opens without a warning.""" + path = tmp_path / "group.zarr" + group = zarr.open_group(path, mode="w", zarr_format=zarr_format) + names = ("a", "b") + for name in names: + group.create_array(name, shape=(2,), chunks=(2,), dtype="int32") + zarr.consolidate_metadata(path) + + # What zarr-python wrote for `chunks=(0,)`, in both copies. + for name in names: + if zarr_format == 2: + _rewrite_doc(path / name, 2, lambda doc: doc.update(chunks=[0])) + zmetadata = json.loads((path / ".zmetadata").read_text()) + zmetadata["metadata"][f"{name}/.zarray"]["chunks"] = [0] + (path / ".zmetadata").write_text(json.dumps(zmetadata)) + else: + _rewrite_doc( + path / name, + 3, + lambda doc: doc["chunk_grid"]["configuration"].update(chunk_shape=[0]), + ) + _rewrite_doc( + path, + 3, + lambda doc, name=name: doc["consolidated_metadata"]["metadata"][name]["chunk_grid"][ + "configuration" + ].update(chunk_shape=[0]), + ) + + for use_consolidated in (True, False): + with warnings.catch_warnings(record=True) as record: + warnings.simplefilter("always", ZarrUserWarning) + group = zarr.open_group(path, mode="r+", use_consolidated=use_consolidated) + arrays = [group[name] for name in names] + assert all("zarr.consolidate_metadata" in str(w.message) for w in record) + assert sorted(str(w.message).split(": ")[0] for w in record) == [ + f"Array '{group.store_path / name}'" for name in names + ] + for array in arrays: + assert isinstance(array, zarr.Array) + assert array.chunks == (2,) + array.update_attributes({}) + zarr.consolidate_metadata(path) + with warnings.catch_warnings(): + warnings.simplefilter("error", ZarrUserWarning) + warnings.filterwarnings("ignore", "Consolidated metadata is currently not part") + for use_consolidated in (True, False): + reopened = zarr.open_group(path, mode="r", use_consolidated=use_consolidated) + for name in names: + array = reopened[name] + assert isinstance(array, zarr.Array) + assert array.chunks == (2,) diff --git a/tests/test_unified_chunk_grid.py b/tests/test_unified_chunk_grid.py index b8289d2135..6a43cb4f3b 100644 --- a/tests/test_unified_chunk_grid.py +++ b/tests/test_unified_chunk_grid.py @@ -8,6 +8,7 @@ from __future__ import annotations +import re from typing import TYPE_CHECKING, Any import numpy as np @@ -149,9 +150,9 @@ def test_rectilinear_feature_flag_enabled() -> None: (10, 100, 1, 10, 10, 10, 10), (10, 100, 9, 10, 10, 10, 90), (10, 95, 9, 10, 10, 5, 90), # boundary chunk - (0, 0, None, 0, None, None, None), # zero-size + (10, 0, None, 0, None, None, None), # zero-extent: no chunks, size still >= 1 ], - ids=["start", "middle", "end", "boundary", "zero-size"], + ids=["start", "middle", "end", "boundary", "zero-extent"], ) def test_fixed_dimension( size: int, @@ -190,11 +191,16 @@ def test_fixed_dimension_indices_to_chunks() -> None: @pytest.mark.parametrize( ("size", "extent", "match"), - [(-1, 100, "must be >= 0"), (10, -1, "must be >= 0")], - ids=["negative-size", "negative-extent"], + [ + (-1, 100, "size must be >= 1"), + (0, 100, "size must be >= 1"), + (0, 0, "size must be >= 1"), + (10, -1, "extent must be >= 0"), + ], + ids=["negative-size", "zero-size", "zero-size-zero-extent", "negative-extent"], ) -def test_fixed_dimension_rejects_negative(size: int, extent: int, match: str) -> None: - """FixedDimension raises ValueError for negative size or extent""" +def test_fixed_dimension_rejects_invalid(size: int, extent: int, match: str) -> None: + """FixedDimension raises ValueError for a size below 1 or a negative extent.""" with pytest.raises(ValueError, match=match): FixedDimension(size=size, extent=extent) @@ -489,6 +495,7 @@ def test_chunk_grid_iter() -> None: [ ([[10, 3]], [10, 10, 10]), ([[10, 2], [20, 1]], [10, 10, 20]), + ([[True, 2], [3, 1]], [1, 1, 3]), ], ) def test_rle_expand(compressed: list[Any], expected: list[int]) -> None: @@ -542,19 +549,49 @@ def test_rle_expand_rejects_invalid(rle_input: list[Any], match: str) -> None: expand_rle(rle_input) -# -- expand_rle handles JSON floats -- - - -def test_expand_rle_bare_integer_floats_accepted() -> None: - """JSON parsers may emit 10.0 for the integer 10; expand_rle should handle it.""" - result = expand_rle([10.0, 20.0]) # type: ignore[list-item] - assert result == [10, 20] +@pytest.mark.parametrize( + ("rle_input", "match"), + [ + ([10.5], "Chunk edge length must be an int, got 10.5"), + ([10.0], "Chunk edge length must be an int, got 10.0"), + ([[10.0, 3]], "Chunk edge length must be an int, got 10.0"), + (["10"], "Chunk edge length must be an int, got '10'"), + ([[10, 3.5]], "RLE repeat count must be an int, got 3.5"), + ([[10, 3.0]], "RLE repeat count must be an int, got 3.0"), + ([np.int64(10)], "Chunk edge length must be an int, got np.int64(10)"), + ([[10, np.int64(3)]], "RLE repeat count must be an int, got np.int64(3)"), + ], + ids=[ + "fractional-edge", + "float-edge", + "float-rle-size", + "string-edge", + "fractional-count", + "float-count", + "numpy-int-edge", + "numpy-int-count", + ], +) +def test_rle_expand_rejects_non_int(rle_input: list[Any], match: str) -> None: + """expand_rle reads `int`s (and `bool`s) only, not floats or NumPy integers.""" + with pytest.raises(TypeError, match=re.escape(match)): + expand_rle(rle_input) -def test_expand_rle_pair_with_float_count() -> None: - """expand_rle accepts float repeat counts that are integer-valued""" - result = expand_rle([[10, 3.0]]) # type: ignore[list-item] - assert result == [10, 10, 10] +@pytest.mark.parametrize( + ("rle_input", "match"), + [ + ([0], "chunk edge length must be >= 1"), + ([10.5], "chunk edge length must be an int,"), + ([[5, 0]], "RLE repeat count must be >= 1"), + ([[5, 2, 1]], r"RLE entries must be an integer or \[size, count\]"), + ], + ids=["zero-edge", "fractional-edge", "zero-rle-count", "rle-entry-of-three"], +) +def test_rle_expand_names_dimension(rle_input: list[Any], match: str) -> None: + """Given the dimension `axis` the edges belong to, every error of expand_rle names it.""" + with pytest.raises((TypeError, ValueError), match=f"^Dimension 2: {match}"): + expand_rle(rle_input, axis=2) # --------------------------------------------------------------------------- @@ -1421,46 +1458,22 @@ def test_edge_case_chunk_grid_boundary_shape() -> None: # -- Zero-size and zero-extent -- -@pytest.mark.parametrize( - ("size", "extent"), - [(0, 0), (0, 5), (10, 0)], - ids=["zero-size-zero-extent", "zero-size-nonzero-extent", "zero-extent-nonzero-size"], -) -def test_edge_case_zero_size_or_extent(size: int, extent: int) -> None: - """FixedDimension with zero size or extent has zero chunks and getitem returns None""" - d = FixedDimension(size=size, extent=extent) +@pytest.mark.parametrize("size", [1, 10], ids=["size-1", "size-10"]) +def test_fixed_dimension_zero_extent(size: int) -> None: + """A zero-length axis has zero chunks and behaves like an empty grid.""" + d = FixedDimension(size=size, extent=0) assert d.nchunks == 0 + assert d.ngridcells == 0 + assert d.data_size(0) == 0 + assert d.with_extent(0) == d + assert d.with_extent(3) == FixedDimension(size=size, extent=3) + empty = np.array([], dtype=np.intp) + np.testing.assert_array_equal(d.indices_to_chunks(empty), empty) + with pytest.raises(IndexError): + d.index_to_chunk(0) g = ChunkGrid(dimensions=(d,)) assert g[0] is None - - -def test_edge_case_zero_size_data_and_indices() -> None: - """FixedDimension(size=0) handles data_size, index_to_chunk, and indices_to_chunks safely.""" - d = FixedDimension(size=0, extent=0) - # Zero-sized chunks have zero data - assert d.data_size(0) == 0 - # Vectorized lookup maps every index to chunk 0 (avoids division by zero) - indices = np.array([0, 0, 0], dtype=np.intp) - np.testing.assert_array_equal(d.indices_to_chunks(indices), np.zeros(3, dtype=np.intp)) - - -def test_edge_case_zero_size_nonzero_extent_index() -> None: - """FixedDimension(size=0, extent>0) maps valid indices to chunk 0 without dividing by zero.""" - d = FixedDimension(size=0, extent=5) - assert d.nchunks == 0 - # index_to_chunk avoids division by zero and returns 0 - assert d.index_to_chunk(0) == 0 - assert d.index_to_chunk(4) == 0 - - -def test_edge_case_zero_size_data_and_index() -> None: - """FixedDimension(size=0) returns zero for data_size and maps indices to chunk 0.""" - d = FixedDimension(size=0, extent=0) - # data_size returns 0 for a zero-sized chunk - assert d.data_size(0) == 0 - # vectorized indices_to_chunks returns zeros - indices = np.array([0, 0, 0], dtype=np.intp) - np.testing.assert_array_equal(d.indices_to_chunks(indices), np.zeros(3, dtype=np.intp)) + assert list(g) == [] # -- 0-d grid --