Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
70 changes: 13 additions & 57 deletions src/zarr/storage/_fsspec.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,6 @@
import json
import warnings
from contextlib import suppress
from logging import getLogger
from typing import TYPE_CHECKING, Any

from packaging.version import parse as parse_version
Expand All @@ -19,8 +18,6 @@
from zarr.errors import ZarrUserWarning
from zarr.storage._utils import _dereference_path

logger = getLogger(__name__)

if TYPE_CHECKING:
from collections.abc import AsyncIterator, Iterable

Expand All @@ -38,26 +35,6 @@
)


async def _close_fs(fs: AsyncFileSystem) -> None:
"""
Best-effort async close of an fsspec async filesystem owned by FsspecStore.

For filesystems that expose `set_session()` (e.g. s3fs) the underlying
aiohttp `ClientSession` is closed explicitly, which prevents
"Unclosed client session" `ResourceWarning`s from aiohttp. For all
other filesystem types the call is a no-op (not every implementation
manages an HTTP session directly).

Note that `set_session()` lazily creates a session if none exists yet, so
closing a store that never performed any I/O may instantiate a session
purely to close it. This is accepted best-effort behavior; fsspec does not
expose a stable, cross-implementation way to test for an existing session.
"""
if hasattr(fs, "set_session"):
session = await fs.set_session()
await session.close()


def _make_async(fs: AbstractFileSystem) -> AsyncFileSystem:
"""Convert a sync FSSpec filesystem to an async FFSpec filesystem

Expand Down Expand Up @@ -126,6 +103,15 @@ class FsspecStore(Store):
ZarrUserWarning
If the file system (fs) was not created with `asynchronous=True`.

Notes
-----
Closing the store does not close the underlying filesystem or its network
session. fsspec caches and shares filesystem instances across callers, so
the store cannot know whether it is the only user, and closing a shared
session would break other stores. The filesystem's lifecycle belongs to
whoever created it; use fsspec's own tools (e.g. `clear_instance_cache`)
to release it.

See Also
--------
FsspecStore.from_upath
Expand All @@ -152,9 +138,6 @@ def __init__(
self.fs = fs
self.path = path
self.allowed_exceptions = allowed_exceptions
# True only when this store created fs itself (from_url / from_mapper with new instance).
# Callers who supply their own fs remain responsible for its lifecycle.
self._owns_fs: bool = False

if not self.fs.async_impl:
raise TypeError("Filesystem needs to support async operations.")
Expand Down Expand Up @@ -220,17 +203,13 @@ def from_mapper(
-------
FsspecStore
"""
original_fs = fs_map.fs
fs = _make_async(original_fs)
store = cls(
fs = _make_async(fs_map.fs)
return cls(
fs=fs,
path=fs_map.root,
read_only=read_only,
allowed_exceptions=allowed_exceptions,
)
# _make_async returns a new instance when converting sync→async; own it.
store._owns_fs = fs is not original_fs
return store

@classmethod
def from_url(
Expand Down Expand Up @@ -272,39 +251,16 @@ def from_url(
if not fs.async_impl:
fs = _make_async(fs)

store = cls(fs=fs, path=path, read_only=read_only, allowed_exceptions=allowed_exceptions)
store._owns_fs = True
return store
return cls(fs=fs, path=path, read_only=read_only, allowed_exceptions=allowed_exceptions)

def with_read_only(self, read_only: bool = False) -> FsspecStore:
# docstring inherited
new_store = type(self)(
return type(self)(
fs=self.fs,
path=self.path,
allowed_exceptions=self.allowed_exceptions,
read_only=read_only,
)
# The derived store shares the same fs. Transfer ownership so the
# surviving store closes it, and clear ours to avoid a double-close.
# Otherwise the common `from_url(...).with_read_only()` chain would
# drop the only owner (the unreferenced source) and leak the session.
new_store._owns_fs = self._owns_fs
self._owns_fs = False
return new_store

def close(self) -> None:
# docstring inherited
if self._owns_fs:
from zarr.core.sync import sync as zarr_sync

# Best-effort: a failure to release the session must not block close(),
# but log it so a genuine regression in the close path stays observable
# rather than silently reverting to the leaking behavior.
try:
zarr_sync(_close_fs(self.fs))
except Exception:
logger.debug("Failed to close owned filesystem %r", self.fs, exc_info=True)
super().close()

async def clear(self) -> None:
# docstring inherited
Expand Down
134 changes: 34 additions & 100 deletions tests/test_store/test_fsspec.py
Original file line number Diff line number Diff line change
Expand Up @@ -276,75 +276,20 @@ async def test_delete_dir_unsupported_deletes(self, store: FsspecStore) -> None:
):
await store.delete_dir("test_prefix")

# ── Filesystem lifecycle (ownership) ──────────────────────────────────────
# ── Filesystem lifecycle ──────────────────────────────────────────────────

def test_from_url_owns_filesystem(self, endpoint_url: str) -> None:
"""FsspecStore.from_url() creates the async fs; it must own it."""
async def test_close_marks_store_closed(self, endpoint_url: str) -> None:
"""close() must succeed and mark the store not-open."""
store = FsspecStore.from_url(
f"s3://{test_bucket_name}/lifecycle/",
storage_options={"endpoint_url": endpoint_url, "anon": False},
)
assert store._owns_fs
store.close()

async def test_from_url_close_releases_store(self, endpoint_url: str) -> None:
"""
close() on a from_url() store must succeed without error and mark the
store as closed. For the owned filesystem, _close_fs() is invoked to
release the underlying S3 client / aiohttp connection pool.
"""
store = FsspecStore.from_url(
f"s3://{test_bucket_name}/lifecycle/",
storage_options={"endpoint_url": endpoint_url, "anon": False},
)
# Materialise the S3 client and connection pool.
await store.set("probe", cpu.Buffer.from_bytes(b"x"))

store.close()

assert not store._is_open

def test_direct_construction_does_not_own_filesystem(self, endpoint_url: str) -> None:
"""Direct FsspecStore() must not claim ownership — the caller owns the fs."""
try:
from fsspec import url_to_fs
except ImportError:
from fsspec.core import url_to_fs
fs, path = url_to_fs(
f"s3://{test_bucket_name}", endpoint_url=endpoint_url, anon=False, asynchronous=True
)
store = FsspecStore(fs=fs, path=path)
assert not store._owns_fs

@pytest.mark.skipif(
parse_version(fsspec.__version__) < parse_version("2024.03.01"),
reason="Prior bug in from_upath",
)
def test_from_upath_does_not_own_filesystem(self, endpoint_url: str) -> None:
"""from_upath() uses the UPath's existing fs; the store must not own it."""
upath = pytest.importorskip("upath")
path = upath.UPath(
f"s3://{test_bucket_name}/foo/bar/",
endpoint_url=endpoint_url,
anon=False,
asynchronous=True,
)
store = FsspecStore.from_upath(path)
assert not store._owns_fs

def test_from_mapper_does_not_own_already_async_filesystem(self, endpoint_url: str) -> None:
"""from_mapper() with an already-async fs must not claim ownership."""
s3_filesystem = s3fs.S3FileSystem(
asynchronous=True,
endpoint_url=endpoint_url,
anon=False,
skip_instance_cache=True,
)
mapper = s3_filesystem.get_mapper(f"s3://{test_bucket_name}/")
store = FsspecStore.from_mapper(mapper)
# _make_async returns the same instance for an already-async fs.
assert not store._owns_fs


def array_roundtrip(store: FsspecStore) -> None:
"""
Expand Down Expand Up @@ -574,80 +519,69 @@ def test_open_s3map_raises(endpoint_url: str) -> None:
zarr.open(store=mapper, storage_options={"anon": True}, mode="w", shape=(3, 3))


async def test_close_fs_closes_s3_client() -> None:
"""
_close_fs() must call set_session() and then close() on the returned
S3 client. This is verified with mocks to avoid a real S3 connection.
"""
from unittest.mock import AsyncMock
async def test_close_does_not_close_filesystem_session() -> None:
"""close() must not touch the filesystem's session.

from zarr.storage._fsspec import _close_fs
fsspec caches and shares filesystem instances across callers, so the
session is not the store's to close. HTTP is used because its aiohttp
session is observably closed for good; s3fs transparently reconnects, which
would hide a regression. No request is issued — set_session() only
constructs the session.
"""
pytest.importorskip("aiohttp")
store = FsspecStore.from_url("http://example.com/a")
session = await store.fs.set_session()

mock_client = AsyncMock()
mock_fs = AsyncMock()
mock_fs.set_session = AsyncMock(return_value=mock_client)
store.close()

await _close_fs(mock_fs)
assert not session.closed

mock_fs.set_session.assert_called_once()
mock_client.close.assert_called_once()

async def test_close_does_not_break_a_sibling_store() -> None:
"""Closing one store must not close a session another store is using.

async def test_close_fs_no_op_for_fs_without_set_session() -> None:
"""_close_fs() must be a no-op for filesystems that don't expose set_session()."""
from unittest.mock import AsyncMock
Two stores from different URLs on one host are handed the same cached
filesystem; a store that closed it on close() would take the sibling's
session down too. This is the regression guard for that bug.
"""
pytest.importorskip("aiohttp")
s1 = FsspecStore.from_url("http://example.com/a")
s2 = FsspecStore.from_url("http://example.com/b")
session = await s2.fs.set_session()

from zarr.storage._fsspec import _close_fs
s1.close()

mock_fs = AsyncMock(spec=[]) # empty spec — no set_session attribute
await _close_fs(mock_fs) # must not raise
assert not session.closed


@pytest.mark.skipif(
parse_version(fsspec.__version__) < parse_version("2024.12.0"),
reason="No AsyncFileSystemWrapper",
)
def test_from_mapper_owns_wrapped_sync_filesystem(tmp_path: pathlib.Path) -> None:
"""
from_mapper() with a sync fs must wrap it in AsyncFileSystemWrapper and
claim ownership so that close() cleans it up.

The local filesystem is synchronous; _make_async() produces a new
AsyncFileSystemWrapper instance — a different object from the original fs.
"""
def test_from_mapper_wraps_sync_filesystem(tmp_path: pathlib.Path) -> None:
"""from_mapper() with a sync fs wraps it in an AsyncFileSystemWrapper."""
import fsspec as _fsspec
from fsspec.implementations.asyn_wrapper import AsyncFileSystemWrapper

fs = _fsspec.filesystem("file", auto_mkdir=True)
mapper = fs.get_mapper(str(tmp_path))
store = FsspecStore.from_mapper(mapper)
assert isinstance(store.fs, AsyncFileSystemWrapper)
assert store._owns_fs


@pytest.mark.skipif(
parse_version(fsspec.__version__) < parse_version("2024.12.0"),
reason="No AsyncFileSystemWrapper",
)
def test_with_read_only_transfers_filesystem_ownership(tmp_path: pathlib.Path) -> None:
"""
with_read_only() must transfer fs ownership to the derived store and clear
it on the source, so the surviving store closes the shared fs exactly once.

In the common ``from_url(...).with_read_only()`` chain the source store is
immediately unreferenced; if ownership were not transferred, the only owner
would be garbage-collected without close() and the session would leak.
"""
def test_with_read_only_shares_filesystem(tmp_path: pathlib.Path) -> None:
"""with_read_only() returns a store sharing the source's filesystem."""
source = FsspecStore.from_url(f"file://{tmp_path}", storage_options={"auto_mkdir": False})
assert source._owns_fs

derived = source.with_read_only(read_only=True)

# Ownership moved to the survivor; the source no longer owns it (no double-close).
assert derived._owns_fs
assert not source._owns_fs
# The derived store shares the same underlying fs.
assert derived.fs is source.fs
assert derived.read_only
assert not source.read_only


@pytest.mark.parametrize("asynchronous", [True, False])
Expand Down
Loading