Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
31 commits
Select commit Hold shift + click to select a range
7154260
fix: ensure `ZipStore` is open before acquiring lock
oelhammouchi Nov 20, 2025
64940a1
Merge branch 'main' into fix-zip-lock
d-v-b Nov 21, 2025
2aa6e14
Merge branch 'main' into fix-zip-lock
d-v-b Nov 27, 2025
d59a4ae
Merge branch 'zarr-developers:main' into fix-zip-lock
oelhammouchi Dec 25, 2025
35f8019
refactor: Cleanup open check in ZipStore
Dec 25, 2025
05970a2
test: Add test to check that lock is present when required
Dec 25, 2025
864e445
refactor: Add return type to `_sync_ensure_open`
Dec 25, 2025
4487ab3
Merge branch 'main' into fix-zip-lock
d-v-b Jan 8, 2026
f85b29c
docs: Add release note
Jan 13, 2026
d77583a
Merge branch 'main' into fix-zip-lock
d-v-b Jan 13, 2026
8688953
Merge branch 'main' into fix-zip-lock
d-v-b Mar 11, 2026
6422ee7
Merge remote-tracking branch 'upstream/main' into fix/zipstore-open-l…
d-v-b Sep 29, 2026
adbed18
fix(storage): open ZipStore archive through one accessor and reopen i…
d-v-b Sep 29, 2026
119e7a9
docs: rename changelog fragment to PR number
d-v-b Sep 29, 2026
d24126c
test(storage): make ZipStore reopen tests exercise the reopen
d-v-b Sep 29, 2026
4f166dc
fix(storage): hold the ZipStore lock for all of close()
d-v-b Sep 29, 2026
fdf16c7
fix(storage): hold the ZipStore lock across move() and every open
d-v-b Sep 29, 2026
62cb610
docs: note move() in the ZipStore changelog fragment
d-v-b Sep 29, 2026
edf2a0c
fix(storage): accept mode "x" in ZipStoreAccessModeLiteral
d-v-b Sep 29, 2026
698c2c0
fix(storage): refuse to reopen a ZipStore on a write-only file object
d-v-b Sep 29, 2026
45cd41f
fix(storage): make ZipStore close() exception-safe and check first us…
d-v-b Sep 29, 2026
ee6033f
fix(storage): refuse to reopen a ZipStore whose close() failed
d-v-b Sep 29, 2026
c93bfda
fix(storage): leave a ZipStore closed when clear() fails partway
d-v-b Sep 29, 2026
c3c3420
fix(storage): let clear() recover a ZipStore whose close() failed
d-v-b Sep 29, 2026
9cb2c53
test(storage): close ZipStores over file objects when a reuse test fails
d-v-b Sep 29, 2026
a93ee73
fix(storage): open unpickled ZipStores lazily and refuse to reopen an…
d-v-b Sep 29, 2026
badf49d
docs: describe the unpickling fix from main's side in the ZipStore fr…
d-v-b Sep 29, 2026
19f118d
fix(storage): let ZipStore.clear() replace a damaged or missing archive
d-v-b Sep 29, 2026
a0d2b42
fix(storage): narrow ZipStore.clear()'s first open to mode "x" and re…
d-v-b Sep 29, 2026
0d25f13
refactor(storage): drop a redundant condition from ZipStore.clear()
d-v-b Sep 29, 2026
5b340cd
Merge branch 'main' into fix/zipstore-open-lifecycle
d-v-b Sep 29, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions changes/4450.bugfix.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
`ZipStore` now opens its archive on first use from every method, so calling `get`, `get_partial_values`, `set_if_not_exists`, or `clear` on a store that was never opened no longer raises `AttributeError`. A store created with `mode="w"` or `mode="x"` now keeps its entries when it is used again after `close()`, moved with `move()`, or unpickled, instead of truncating the archive or refusing to open it. Previously, even a read on a closed `mode="w"` store erased the archive, and such a store backed by a file object dropped its earlier entries when written to again after `close()`. A store backed by a readable, seekable file object now keeps them, and one backed by any other file object raises `io.UnsupportedOperation` instead, because it cannot read the archive back. Threads that use a store for the first time at the same moment no longer open the archive twice and lose writes, and neither does a thread that uses a store while another thread closes or moves it. If closing the archive fails, for example because the disk is full, every later use of the store that would reopen the archive now raises `RuntimeError`, instead of reopening an archive that may lack its central directory and silently starting a new one; for a writable store backed by a path, `clear()` replaces the archive and makes the store usable again. A `clear()` that fails partway now leaves the store closed instead of open on a closed archive. An unpickled store now opens its archive on first use, so a copy that is never used, such as those dask makes while building a graph, no longer touches the file; before, such a copy reopened the archive and corrupted it, so writing a dask array to a `ZipStore` produced an unreadable file. A store that reopens its archive and finds it is not a zip, as happens while another copy is still in its first writing session, raises `zipfile.BadZipFile` instead of starting a new archive; other concurrent writers are not detected. On a store that was never opened, `move()` now applies the first-open mode at the original path, so mode `"x"` refuses an existing file instead of moving it. The type of the `mode` parameter now includes `"x"`, which the store already accepted at runtime.
176 changes: 134 additions & 42 deletions src/zarr/storage/_zip.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@
if TYPE_CHECKING:
from collections.abc import AsyncIterator, Iterable

ZipStoreAccessModeLiteral = Literal["r", "w", "a"]
ZipStoreAccessModeLiteral = Literal["r", "w", "a", "x"]


class _RawReaderAdapter(io.RawIOBase):
Expand Down Expand Up @@ -72,6 +72,17 @@ def readinto(self, b: Any) -> int:
return n


def _is_unfinished_archive(path: Path) -> bool:
"""Whether the file at `path` exists but is not a zip archive."""
try:
with path.open("rb") as f:
# is_zipfile swallows OSError, so open the file here to let a
# directory or a permission error surface as itself
return not zipfile.is_zipfile(f)
except FileNotFoundError:
return False


class ZipStore(Store):
"""
Store using a ZIP file.
Expand All @@ -85,10 +96,20 @@ class ZipStore(Store):
can only be used for reading (`mode="r"`). The file object must stay
open for the lifetime of the store, and operations that require a
filesystem location (`clear`, `move`, pickling) are not supported.
Using the store again after `close()` reopens the archive, which
requires a file object that is readable and seekable; otherwise it
raises `io.UnsupportedOperation`.
mode : str, optional
One of 'r' to read an existing file, 'w' to truncate and write a new
file, 'a' to append to an existing file, or 'x' to exclusively create
and write a new file.
and write a new file. 'w' and 'x' apply to the first open only; the
store reopens its archive with 'a' after `close()`, `move()`, or
unpickling, so the entries it already wrote are kept. An unpickled
store opens its archive on first use. If `close()` raises, the
archive may be incomplete, and every later use that would reopen it
raises `RuntimeError` instead; for a writable store backed
by a path, `clear()` replaces the archive and makes the store usable
again.
compression : int, optional
Compression method to use when writing to the archive.
allowZip64 : bool, optional
Expand Down Expand Up @@ -162,24 +183,78 @@ def __init__(
self._zmode = mode
self.compression = compression
self.allowZip64 = allowZip64
self._lock = threading.RLock()
self._was_opened = False
self._close_failed = False

def _sync_open(self) -> None:
if self._is_open:
raise ValueError("store is already open")

self._lock = threading.RLock()
if self._close_failed:
# the central directory may be missing; appending to such a file
# makes zipfile start a new archive and drop the earlier entries
raise RuntimeError(
f"closing the archive of {self!r} failed, so it may be incomplete; "
"the store will not reopen it"
)
if (
self.path is not None
and self._was_opened
and self._zmode == "a"
and _is_unfinished_archive(self.path)
):
# the file lacks a central directory, e.g. because another copy of
# this store still has it open for writing; appending would start
# a new archive and drop the entries already written
raise zipfile.BadZipFile(
f"{self!r} cannot reopen {self.path}: it is not a zip archive. It "
"may be incomplete, for example because another store still has it "
"open for writing."
)
if (
self.path is None
and self._fileobj is not None
and self._was_opened
and not (self._fileobj.readable() and self._fileobj.seekable())
):
# reopening appends, which needs to read the archive back; on a
# write-only file zipfile would start a new archive and drop the
# earlier entries
raise io.UnsupportedOperation(
"a ZipStore backed by a file object that is not readable and "
"seekable cannot be used again after close(), because the "
"archive it wrote cannot be read back"
)

self._zf = zipfile.ZipFile(
self.path if self.path is not None else self._fileobj, # type: ignore[arg-type]
mode=self._zmode,
compression=self.compression,
allowZip64=self.allowZip64,
)
# "w" truncates and "x" refuses an existing file. Both apply only to the
# first open: reopening after close(), move(), or unpickling must keep
# the entries already written.
if self._zmode in ("w", "x"):
self._zmode = "a"

self._was_opened = True
self._is_open = True

def _zipfile(self) -> zipfile.ZipFile:
"""Return the archive, opening it on first use."""
with self._lock:
if not self._is_open:
self._sync_open()
return self._zf

async def _open(self) -> None:
self._sync_open()
with self._lock:
self._sync_open()

async def _ensure_open(self) -> None:
# the base class checks _is_open outside the lock
self._zipfile()

def __getstate__(self) -> dict[str, Any]:
if self.path is None:
Expand All @@ -189,24 +264,39 @@ def __getstate__(self) -> dict[str, Any]:
"cannot pickle a ZipStore backed by a file-like object; "
"construct the store from a path instead"
)
# We need a copy to not modify the state of the original store
state = self.__dict__.copy()
# We need a copy to not modify the state of the original store. The
# lock keeps it from catching _sync_open between opening with "w" and
# switching to "a", which would make unpickling truncate the archive.
with self._lock:
state = self.__dict__.copy()
for attr in ["_zf", "_lock"]:
state.pop(attr, None)
return state

def __setstate__(self, state: dict[str, Any]) -> None:
self.__dict__ = state
self.__dict__.setdefault("_close_failed", False)
self.__dict__.setdefault("_was_opened", False)
self._lock = threading.RLock()
# open on first use: a copy that is never used, like those dask makes
# while building a graph, must not open the file, because zipfile
# writes a central directory when it closes an archive it could not read
self._is_open = False
self._sync_open()

def close(self) -> None:
# docstring inherited
if not self._is_open:
return
super().close()
# hold the lock until the archive is closed: a thread that reopened it
# before its central directory was written would lose the entries
with self._lock:
self._zf.close()
if not self._is_open:
return
try:
self._zf.close()
except BaseException:
self._close_failed = True
raise
finally:
super().close()

async def clear(self) -> None:
# docstring inherited
Expand All @@ -216,11 +306,20 @@ async def clear(self) -> None:
raise NotImplementedError(
"clear() is not supported for a ZipStore backed by a file-like object"
)
self._zf.close()
os.remove(self.path)
self._zf = zipfile.ZipFile(
self.path, mode="w", compression=self.compression, allowZip64=self.allowZip64
)
if self._zmode == "x":
# "x" lasts only until the first open; opening now keeps it from
# deleting a file this store never claimed
self._zipfile()
self.close()
# the file may be damaged or gone; clear() replaces it either way
self.path.unlink(missing_ok=True)
# replacing the file cannot drop entries, so clear() is the one way
# to recover a store whose close() failed
self._close_failed = False
# if this open fails the store stays closed, and the next use
# creates the archive again
self._zmode = "w"
self._sync_open()

def __str__(self) -> str:
if self.path is None:
Expand All @@ -243,11 +342,9 @@ def _get(
prototype: BufferPrototype,
byte_range: ByteRequest | None = None,
) -> Buffer | None:
if not self._is_open:
self._sync_open()
# docstring inherited
try:
with self._zf.open(key) as f: # will raise KeyError
with self._zipfile().open(key) as f: # will raise KeyError
if byte_range is None:
return prototype.buffer.from_bytes(f.read())
elif isinstance(byte_range, RangeByteRequest):
Expand Down Expand Up @@ -288,8 +385,6 @@ async def get_partial_values(
return out

def _set(self, key: str, value: Buffer) -> None:
if not self._is_open:
self._sync_open()
# generally, this should be called inside a lock
keyinfo = zipfile.ZipInfo(filename=key, date_time=time.localtime(time.time())[:6])
keyinfo.compress_type = self.compression
Expand All @@ -298,13 +393,11 @@ def _set(self, key: str, value: Buffer) -> None:
keyinfo.external_attr |= 0x10 # MS-DOS directory flag
else:
keyinfo.external_attr = 0o644 << 16 # ?rw-r--r--
self._zf.writestr(keyinfo, value.to_bytes())
self._zipfile().writestr(keyinfo, value.to_bytes())

async def set(self, key: str, value: Buffer) -> None:
# docstring inherited
self._check_writable()
if not self._is_open:
self._sync_open()
if not isinstance(value, Buffer):
raise TypeError(
f"ZipStore.set(): `value` must be a Buffer instance. Got an instance of {type(value)} instead."
Expand All @@ -315,7 +408,7 @@ async def set(self, key: str, value: Buffer) -> None:
async def set_if_not_exists(self, key: str, value: Buffer) -> None:
self._check_writable()
with self._lock:
members = self._zf.namelist()
members = self._zipfile().namelist()
if key not in members:
self._set(key, value)

Expand All @@ -337,23 +430,20 @@ async def delete(self, key: str) -> None:

async def exists(self, key: str) -> bool:
# docstring inherited
if not self._is_open:
self._sync_open()
with self._lock:
try:
self._zf.getinfo(key)
self._zipfile().getinfo(key)
except KeyError:
return False
else:
return True

async def list(self) -> AsyncIterator[str]:
# docstring inherited
if not self._is_open:
self._sync_open()
with self._lock:
for key in self._zf.namelist():
yield key
# namelist() is a copy; holding the lock across yield would block
# other threads for as long as the caller iterates
for key in self._zipfile().namelist():
yield key

async def list_prefix(self, prefix: str) -> AsyncIterator[str]:
# docstring inherited
Expand All @@ -363,11 +453,9 @@ async def list_prefix(self, prefix: str) -> AsyncIterator[str]:

async def list_dir(self, prefix: str) -> AsyncIterator[str]:
# docstring inherited
if not self._is_open:
self._sync_open()
prefix = prefix.rstrip("/")

keys = self._zf.namelist()
keys = self._zipfile().namelist()
seen = set()
if prefix == "":
keys_unique = {k.split("/")[0] for k in keys}
Expand All @@ -393,8 +481,12 @@ async def move(self, path: Path | str) -> None:
)
if isinstance(path, str):
path = Path(path)
self.close()
os.makedirs(path.parent, exist_ok=True)
shutil.move(self.path, path)
self.path = path
await self._open()
# hold the lock so that no thread reopens the old path mid-move
with self._lock:
# opening first keeps mode "x" from moving a file it may not claim
self._zipfile()
self.close()
os.makedirs(path.parent, exist_ok=True)
shutil.move(self.path, path)
self.path = path
self._sync_open()
Loading
Loading