diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 6eacae1..b04a413 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -33,6 +33,21 @@ Requires FUSE system libraries (`apt-get install fuse` on Debian/Ubuntu): tox -e py3 -- --libfuse ``` +### Backend tests + +Remote files are read by a backend (`fsspec`, always installed, and the +optional `remfile`). `remfile` comes from the `full` extra, which tox +installs only in envs whose name carries `full` or `libfuse`; tests that +need it are skipped elsewhere via the `requires_remfile` marker. + +```bash +# with remfile (and --libfuse) +tox -e py314-full + +# only the backend tests, in whatever the current env has installed +tox -e py3 -- datalad_fuse/tests/test_backends.py +``` + ### Forgejo-aneksajo integration tests These tests start an ephemeral [Forgejo-aneksajo](https://codeberg.org/forgejo-aneksajo/forgejo-aneksajo) diff --git a/README.md b/README.md index aed11f5..627bffd 100644 --- a/README.md +++ b/README.md @@ -5,12 +5,13 @@ `datalad-fuse` lets you read files of [DataLad](https://www.datalad.org) datasets and [git-annex](https://git-annex.branchable.com) repositories without downloading them first: only the parts of files that are actually -read are fetched, via [fsspec](https://filesystem-spec.readthedocs.io), from -the URLs that git-annex knows. Unlike `datalad get`, which downloads whole -files before you can use them, this pays off when you need only parts of -large files, e.g. a few arrays from NWB/HDF5 files or the headers of many -images. Use it through a FUSE mount, so that any program can open the files, -or directly from Python. +read are fetched, from the URLs that git-annex knows, via +[fsspec](https://filesystem-spec.readthedocs.io) or — for NWB/HDF5 files, if +installed — [remfile](https://github.com/flatironinstitute/remfile). Unlike +`datalad get`, which downloads whole files before you can use them, this pays +off when you need only parts of large files, e.g. a few arrays from NWB/HDF5 +files or the headers of many images. Use it through a FUSE mount, so that +any program can open the files, or directly from Python. **Documentation: https://datalad-fuse.readthedocs.io** @@ -18,6 +19,9 @@ or directly from Python. python3 -m pip install datalad-fuse +Add the `remfile` backend, which is faster for NWB/HDF5 files, with +`python3 -m pip install "datalad-fuse[remfile]"`. + [git-annex](https://git-annex.branchable.com/install/) is required, and FUSE for mounting datasets (e.g. `sudo apt-get install fuse3 libfuse2t64` on Ubuntu 24.04). See the @@ -52,7 +56,7 @@ from contextlib import closing import h5py import pynwb -from datalad_fuse.fsspec import DatasetAdapter +from datalad_fuse.adapter import DatasetAdapter # caching=False: keep fetched data in memory only with closing(DatasetAdapter("000582", caching=False)) as dsa: diff --git a/datalad_fuse/__init__.py b/datalad_fuse/__init__.py index dace56d..d038856 100644 --- a/datalad_fuse/__init__.py +++ b/datalad_fuse/__init__.py @@ -14,7 +14,7 @@ ) from datalad.interface.base import Interface, build_doc, eval_results from datalad.interface.results import get_status_dict -from datalad.support.constraints import EnsureNone +from datalad.support.constraints import EnsureNone, EnsureStr from datalad.support.param import Parameter from ._version import get_versions @@ -94,7 +94,17 @@ class FuseFS(Interface): args=("--caching",), choices=["none", "ondisk"], default="none", - doc="Whether to cache fsspec'ed files on disk on not at all", + doc="Whether to cache data fetched from remote URLs on disk or not at all", + ), + "backends": Parameter( + args=("--backends",), + doc=( + "Comma-separated list of backends to try for remote file" + " access, in priority order. Available: remfile, fsspec." + " Default: remfile,fsspec (remfile for HDF5 files," + " fsspec for everything else)" + ), + constraints=EnsureStr() | EnsureNone(), ), # TODO: (might better become config vars?) # --cache=persist @@ -112,6 +122,7 @@ def __call__( mode_transparent: bool = False, allow_other: bool = False, caching: str | None = None, + backends: str | None = None, ) -> Iterator[Dict[str, Any]]: from fuse import FUSE @@ -133,6 +144,7 @@ def __call__( ds.path, mode_transparent=mode_transparent, caching=caching == "ondisk", + backends=backends, ), mount_path, foreground=foreground, diff --git a/datalad_fuse/adapter.py b/datalad_fuse/adapter.py new file mode 100644 index 0000000..df3d6fa --- /dev/null +++ b/datalad_fuse/adapter.py @@ -0,0 +1,800 @@ +"""Backend-agnostic adapter layer for remote file access.""" + +from __future__ import annotations + +from collections.abc import Iterator +from datetime import datetime, timezone +from enum import Enum +from itertools import chain +import json +import logging +import os +import os.path +from pathlib import Path +import subprocess +from types import TracebackType +from typing import IO, Any, Callable, Optional, Tuple, cast +from urllib.parse import urlparse +import urllib.request + +import boto3 +from botocore import UNSIGNED +from botocore.config import Config as BotocoreConfig +from datalad import cfg +from datalad.distribution.dataset import Dataset +from datalad.support.annexrepo import AnnexRepo +from datalad.utils import get_dataset_root +import methodtools + +from .backends import DEFAULT_BACKENDS, Backend +from .consts import CACHE_SIZE +from .fsspec import FsspecBackend +from .remfile import RemfileBackend +from .utils import AnnexKey, is_annex_dir_or_key + +lgr = logging.getLogger("datalad.fuse.adapter") + + +class FileState(Enum): + """State of a file in a dataset, as returned by ``get_file_state()``""" + + #: The file is not annexed (e.g. committed to git directly); it is read + #: from disk. + NOT_ANNEXED = 1 + #: The file is annexed but its content is not present locally; it is read + #: from a remote URL. + NO_CONTENT = 2 + #: The file is annexed and its content is present locally; it is read from + #: disk. + HAS_CONTENT = 3 + + +# --------------------------------------------------------------------------- +# Backend creation helpers +# --------------------------------------------------------------------------- + + +def resolve_backends( + backends: Optional[str] = None, config: Optional[Any] = None +) -> tuple[str, bool]: + """Resolve backends spec from *backends* argument, config, or default. + + *config* is a dataset's :class:`~datalad.config.ConfigManager`, so that a + per-dataset ``datalad.fusefs.backends`` in ``.git/config`` or + ``.datalad/config`` is honored; it inherits the global overrides (``-c``), + so it is a superset of the global ``datalad.cfg``. Falls back to the + global config when not given. + + Returns ``(spec, explicit)`` where ``explicit`` is True when the user (or + config) supplied the spec and False when falling back to + :data:`~datalad_fuse.backends.DEFAULT_BACKENDS`. Callers use + ``explicit`` to decide whether a missing backend is a warning (explicit) + or a silent skip (default). + """ + if backends is not None: + return backends, True + from_cfg = (config if config is not None else cfg).get( + "datalad.fusefs.backends", None + ) + if from_cfg is not None: + return str(from_cfg), True + return DEFAULT_BACKENDS, False + + +def create_backends( + spec: str, path: str | Path, caching: bool, explicit: bool = True +) -> list[Backend]: + """Instantiate backends from a comma-separated *spec*. + + Backends that cannot be imported are skipped; missing backends requested + via an *explicit* spec (user argument or config) are logged as warnings, + while those missing from the default spec are logged at debug level only. + Raises ``ValueError`` if no usable backend remains. + """ + backends: list[Backend] = [] + for name in spec.split(","): + name = name.strip() + if not name: + # tolerate stray commas / whitespace, e.g. "remfile,,fsspec" + continue + try: + if name == "fsspec": + backends.append(FsspecBackend(path, caching)) + elif name == "remfile": + backends.append(RemfileBackend(path, caching)) + else: + raise ValueError(f"Unknown backend: {name!r}") + except ImportError as e: + if explicit: + lgr.warning( + "Backend %r requested but not available (%s); skipping. " + "Install its package or adjust the --backends spec.", + name, + e, + ) + else: + lgr.debug("Backend %r not available (not installed), skipping", name) + if not backends: + raise ValueError( + f"No usable backends from spec {spec!r}. " + "Install missing packages or adjust --backends / " + "datalad.fusefs.backends config." + ) + return backends + + +def replayable(it: Iterator[str]) -> Callable[[], Iterator[str]]: + """Wrap *it* so it can be iterated repeatedly, consuming it only once. + + Each returned iterator replays what has already been pulled from *it* + before pulling anything new, so the underlying generator stays lazy. + """ + seen: list[str] = [] + + def replay() -> Iterator[str]: + yield from seen + for item in it: + seen.append(item) + yield item + + return replay + + +def is_http_url(s: str) -> bool: + return s.lower().startswith(("http://", "https://")) + + +_aneksajo_cache: dict[str, bool] = {} + + +def _is_aneksajo(base_url: str) -> bool: + """Check if a URL points to a Forgejo-aneksajo instance. + + Probes ``{scheme}://{host}/api/forgejo/v1/version`` and checks whether + the version string contains ``git-annex``, which indicates the + forgejo-aneksajo fork. + + Results are cached per ``scheme://host:port`` for the process lifetime. + """ + parsed = urlparse(base_url) + # Cache key without userinfo so credentials don't fragment the cache + host = parsed.hostname or "" + port_suffix = f":{parsed.port}" if parsed.port else "" # noqa: E231 + cache_key = f"{parsed.scheme}://{host}{port_suffix}" # noqa: E231 + + if cache_key in _aneksajo_cache: + return _aneksajo_cache[cache_key] + + try: + api_url = f"{cache_key}/api/forgejo/v1/version" + req = urllib.request.Request(api_url, method="GET") + req.add_header("Accept", "application/json") + with urllib.request.urlopen(req, timeout=10) as resp: + data = json.loads(resp.read().decode()) + result = "git-annex" in data.get("version", "") + except Exception: + lgr.debug("_is_aneksajo(%s) probe failed", cache_key, exc_info=True) + result = False + + _aneksajo_cache[cache_key] = result + lgr.debug("_is_aneksajo(%s) = %s", cache_key, result) + return result + + +# --------------------------------------------------------------------------- +# Dataset / Adapter layer +# --------------------------------------------------------------------------- + + +class DatasetAdapter: + """Read access to the files of a single dataset. + + Files that are not annexed, and annexed files whose content is present + locally, are opened from disk. Annexed files without local content are + opened from one of the http(s) URLs found for their git-annex key (see + :meth:`get_urls`), reading only the needed parts of the file. + + Parameters + ---------- + path : str or Path + Top directory of the dataset (any git or git-annex repository). + caching : bool + If true, keep the data fetched from remote URLs in a sparse on-disk + cache under ``/.git/datalad/cache/`` (one subdirectory per + backend), to be reused by subsequent reads (for up to a week with the + ``fsspec`` backend). If false, data are only buffered in memory while + a file is open. + mode_transparent : bool + If true, paths of key files under ``.git/annex/objects/`` (the targets + of annexed symlinks) are opened as annexed content, fetched from a + remote URL if not present locally. + backends : str, optional + Comma-separated, priority-ordered backends to try, e.g. + ``"remfile,fsspec"``. Defaults to the ``datalad.fusefs.backends`` + configuration option, or to + :data:`~datalad_fuse.backends.DEFAULT_BACKENDS`. + + Notes + ----- + Call :meth:`close` (or use :func:`contextlib.closing`) when done, to stop + the ``git annex`` processes started for the dataset. + """ + + def __init__( + self, + path: str | Path, + caching: bool, + mode_transparent: bool = False, + backends: Optional[str] = None, + ) -> None: + self.path = Path(path) + self.mode_transparent = mode_transparent + ds = Dataset(path) + self.annex: Optional[AnnexRepo] + if isinstance(ds.repo, AnnexRepo): + self.annex = ds.repo + else: + self.annex = None + self.commit_dt = datetime.fromtimestamp( + ds.repo.get_commit_date(), tz=timezone.utc + ) + spec, explicit = resolve_backends(backends, config=ds.config) + self._backends = create_backends(spec, path, caching, explicit=explicit) + + def close(self) -> None: + """Stop the batched ``git annex`` processes started for the dataset""" + if self.annex is not None: + self.annex._batched.clear() + + @methodtools.lru_cache(maxsize=CACHE_SIZE) + def get_file_state(self, relpath: str) -> tuple[FileState, Optional[AnnexKey]]: + """Determine whether a file is annexed and has its content present + + Results are cached (for the most recently queried files). + + Parameters + ---------- + relpath : str + Path of the file relative to the top directory of the dataset. + + Returns + ------- + tuple of (FileState, AnnexKey or None) + The state of the file, and its git-annex key if it is annexed. + """ + p = self.path / relpath + lgr.debug("get_file_state: %s", relpath) + + def handle_path_under_annex_objects( + p: Path, + ) -> tuple[FileState, Optional[AnnexKey]]: + iadok = is_annex_dir_or_key(p) + if isinstance(iadok, AnnexKey): + if p.exists(): + return (FileState.HAS_CONTENT, iadok) + else: + return (FileState.NO_CONTENT, iadok) + else: + return (FileState.NOT_ANNEXED, None) + + # Shortcut handling of content under .git, in particular - annex key paths + if self.mode_transparent and relpath.startswith(".git/"): + return handle_path_under_annex_objects(p) + + # A regular file or git link for which we need to explicitly ask annex about + if not p.is_symlink(): + if p.stat().st_size < 1024 and self.annex is not None: + if self.annex.is_under_annex(relpath, batch=True): + key = AnnexKey.parse(self.annex.get_file_key(relpath, batch=True)) + if self.annex.file_has_content(relpath, batch=True): + return (FileState.HAS_CONTENT, key) + else: + return (FileState.NO_CONTENT, key) + return (FileState.NOT_ANNEXED, None) + + return handle_path_under_annex_objects( + Path(os.path.normpath(p.parent / os.readlink(p))) + ) + + def get_urls(self, key: str) -> Iterator[str]: + """Yield candidate http(s) URLs for the content of an annex key + + URLs are yielded in the order in which they are tried by :meth:`open`: + + 1. http(s) URLs recorded in git-annex for the key, as reported by + ``git annex whereis`` (e.g. those of the ``web`` special remote); + 2. ``annex/objects/...`` locations on the http(s) git remotes that + ``git annex whereis`` lists as having the key, including the + ``annex/objects`` endpoint of Forgejo-aneksajo instances. + + URLs on S3 special remotes with ``exporttree=yes`` are not included; + :meth:`open` falls back to them via :meth:`get_exporttree_urls`. + + Parameters + ---------- + key : str + A git-annex key, e.g. ``str(AnnexKey)``. + """ + assert self.annex is not None + # TODO: switch to batch=True whenever + # https://github.com/datalad/datalad/pull/6379 is merged/released. + # Will need a recent git-annex to work! + whereis = self.annex.whereis(key, output="full", batch=False, key=True) + remote_uuids = [] + for ru, v in whereis.items(): + remote_uuids.append(ru) + for u in v["urls"]: + if is_http_url(u): + yield u + + path_mixed = self.annex._batched.get( + "examinekey", + annex_options=["--format=annex/objects/${hashdirmixed}${key}/${key}\\n"], + path=self.annex.path, + )(key) + path_lower = self.annex._batched.get( + "examinekey", + annex_options=["--format=annex/objects/${hashdirlower}${key}/${key}\\n"], + path=self.annex.path, + )(key) + + uuid2remote_url = {} + aneksajo_uuids: set[str] = set() + for r in self.annex.get_remotes(): + if (ru := self.annex.config.get(f"remote.{r}.annex-uuid")) is None: + continue + if (remote_url := self.annex.config.get(f"remote.{r}.url")) is None: + continue + remote_url = self.annex.config.rewrite_url(remote_url) + uuid2remote_url[ru] = remote_url + # Detect Forgejo-aneksajo instances via API probe (cached). + # TODO: pushurl could be different from url, should also check + # remote.{r}.pushurl config + # TODO: SSH remote URLs not yet supported -- would need to + # derive the HTTP base URL from the SSH URL + if is_http_url(remote_url) and _is_aneksajo(remote_url): + aneksajo_uuids.add(ru) + + for ru in remote_uuids: + try: + base_url = uuid2remote_url[ru] + except KeyError: + continue + if is_http_url(base_url): + base_stripped = base_url.rstrip("/") + # Forgejo/Gitea with aneksajo: use annex/objects endpoint + # which supports HEAD and Range requests. + # See https://codeberg.org/forgejo-aneksajo/forgejo-aneksajo/issues/111 + if ru in aneksajo_uuids and base_stripped.endswith(".git"): + forge_base = base_stripped[:-4].rstrip("/") + yield forge_base + "/" + path_lower + if base_stripped.lower().endswith("/.git"): + paths = [path_mixed, path_lower] + else: + paths = [ + path_lower, + path_mixed, + f".git/{path_lower}", + f".git/{path_mixed}", + ] + for p in paths: + yield base_stripped + "/" + p + + @methodtools.lru_cache(maxsize=1) + def _get_exporttree_remotes(self) -> list[dict[str, str]]: + """Get S3 exporttree remotes with public URLs. + + Parses the git-annex branch remote.log once (cached per + DatasetAdapter instance) to find S3 special remotes configured + with ``exporttree=yes`` and a usable ``publicurl``. + + Workaround for legacy datasets that lack proper versioned S3 URLs in + their git-annex metadata. See + https://github.com/OpenNeuroOrg/openneuro/issues/3875 + """ + try: + result = subprocess.run( + ["git", "-C", str(self.path), "show", "git-annex:remote.log"], + capture_output=True, + text=True, + check=True, + ) + except subprocess.CalledProcessError: + lgr.debug("Could not read git-annex:remote.log for %s", self.path) + return [] + + remotes: list[dict[str, str]] = [] + for line in result.stdout.strip().splitlines(): + if not line or line.startswith("#"): + continue + parts = line.split() + if len(parts) < 2: + continue + uuid = parts[0] + config: dict[str, str] = {} + for token in parts[1:]: + if "=" in token: + k, v = token.split("=", 1) + config[k] = v + if ( + config.get("type") == "S3" + and config.get("exporttree") == "yes" + and config.get("publicurl", "no").startswith("http") + ): + remotes.append( + { + "uuid": uuid, + "publicurl": config["publicurl"].rstrip("/"), + "fileprefix": config.get("fileprefix", ""), + "bucket": config.get("bucket", ""), + "host": config.get("host", "s3.amazonaws.com"), + } + ) + return remotes + + @staticmethod + def _list_s3_versions( + bucket: str, + object_key: str, + host: str = "s3.amazonaws.com", + ) -> list[dict[str, Any]]: + """List all S3 object versions for *object_key* in *bucket*. + + Uses anonymous credentials (for public buckets). + """ + try: + endpoint_url = f"https://{host}" # noqa: E231 + client = boto3.client( + "s3", + endpoint_url=endpoint_url, + config=BotocoreConfig(signature_version=UNSIGNED), + ) + response = client.list_object_versions(Bucket=bucket, Prefix=object_key) + except Exception as e: + lgr.debug("Failed to list S3 versions for %s/%s: %s", bucket, object_key, e) + return [] + + versions: list[dict[str, Any]] = [] + for v in response.get("Versions", []): + # Only include exact key matches (prefix query may return others) + if v.get("Key") == object_key: + versions.append( + { + "VersionId": v.get("VersionId", ""), + "Size": v.get("Size", 0), + "ETag": v.get("ETag", ""), + "IsLatest": v.get("IsLatest", False), + } + ) + return versions + + @staticmethod + def _match_s3_version( + versions: list[dict[str, Any]], expected_size: int + ) -> Optional[str]: + """Match the correct S3 object version by file size. + + Returns the matched versionId, or ``None`` if no version matches. + Raises ``ValueError`` if multiple versions match by size but have + different ETags (ambiguous content — refuse to guess). + """ + matches = [v for v in versions if v["Size"] == expected_size] + if not matches: + return None + if len(matches) == 1: + return str(matches[0]["VersionId"]) + # Multiple matches — check ETags + etags = {v["ETag"] for v in matches} + if len(etags) == 1: + # Same content uploaded multiple times; prefer the latest + for v in matches: + if v["IsLatest"]: + return str(v["VersionId"]) + return str(matches[0]["VersionId"]) + raise ValueError( + f"Ambiguous S3 versions: {len(matches)} versions match size " + f"{expected_size} but have {len(etags)} distinct ETags. " + f"Cannot determine correct version." + ) + + def get_exporttree_urls(self, relpath: str, key: AnnexKey) -> Iterator[str]: + """Yield versioned URLs for *relpath* on S3 exporttree remotes. + + Workaround for datasets lacking proper versioned URLs in git-annex + metadata. Constructs URLs from the remote's ``publicurl`` + + ``fileprefix`` and resolves the correct S3 object version by matching + ``key.size``. + """ + remotes = self._get_exporttree_remotes() + if not remotes: + return + + for remote in remotes: + publicurl = remote["publicurl"] + fileprefix = remote["fileprefix"] + bucket = remote["bucket"] + host = remote["host"] + object_key = f"{fileprefix}{relpath}" + base_url = f"{publicurl}/{object_key}" + + if key.size is not None: + versions = self._list_s3_versions(bucket, object_key, host) + if versions: + try: + version_id = self._match_s3_version(versions, key.size) + except ValueError as e: + lgr.warning("%s: %s", relpath, e) + continue + if version_id: + yield f"{base_url}?versionId={version_id}" + continue + else: + lgr.debug( + "%s: no S3 version matches size %d at %s", + relpath, + key.size, + base_url, + ) + continue + + # Fallback: no size info or version listing failed — + # try unversioned URL (returns latest version). + lgr.warning( + "%s: falling back to unversioned S3 URL %s " + "(cannot verify correct version)", + relpath, + base_url, + ) + yield base_url + + def open( + self, + relpath: str, + mode: str = "rb", + encoding: str = "utf-8", + errors: Optional[str] = None, + ) -> IO: + """Open a file of the dataset for reading + + Annexed files without local content are opened by the first configured + backend that can handle them (see + :meth:`~datalad_fuse.backends.Backend.can_handle`), trying each + candidate URL in turn; if none of them works, the next backend is + tried. + + Parameters + ---------- + relpath : str + Path of the file relative to the top directory of the dataset. + mode : str + ``"rb"`` (default) to get a binary file object, ``"r"`` or + ``"rt"`` to get a text file object. + encoding : str + Encoding to use in text mode. + errors : str, optional + How to handle encoding errors in text mode, as for :func:`open`. + + Returns + ------- + file object + A seekable, read-only file object. Files read from disk are + regular Python file objects; files read from a URL are whatever + the backend that opened them returns (an fsspec file object for + the ``fsspec`` backend, a + :class:`~datalad_fuse.remfile.RemfileWrapper` for ``remfile``). + + Raises + ------ + NotImplementedError + If ``mode`` is not one of the supported read modes. + IOError + If the content of an annexed file is not present locally and no + backend could open any of its candidate URLs. + """ + if mode not in ("r", "rb", "rt"): + raise NotImplementedError("Only modes 'r', 'rb', and 'rt' are supported") + if mode == "rb": + kwargs: dict[str, Any] = {} + else: + kwargs = {"encoding": encoding, "errors": errors} + fstate, key = self.get_file_state(relpath) + if fstate is FileState.NOT_ANNEXED: + lgr.debug("%s: not under annex", relpath) + else: + lgr.debug( + "%s: under annex, %s content", + relpath, + "has" if fstate is FileState.HAS_CONTENT else "does not have", + ) + if fstate is FileState.NO_CONTENT: + # Primary URLs from git-annex whereis / remote paths, then S3 + # exporttree fallback URLs (legacy openneuro datasets that lack + # proper versioned URLs in annex metadata). Lazy, so the whereis / + # examinekey / boto3 calls only happen when a backend actually asks + # for URLs, and replayable, so a backend falling through to the next + # one does not repeat them. + fallback_urls: Iterator[str] = ( + self.get_exporttree_urls(relpath, key) if key is not None else iter([]) + ) + urls = replayable(chain(self.get_urls(str(key)), fallback_urls)) + # Walk the backend chain; fall through to next backend on failure + last_error: Optional[Exception] = None + for backend in self._backends: + if not backend.can_handle(key, mode, relpath): + lgr.debug( + "%s: backend %s cannot handle (suffix=%s, mode=%s)", + relpath, + backend.name, + key.suffix if key else None, + mode, + ) + continue + lgr.debug("%s: opening via backend %s", relpath, backend.name) + for url in urls(): + try: + lgr.debug( + "%s: trying URL %s (backend=%s)", + relpath, + url, + backend.name, + ) + return backend.open_url(url, mode, **kwargs) + except FileNotFoundError as e: + lgr.debug( + "Failed to open %s at URL %s: %s", + relpath, + url, + str(e), + ) + last_error = e + except Exception as e: + lgr.debug( + "%s: backend %s failed at URL %s: %s", + relpath, + backend.name, + url, + e, + ) + last_error = e + # All URLs failed for this backend — try the next one + lgr.debug( + "%s: backend %s exhausted all URLs, trying next", + relpath, + backend.name, + ) + # No backend succeeded + raise IOError( + f"Could not open {relpath} within {self.path}" + f" (backends={','.join(b.name for b in self._backends)})" + ) from last_error + else: + lgr.debug("%s: opening directly", relpath) + return open(self.path / relpath, mode, **kwargs) # type: ignore[return-value] + + def clear(self) -> None: + """Remove the on-disk caches of the dataset (only if ``caching``)""" + for backend in self._backends: + backend.clear() + + +class RemoteFilesystemAdapter: + """Read access to the files of a dataset and its installed subdatasets. + + Each path is mapped to the (sub)dataset containing it, and a + :class:`DatasetAdapter` is created for that dataset on first use. Use it + as a context manager, so that the ``git annex`` processes started for the + datasets are stopped on exit. + + Parameters + ---------- + root : str or Path + Top directory of the (super)dataset. + caching : bool + Passed to each :class:`DatasetAdapter`. + mode_transparent : bool + Passed to each :class:`DatasetAdapter`. + backends : str, optional + Passed to each :class:`DatasetAdapter`. + + Notes + ----- + Use an absolute ``root``, and absolute paths under it for the methods. + Paths relative to ``root`` or to the current directory are not supported. + """ + + def __init__( + self, + root: str | Path, + caching: bool, + mode_transparent: bool = False, + backends: Optional[str] = None, + ) -> None: + self.root = Path(root) + self.mode_transparent = mode_transparent + self.caching = caching + self.backends = backends + self.datasets: dict[Path, DatasetAdapter] = {} + + def __enter__(self) -> RemoteFilesystemAdapter: + return self + + def __exit__( + self, + _exc_type: Optional[type[BaseException]], + _exc_val: Optional[BaseException], + _exc_tb: Optional[TracebackType], + ) -> None: + for ds in self.datasets.values(): + ds.close() + self.datasets.clear() + + @methodtools.lru_cache(maxsize=CACHE_SIZE) + # TODO: optimize "caching" more since for all files under the same directory + # they all would belong to the same dataset + def get_dataset_path(self, path: str | Path) -> Path: + """Return the top directory of the (sub)dataset containing ``path``""" + path = Path(self.root, path) + dspath = get_dataset_root(path) + if dspath is None: + raise ValueError(f"Path not under DataLad: {path}") + dspath = Path(dspath) + assert isinstance(dspath, Path) + try: + dspath.relative_to(self.root) + except ValueError: + raise ValueError(f"Path not under root dataset: {path}") + return dspath + + def resolve_dataset(self, filepath: str | Path) -> tuple[DatasetAdapter, str]: + """Return the adapter for the dataset containing ``filepath`` + + Returns + ------- + tuple of (DatasetAdapter, str) + The adapter, and the path of ``filepath`` relative to the + dataset's top directory. + """ + dspath = self.get_dataset_path(filepath) + try: + dsap = self.datasets[dspath] + except KeyError: + dsap = self.datasets[dspath] = DatasetAdapter( + dspath, + mode_transparent=self.mode_transparent, + caching=self.caching, + backends=self.backends, + ) + relpath = str(Path(filepath).relative_to(dspath)) + return dsap, relpath + + def open( + self, + filepath: str | Path, + mode: str = "rb", + encoding: str = "utf-8", + errors: Optional[str] = None, + ) -> IO: + """Open a file for reading; see :meth:`DatasetAdapter.open`""" + dsap, relpath = self.resolve_dataset(filepath) + lgr.debug( + "%s: path resolved to %s in dataset at %s", filepath, relpath, dsap.path + ) + return dsap.open(relpath, mode=mode, encoding=encoding, errors=errors) + + def get_file_state( + self, filepath: str | Path + ) -> tuple[FileState, Optional[AnnexKey]]: + """Return state and key of a file; see `DatasetAdapter.get_file_state`""" + dsap, relpath = self.resolve_dataset(filepath) + return cast(Tuple[FileState, Optional[AnnexKey]], dsap.get_file_state(relpath)) + + def is_under_annex(self, filepath: str | Path) -> bool: + """Tell whether a file is annexed""" + dsap, relpath = self.resolve_dataset(filepath) + fstate, _ = dsap.get_file_state(relpath) + return fstate is not FileState.NOT_ANNEXED + + def get_commit_datetime(self, filepath: str | Path) -> datetime: + """Return the date of ``HEAD`` in the dataset containing ``filepath``""" + dsap, _ = self.resolve_dataset(filepath) + return dsap.commit_dt diff --git a/datalad_fuse/backends.py b/datalad_fuse/backends.py new file mode 100644 index 0000000..4d50bd4 --- /dev/null +++ b/datalad_fuse/backends.py @@ -0,0 +1,35 @@ +"""Abstract base for remote file access backends and shared constants.""" + +from __future__ import annotations + +from abc import ABC, abstractmethod +from typing import IO, Any, Optional + +from .utils import AnnexKey + +#: Backends used when neither ``--backends`` nor the +#: ``datalad.fusefs.backends`` configuration option is set. +DEFAULT_BACKENDS = "remfile,fsspec" + + +class Backend(ABC): + """Base class for remote file access backends.""" + + name: str + + @abstractmethod + def can_handle( + self, key: Optional[AnnexKey], mode: str, relpath: Optional[str] = None + ) -> bool: + """Return True if this backend should be used for *key* in *mode*. + + *relpath* is the path within the dataset, for backends that dispatch on + the file name when the annex key carries no suffix (URL/VURL keys). + """ + + @abstractmethod + def open_url(self, url: str, mode: str = "rb", **kwargs: Any) -> IO: + """Open *url* and return a file-like object.""" + + def clear(self) -> None: # noqa: B027 + """Clear any caches held by this backend. Default: no-op.""" diff --git a/datalad_fuse/fsspec.py b/datalad_fuse/fsspec.py index 63b2b0c..b83a92f 100644 --- a/datalad_fuse/fsspec.py +++ b/datalad_fuse/fsspec.py @@ -1,739 +1,76 @@ +"""FsspecBackend — remote file access via fsspec's HTTPFileSystem. + +This module also provides backward-compatible imports for names that have +moved to :mod:`datalad_fuse.backends`, :mod:`datalad_fuse.remfile`, and +:mod:`datalad_fuse.adapter` as part of the multi-backend refactoring. +Importing those names from here still works but emits a +:class:`DeprecationWarning`. +""" + from __future__ import annotations -from collections.abc import Iterator -from datetime import datetime, timezone -from enum import Enum -import json +import importlib import logging import os -import os.path from pathlib import Path -import subprocess -from types import SimpleNamespace, TracebackType -from typing import IO, Any, Optional, Tuple, cast -from urllib.parse import urlparse -import urllib.request +from types import SimpleNamespace +from typing import IO, Any, Optional +import warnings import aiohttp from aiohttp_retry import ListRetry, RetryClient -import boto3 -from botocore import UNSIGNED -from botocore.config import Config as BotocoreConfig -from datalad.distribution.dataset import Dataset -from datalad.support.annexrepo import AnnexRepo -from datalad.utils import get_dataset_root from fsspec.exceptions import BlocksizeMismatchError from fsspec.implementations.cached import CachingFileSystem from fsspec.implementations.http import HTTPFileSystem -import methodtools -from .consts import CACHE_SIZE -from .utils import AnnexKey, is_annex_dir_or_key +from .backends import Backend as _Backend +from .utils import AnnexKey lgr = logging.getLogger("datalad.fuse.fsspec") -class FileState(Enum): - """State of a file in a dataset, as returned by ``get_file_state()``""" - - #: The file is not annexed (e.g. committed to git directly); it is read - #: from disk. - NOT_ANNEXED = 1 - #: The file is annexed but its content is not present locally; it is read - #: from a remote URL. - NO_CONTENT = 2 - #: The file is annexed and its content is present locally; it is read from - #: disk. - HAS_CONTENT = 3 - - -class DatasetAdapter: - """Read access to the files of a single dataset. - - Files that are not annexed, and annexed files whose content is present - locally, are opened from disk. Annexed files without local content are - opened from one of the http(s) URLs found for their git-annex key (see - :meth:`get_urls`), reading only the needed parts of the file. - - Parameters - ---------- - path : str or Path - Top directory of the dataset (any git or git-annex repository). - caching : bool - If true, keep the data fetched from remote URLs in a sparse on-disk - cache under ``/.git/datalad/cache/fsspec/``, to be reused by - subsequent reads (for up to a week). If false, data are only - buffered in memory while a file is open. - mode_transparent : bool - If true, paths of key files under ``.git/annex/objects/`` (the targets - of annexed symlinks) are opened as annexed content, fetched from a - remote URL if not present locally. +class FsspecBackend(_Backend): + """Backend using fsspec's HTTPFileSystem (optionally with disk caching).""" - Notes - ----- - Call :meth:`close` (or use :func:`contextlib.closing`) when done, to stop - the ``git annex`` processes started for the dataset. - """ + name = "fsspec" - def __init__( - self, path: str | Path, caching: bool, mode_transparent: bool = False - ) -> None: - self.path = Path(path) - self.mode_transparent = mode_transparent - ds = Dataset(path) - self.annex: Optional[AnnexRepo] - if isinstance(ds.repo, AnnexRepo): - self.annex = ds.repo - else: - self.annex = None - self.commit_dt = datetime.fromtimestamp( - ds.repo.get_commit_date(), tz=timezone.utc - ) - self.caching = caching + def __init__(self, path: str | Path, caching: bool) -> None: fs = HTTPFileSystem(get_client=get_client) - if self.caching: - self.fs = CachingFileSystem( + if caching: + self.fs: HTTPFileSystem | CachingFileSystem = CachingFileSystem( fs=fs, - # target_protocol='blockcache', cache_storage=os.path.join(path, ".git", "datalad", "cache", "fsspec"), - # cache_check=600, - # block_size=1024, - # check_files=True, - # expiry_times=True, - # same_names=True ) else: self.fs = fs + self._caching = caching - def close(self) -> None: - """Stop the batched ``git annex`` processes started for the dataset""" - if self.annex is not None: - self.annex._batched.clear() - - @methodtools.lru_cache(maxsize=CACHE_SIZE) - def get_file_state(self, relpath: str) -> tuple[FileState, Optional[AnnexKey]]: - """Determine whether a file is annexed and has its content present - - Results are cached (for the most recently queried files). - - Parameters - ---------- - relpath : str - Path of the file relative to the top directory of the dataset. - - Returns - ------- - tuple of (FileState, AnnexKey or None) - The state of the file, and its git-annex key if it is annexed. - """ - p = self.path / relpath - lgr.debug("get_file_state: %s", relpath) - - def handle_path_under_annex_objects( - p: Path, - ) -> tuple[FileState, Optional[AnnexKey]]: - iadok = is_annex_dir_or_key(p) - if isinstance(iadok, AnnexKey): - if p.exists(): - return (FileState.HAS_CONTENT, iadok) - else: - return (FileState.NO_CONTENT, iadok) - else: - return (FileState.NOT_ANNEXED, None) - - # Shortcut handling of content under .git, in particular - annex key paths - if self.mode_transparent and relpath.startswith(".git/"): - return handle_path_under_annex_objects(p) - - # A regular file or git link for which we need to explicitly ask annex about - if not p.is_symlink(): - if p.stat().st_size < 1024 and self.annex is not None: - if self.annex.is_under_annex(relpath, batch=True): - key = AnnexKey.parse(self.annex.get_file_key(relpath, batch=True)) - if self.annex.file_has_content(relpath, batch=True): - return (FileState.HAS_CONTENT, key) - else: - return (FileState.NO_CONTENT, key) - return (FileState.NOT_ANNEXED, None) - - return handle_path_under_annex_objects( - Path(os.path.normpath(p.parent / os.readlink(p))) - ) - - def get_urls(self, key: str) -> Iterator[str]: - """Yield candidate http(s) URLs for the content of an annex key - - URLs are yielded in the order in which they are tried by :meth:`open`: - - 1. http(s) URLs recorded in git-annex for the key, as reported by - ``git annex whereis`` (e.g. those of the ``web`` special remote); - 2. ``annex/objects/...`` locations on the http(s) git remotes that - ``git annex whereis`` lists as having the key, including the - ``annex/objects`` endpoint of Forgejo-aneksajo instances. - - URLs on S3 special remotes with ``exporttree=yes`` are not included; - :meth:`open` falls back to them via :meth:`get_exporttree_urls`. - - Parameters - ---------- - key : str - A git-annex key, e.g. ``str(AnnexKey)``. - """ - assert self.annex is not None - # TODO: switch to batch=True whenever - # https://github.com/datalad/datalad/pull/6379 is merged/released. - # Will need a recent git-annex to work! - whereis = self.annex.whereis(key, output="full", batch=False, key=True) - remote_uuids = [] - for ru, v in whereis.items(): - remote_uuids.append(ru) - for u in v["urls"]: - if is_http_url(u): - yield u - - path_mixed = self.annex._batched.get( - "examinekey", - annex_options=["--format=annex/objects/${hashdirmixed}${key}/${key}\\n"], - path=self.annex.path, - )(key) - path_lower = self.annex._batched.get( - "examinekey", - annex_options=["--format=annex/objects/${hashdirlower}${key}/${key}\\n"], - path=self.annex.path, - )(key) - - uuid2remote_url = {} - aneksajo_uuids: set[str] = set() - for r in self.annex.get_remotes(): - if (ru := self.annex.config.get(f"remote.{r}.annex-uuid")) is None: - continue - if (remote_url := self.annex.config.get(f"remote.{r}.url")) is None: - continue - remote_url = self.annex.config.rewrite_url(remote_url) - uuid2remote_url[ru] = remote_url - # Detect Forgejo-aneksajo instances via API probe (cached). - # TODO: pushurl could be different from url, should also check - # remote.{r}.pushurl config - # TODO: SSH remote URLs not yet supported -- would need to - # derive the HTTP base URL from the SSH URL - if is_http_url(remote_url) and _is_aneksajo(remote_url): - aneksajo_uuids.add(ru) - - for ru in remote_uuids: - try: - base_url = uuid2remote_url[ru] - except KeyError: - continue - if is_http_url(base_url): - base_stripped = base_url.rstrip("/") - # Forgejo/Gitea with aneksajo: use annex/objects endpoint - # which supports HEAD and Range requests. - # See https://codeberg.org/forgejo-aneksajo/forgejo-aneksajo/issues/111 - if ru in aneksajo_uuids and base_stripped.endswith(".git"): - forge_base = base_stripped[:-4].rstrip("/") - yield forge_base + "/" + path_lower - if base_stripped.lower().endswith("/.git"): - paths = [path_mixed, path_lower] - else: - paths = [ - path_lower, - path_mixed, - f".git/{path_lower}", - f".git/{path_mixed}", - ] - for p in paths: - yield base_stripped + "/" + p - - @methodtools.lru_cache(maxsize=1) - def _get_exporttree_remotes(self) -> list[dict[str, str]]: - """Get S3 exporttree remotes with public URLs. - - Parses the git-annex branch remote.log once (cached per - DatasetAdapter instance) to find S3 special remotes configured - with ``exporttree=yes`` and a usable ``publicurl``. - - This is a workaround for legacy datasets that lack proper - versioned S3 URLs in their git-annex metadata. - See https://github.com/OpenNeuroOrg/openneuro/issues/3875 - - Returns - ------- - list of dict - Each dict has keys: ``uuid``, ``publicurl``, ``fileprefix``, - ``bucket``, ``host``. - """ - try: - result = subprocess.run( - ["git", "-C", str(self.path), "show", "git-annex:remote.log"], - capture_output=True, - text=True, - check=True, - ) - except subprocess.CalledProcessError: - lgr.debug("Could not read git-annex:remote.log for %s", self.path) - return [] - - remotes: list[dict[str, str]] = [] - for line in result.stdout.strip().splitlines(): - if not line or line.startswith("#"): - continue - parts = line.split() - if len(parts) < 2: - continue - uuid = parts[0] - config: dict[str, str] = {} - for token in parts[1:]: - if "=" in token: - k, v = token.split("=", 1) - config[k] = v - if ( - config.get("type") == "S3" - and config.get("exporttree") == "yes" - and config.get("publicurl", "no").startswith("http") - ): - remotes.append( - { - "uuid": uuid, - "publicurl": config["publicurl"].rstrip("/"), - "fileprefix": config.get("fileprefix", ""), - "bucket": config.get("bucket", ""), - "host": config.get("host", "s3.amazonaws.com"), - } - ) - return remotes - - @staticmethod - def _list_s3_versions( - bucket: str, - object_key: str, - host: str = "s3.amazonaws.com", - ) -> list[dict[str, Any]]: - """List all S3 object versions for a key. - - Uses ``boto3`` to call ``ListObjectVersions`` with anonymous - credentials (for public buckets). - - Parameters - ---------- - bucket : str - S3 bucket name (e.g., ``openneuro.org``). - object_key : str - Full object key including fileprefix (e.g., - ``ds000113/sub-01/.../bold.nii.gz``). - host : str - S3 endpoint hostname (default: ``s3.amazonaws.com``). - - Returns - ------- - list of dict - Each dict has keys: ``VersionId``, ``Size``, ``ETag``, - ``IsLatest``. - """ - try: - endpoint_url = f"https://{host}" - client = boto3.client( - "s3", - endpoint_url=endpoint_url, - config=BotocoreConfig(signature_version=UNSIGNED), - ) - response = client.list_object_versions( - Bucket=bucket, Prefix=object_key - ) - except Exception as e: - lgr.debug( - "Failed to list S3 versions for %s/%s: %s", - bucket, object_key, e, - ) - return [] - - versions: list[dict[str, Any]] = [] - for v in response.get("Versions", []): - # Only include exact key matches (prefix query may return others) - if v.get("Key") == object_key: - versions.append( - { - "VersionId": v.get("VersionId", ""), - "Size": v.get("Size", 0), - "ETag": v.get("ETag", ""), - "IsLatest": v.get("IsLatest", False), - } - ) - return versions - - @staticmethod - def _match_s3_version( - versions: list[dict[str, Any]], expected_size: int - ) -> Optional[str]: - """Match the correct S3 object version by file size. - - Parameters - ---------- - versions : list of dict - S3 version list from :meth:`_list_s3_versions`. - expected_size : int - Expected file size from ``AnnexKey.size``. - - Returns - ------- - str or None - Matched versionId, or ``None`` if no version matches. - - Raises - ------ - ValueError - If multiple versions match by size but have different ETags - (ambiguous content — refuse to guess). - """ - matches = [v for v in versions if v["Size"] == expected_size] - if not matches: - return None - if len(matches) == 1: - return str(matches[0]["VersionId"]) - # Multiple matches — check ETags - etags = {v["ETag"] for v in matches} - if len(etags) == 1: - # Same content uploaded multiple times; prefer the latest - for v in matches: - if v["IsLatest"]: - return str(v["VersionId"]) - return str(matches[0]["VersionId"]) - raise ValueError( - f"Ambiguous S3 versions: {len(matches)} versions match size " - f"{expected_size} but have {len(etags)} distinct ETags. " - f"Cannot determine correct version." - ) - - def get_exporttree_urls( - self, relpath: str, key: AnnexKey - ) -> Iterator[str]: - """Yield versioned URLs for file on S3 exporttree remotes. - - Workaround for datasets lacking proper versioned URLs in - git-annex metadata. Constructs URLs from the remote's - ``publicurl`` + ``fileprefix`` and resolves the correct S3 - object version by matching ``key.size``. - - Parameters - ---------- - relpath : str - File path relative to dataset root (tree path). - key : AnnexKey - Annex key with expected file size for version matching. - - Yields - ------ - str - Versioned HTTP URLs (``...?versionId=...``) or unversioned - URLs as fallback. - """ - remotes = self._get_exporttree_remotes() - if not remotes: - return - - for remote in remotes: - publicurl = remote["publicurl"] - fileprefix = remote["fileprefix"] - bucket = remote["bucket"] - host = remote["host"] - object_key = f"{fileprefix}{relpath}" - base_url = f"{publicurl}/{object_key}" - - if key.size is not None: - versions = self._list_s3_versions(bucket, object_key, host) - if versions: - try: - version_id = self._match_s3_version( - versions, key.size - ) - except ValueError as e: - lgr.warning( - "%s: %s", relpath, e - ) - continue - if version_id: - yield f"{base_url}?versionId={version_id}" - continue - else: - lgr.debug( - "%s: no S3 version matches size %d at %s", - relpath, - key.size, - base_url, - ) - continue - - # Fallback: no size info or version listing failed — - # try unversioned URL (returns latest version) - lgr.warning( - "%s: falling back to unversioned S3 URL %s " - "(cannot verify correct version)", - relpath, - base_url, - ) - yield base_url - - def open( + def can_handle( self, - relpath: str, - mode: str = "rb", - encoding: str = "utf-8", - errors: Optional[str] = None, - ) -> IO: - """Open a file of the dataset for reading - - Parameters - ---------- - relpath : str - Path of the file relative to the top directory of the dataset. - mode : str - ``"rb"`` (default) to get a binary file object, ``"r"`` or - ``"rt"`` to get a text file object. - encoding : str - Encoding to use in text mode. - errors : str, optional - How to handle encoding errors in text mode, as for :func:`open`. - - Returns - ------- - file object - A seekable, read-only file object. Files read from disk are - regular Python file objects; files read from a URL are fsspec - file objects. + key: Optional[AnnexKey], # noqa: U100 + mode: str, # noqa: U100 + relpath: Optional[str] = None, # noqa: U100 + ) -> bool: + return True # fsspec handles everything - Raises - ------ - NotImplementedError - If ``mode`` is not one of the supported read modes. - IOError - If the content of an annexed file is not present locally and none - of its candidate URLs could be opened. - """ - if mode not in ("r", "rb", "rt"): - raise NotImplementedError("Only modes 'r', 'rb', and 'rt' are supported") - if mode == "rb": - kwargs = {} - else: - kwargs = {"encoding": encoding, "errors": errors} - fstate, key = self.get_file_state(relpath) - if fstate is FileState.NOT_ANNEXED: - lgr.debug("%s: not under annex", relpath) - else: - lgr.debug( - "%s: under annex, %s content", - relpath, - "has" if fstate is FileState.HAS_CONTENT else "does not have", - ) - if fstate is FileState.NO_CONTENT: - lgr.debug("%s: opening via fsspec", relpath) - for url in self.get_urls(str(key)): - try: - lgr.debug("%s: Attempting to open via URL %s", relpath, url) - return self.fs.open(url, mode, **kwargs) # type: ignore - except BlocksizeMismatchError as e: - lgr.warning( - "%s: Blocksize mismatch: %s; deleting cached file and" - " re-opening", - relpath, - e, - ) - self.fs.pop_from_cache(url) - return self.fs.open(url, mode, **kwargs) # type: ignore - except FileNotFoundError as e: - lgr.debug( - "Failed to open file %s at URL %s: %s", relpath, url, str(e) - ) - # Fallback: try S3 exporttree URLs (workaround for datasets - # lacking proper versioned URLs — see openneuro#3875) - if key is not None: - for url in self.get_exporttree_urls(relpath, key): - try: - lgr.debug( - "%s: Attempting exporttree URL %s", relpath, url - ) - return self.fs.open(url, mode, **kwargs) # type: ignore - except BlocksizeMismatchError as e: - lgr.warning( - "%s: Blocksize mismatch: %s; deleting cached file" - " and re-opening", - relpath, - e, - ) - self.fs.pop_from_cache(url) - return self.fs.open(url, mode, **kwargs) # type: ignore - except FileNotFoundError as e: - lgr.debug( - "Failed to open file %s at exporttree URL %s: %s", - relpath, - url, - str(e), - ) - raise IOError( - f"Could not find a usable URL for {relpath} within {self.path}" - ) - else: - lgr.debug("%s: opening directly", relpath) - return open(self.path / relpath, mode, **kwargs) # type: ignore + def open_url(self, url: str, mode: str = "rb", **kwargs: Any) -> IO: + try: + return self.fs.open(url, mode, **kwargs) # type: ignore[no-any-return] + except BlocksizeMismatchError: + # Eviction only makes sense for CachingFileSystem; on a plain + # HTTPFileSystem there is no cache to evict, so re-raise. + if not self._caching: + raise + lgr.warning("Blocksize mismatch for %s; clearing cache and retrying", url) + self.fs.pop_from_cache(url) + return self.fs.open(url, mode, **kwargs) # type: ignore[no-any-return] def clear(self) -> None: - """Remove the on-disk cache of the dataset (only if ``caching``)""" - if self.caching: + if self._caching: self.fs.clear_cache() -class FsspecAdapter: - """Read access to the files of a dataset and its installed subdatasets. - - Each path is mapped to the (sub)dataset containing it, and a - :class:`DatasetAdapter` is created for that dataset on first use. Use it - as a context manager, so that the ``git annex`` processes started for the - datasets are stopped on exit. - - Parameters - ---------- - root : str or Path - Top directory of the (super)dataset. - caching : bool - Passed to each :class:`DatasetAdapter`. - mode_transparent : bool - Passed to each :class:`DatasetAdapter`. - - Notes - ----- - Use an absolute ``root``, and absolute paths under it for the methods. - Paths relative to ``root`` or to the current directory are not supported. - """ - - def __init__( - self, root: str | Path, caching: bool, mode_transparent: bool = False - ) -> None: - self.root = Path(root) - self.mode_transparent = mode_transparent - self.caching = caching - self.datasets: dict[Path, DatasetAdapter] = {} - - def __enter__(self) -> FsspecAdapter: - return self - - def __exit__( - self, - _exc_type: Optional[type[BaseException]], - _exc_val: Optional[BaseException], - _exc_tb: Optional[TracebackType], - ) -> None: - for ds in self.datasets.values(): - ds.close() - self.datasets.clear() - - @methodtools.lru_cache(maxsize=CACHE_SIZE) - # TODO: optimize "caching" more since for all files under the same directory - # they all would belong to the same dataset - def get_dataset_path(self, path: str | Path) -> Path: - """Return the top directory of the (sub)dataset containing ``path``""" - path = Path(self.root, path) - dspath = get_dataset_root(path) - if dspath is None: - raise ValueError(f"Path not under DataLad: {path}") - dspath = Path(dspath) - assert isinstance(dspath, Path) - try: - dspath.relative_to(self.root) - except ValueError: - raise ValueError(f"Path not under root dataset: {path}") - return dspath - - def resolve_dataset(self, filepath: str | Path) -> tuple[DatasetAdapter, str]: - """Return the adapter for the dataset containing ``filepath`` - - Returns - ------- - tuple of (DatasetAdapter, str) - The adapter, and the path of ``filepath`` relative to the - dataset's top directory. - """ - dspath = self.get_dataset_path(filepath) - try: - dsap = self.datasets[dspath] - except KeyError: - dsap = self.datasets[dspath] = DatasetAdapter( - dspath, - mode_transparent=self.mode_transparent, - caching=self.caching, - ) - relpath = str(Path(filepath).relative_to(dspath)) - return dsap, relpath - - def open( - self, - filepath: str | Path, - mode: str = "rb", - encoding: str = "utf-8", - errors: Optional[str] = None, - ) -> IO: - """Open a file for reading; see :meth:`DatasetAdapter.open`""" - dsap, relpath = self.resolve_dataset(filepath) - lgr.debug( - "%s: path resolved to %s in dataset at %s", filepath, relpath, dsap.path - ) - return dsap.open(relpath, mode=mode, encoding=encoding, errors=errors) - - def get_file_state( - self, filepath: str | Path - ) -> tuple[FileState, Optional[AnnexKey]]: - """Return state and key of a file; see `DatasetAdapter.get_file_state`""" - dsap, relpath = self.resolve_dataset(filepath) - return cast(Tuple[FileState, Optional[AnnexKey]], dsap.get_file_state(relpath)) - - def is_under_annex(self, filepath: str | Path) -> bool: - """Tell whether a file is annexed""" - dsap, relpath = self.resolve_dataset(filepath) - fstate, _ = dsap.get_file_state(relpath) - return fstate is not FileState.NOT_ANNEXED - - def get_commit_datetime(self, filepath: str | Path) -> datetime: - """Return the date of ``HEAD`` in the dataset containing ``filepath``""" - dsap, _ = self.resolve_dataset(filepath) - return dsap.commit_dt - - -def is_http_url(s: str) -> bool: - return s.lower().startswith(("http://", "https://")) - - -_aneksajo_cache: dict[str, bool] = {} - - -def _is_aneksajo(base_url: str) -> bool: - """Check if a URL points to a Forgejo-aneksajo instance. - - Probes ``{scheme}://{host}/api/forgejo/v1/version`` and checks whether - the version string contains ``git-annex``, which indicates the - forgejo-aneksajo fork. - - Results are cached per ``scheme://host:port`` for the process lifetime. - """ - parsed = urlparse(base_url) - # Cache key without userinfo so credentials don't fragment the cache - host = parsed.hostname or "" - port_suffix = f":{parsed.port}" if parsed.port else "" # noqa: E231 - cache_key = f"{parsed.scheme}://{host}{port_suffix}" # noqa: E231 - - if cache_key in _aneksajo_cache: - return _aneksajo_cache[cache_key] - - try: - api_url = f"{cache_key}/api/forgejo/v1/version" - req = urllib.request.Request(api_url, method="GET") - req.add_header("Accept", "application/json") - with urllib.request.urlopen(req, timeout=10) as resp: - data = json.loads(resp.read().decode()) - result = "git-annex" in data.get("version", "") - except Exception: - lgr.debug("_is_aneksajo(%s) probe failed", cache_key, exc_info=True) - result = False - - _aneksajo_cache[cache_key] = result - lgr.debug("_is_aneksajo(%s) = %s", cache_key, result) - return result +# -- Async HTTP helpers (fsspec-specific) ------------------------------------ async def on_request_start( @@ -755,3 +92,30 @@ async def get_client(**kwargs: Any) -> RetryClient: ), retry_options=ListRetry(timeouts=[1, 2, 6, 15, 36]), ) + + +# -- Backward compatibility -------------------------------------------------- + +_COMPAT_MAP: dict[str, tuple[str, str]] = { + # Only names that existed in the pre-refactoring fsspec.py + # name -> (module, canonical_name) + "FsspecAdapter": ("datalad_fuse.adapter", "RemoteFilesystemAdapter"), + "DatasetAdapter": ("datalad_fuse.adapter", "DatasetAdapter"), + "FileState": ("datalad_fuse.adapter", "FileState"), + "is_http_url": ("datalad_fuse.adapter", "is_http_url"), +} + + +def __getattr__(name: str) -> Any: + if name in _COMPAT_MAP: + module_path, canonical_name = _COMPAT_MAP[name] + mod = importlib.import_module(module_path) + obj = getattr(mod, canonical_name) + warnings.warn( + f"Importing {name!r} from 'datalad_fuse.fsspec' is deprecated. " + f"Use 'from {module_path} import {canonical_name}' instead.", + DeprecationWarning, + stacklevel=2, + ) + return obj + raise AttributeError(f"module 'datalad_fuse.fsspec' has no attribute {name!r}") diff --git a/datalad_fuse/fsspec_cache_clear.py b/datalad_fuse/fsspec_cache_clear.py index 12a28de..8a90010 100644 --- a/datalad_fuse/fsspec_cache_clear.py +++ b/datalad_fuse/fsspec_cache_clear.py @@ -11,7 +11,7 @@ from datalad.support.constraints import EnsureNone from datalad.support.param import Parameter -from .fsspec import DatasetAdapter +from .adapter import DatasetAdapter @build_doc diff --git a/datalad_fuse/fsspec_head.py b/datalad_fuse/fsspec_head.py index 4ce63d2..def7b36 100644 --- a/datalad_fuse/fsspec_head.py +++ b/datalad_fuse/fsspec_head.py @@ -16,7 +16,7 @@ from datalad.support.constraints import EnsureInt, EnsureNone, EnsureStr from datalad.support.param import Parameter -from .fsspec import FsspecAdapter +from .adapter import RemoteFilesystemAdapter DEFAULT_LINES = 10 @@ -57,13 +57,22 @@ class FsspecHead(Interface): args=("--caching",), choices=["none", "ondisk"], default="none", - doc="Whether to cache fsspec'ed files on disk on not at all", + doc="Whether to cache data fetched from remote URLs on disk or not at all", ), "path": Parameter( args=("path",), doc="Path to an annexed file to show the leading contents of", constraints=EnsureStr(), ), + "backends": Parameter( + args=("--backends",), + doc=( + "Comma-separated list of backends to try for remote file" + " access, in priority order. Available: remfile, fsspec." + " Default: remfile,fsspec" + ), + constraints=EnsureStr() | EnsureNone(), + ), } @staticmethod @@ -76,14 +85,18 @@ def __call__( bytes: Optional[int] = None, mode_transparent: bool = False, caching: str | None = None, + backends: str | None = None, ) -> Iterator[Dict[str, Any]]: ds = require_dataset(dataset, purpose="fetch file data", check_installed=True) if lines is not None and bytes is not None: raise ValueError("'lines' and 'bytes' are mutually exclusive") elif lines is None and bytes is None: lines = DEFAULT_LINES - with FsspecAdapter( - ds.path, mode_transparent=mode_transparent, caching=caching == "ondisk" + with RemoteFilesystemAdapter( + ds.path, + mode_transparent=mode_transparent, + caching=caching == "ondisk", + backends=backends, ) as fsa: if not os.path.isabs(path): path = os.path.join(ds.path, path) diff --git a/datalad_fuse/fuse_.py b/datalad_fuse/fuse_.py index e6e736c..9867800 100644 --- a/datalad_fuse/fuse_.py +++ b/datalad_fuse/fuse_.py @@ -21,8 +21,8 @@ from fuse import FuseOSError, Operations import methodtools +from .adapter import RemoteFilesystemAdapter from .consts import CACHE_SIZE -from .fsspec import FsspecAdapter # Make it relatively small since we are aiming for metadata records ATM # Seems of no real good positive net ATM @@ -66,8 +66,9 @@ def wrapped(self: DataLadFUSE, path: str, *args: P.args, **kwargs: P.kwargs) -> class DataLadFUSE(Operations): # LoggingMixIn, """fusepy file system exposing a dataset, as used by ``datalad fusefs`` - Files are read via an :class:`~datalad_fuse.fsspec.FsspecAdapter`, so - annexed files without local content are read from their remote URLs. + Files are read via a + :class:`~datalad_fuse.adapter.RemoteFilesystemAdapter`, so annexed files + without local content are read from their remote URLs. Unless ``mode_transparent`` is set, annexed files appear as regular files. For files without local content, the size is taken from their annex key and the modification time is the date of the ``HEAD`` commit. Files @@ -80,11 +81,14 @@ class DataLadFUSE(Operations): # LoggingMixIn, top directory of the dataset to expose. caching : bool Whether to cache remote data on disk; see - :class:`~datalad_fuse.fsspec.DatasetAdapter`. + :class:`~datalad_fuse.adapter.DatasetAdapter`. mode_transparent : bool Whether to expose the ``.git`` directories of the datasets (hidden by default). Annexed files without local content then appear as symlinks into ``.git/annex/objects/``. + backends : str, optional + Comma-separated, priority-ordered backends to read remote files with; + see :class:`~datalad_fuse.adapter.DatasetAdapter`. Examples -------- @@ -101,13 +105,20 @@ class DataLadFUSE(Operations): # LoggingMixIn, _counter_offset = 1000 def __init__( - self, root: str, caching: bool, mode_transparent: bool = False + self, + root: str, + caching: bool, + mode_transparent: bool = False, + backends: Optional[str] = None, ) -> None: self.root = op.realpath(root) self.mode_transparent = mode_transparent self.rwlock = Lock() - self._adapter = FsspecAdapter( - root, mode_transparent=mode_transparent, caching=caching + self._adapter = RemoteFilesystemAdapter( + root, + mode_transparent=mode_transparent, + caching=caching, + backends=backends, ) self._fhdict: dict[int, Optional[IO[bytes]]] = {} # fh to fsspec_file, already opened (we are RO for now, so can just open diff --git a/datalad_fuse/remfile.py b/datalad_fuse/remfile.py new file mode 100644 index 0000000..eeee978 --- /dev/null +++ b/datalad_fuse/remfile.py @@ -0,0 +1,209 @@ +"""RemfileBackend — HDF5-optimised remote file access via the *remfile* library.""" + +from __future__ import annotations + +import logging +import os.path +from pathlib import Path, PurePosixPath +import shutil +from types import ModuleType, TracebackType +from typing import IO, Any, Optional, cast +import urllib.request + +from .backends import Backend +from .utils import AnnexKey + +lgr = logging.getLogger("datalad.fuse.remfile") + + +def _get_remfile() -> Optional[ModuleType]: + """Lazy import of remfile; returns None if not installed.""" + try: + import remfile + + return remfile # type: ignore[no-any-return] + except ImportError: + return None + + +class RemfileBackend(Backend): + """Backend using remfile for HDF5-structured files (.nwb, .h5, etc.).""" + + name = "remfile" + + # File extensions this backend handles (HDF5-structured formats) + EXTENSIONS = frozenset({".nwb", ".h5", ".hdf5", ".hdf", ".he5", ".nc", ".nc4"}) + + #: Timeout, in seconds, for the size probe in :meth:`open_url`. + PROBE_TIMEOUT = 10.0 + + def __init__(self, path: str | Path, caching: bool) -> None: + remfile_mod = _get_remfile() + if remfile_mod is None: + raise ImportError("remfile is not installed") + self._remfile: ModuleType = remfile_mod + self._cache_dir = os.path.join(path, ".git", "datalad", "cache", "remfile") + # Without a disk cache, --caching=ondisk would silently be a no-op for + # exactly the large files it matters most for. + self._disk_cache = self._remfile.DiskCache(self._cache_dir) if caching else None + + def can_handle( + self, key: Optional[AnnexKey], mode: str, relpath: Optional[str] = None + ) -> bool: + if mode != "rb": + return False + suffix = key.suffix if key is not None else None + if suffix is None and relpath is not None: + # URL/VURL keys (``addurl --fast``/``--relaxed``) carry no suffix, + # so fall back to the name the user sees in the tree. + suffix = PurePosixPath(relpath).suffix or None + if suffix is None: + return False + return suffix.lower() in self.EXTENSIONS + + def _probe_size(self, url: str) -> Optional[int]: + """Return the size of *url*, or None if the server did not say. + + ``remfile.File`` determines the size itself, but retries transient + failures 8 times with exponential backoff (~25 s). The adapter tries + URLs one after another, and ``DataLadFUSE.open()`` holds the global + rwlock while it does, so a single dead URL would stall an entire mount. + Probing here instead keeps the fall-through to the next URL/backend as + fast as it is for fsspec: one ranged request, one short timeout, no + retries. + """ + req = urllib.request.Request(url, headers={"Range": "bytes=0-0"}) + with urllib.request.urlopen(req, timeout=self.PROBE_TIMEOUT) as resp: + # 206: "Content-Range: bytes 0-0/" ( may be "*") + content_range = resp.headers.get("Content-Range") + if content_range is not None: + total = content_range.rsplit("/", 1)[-1].strip() + if total.isdigit(): + return int(total) + elif resp.status == 200: + # Server ignored the Range header and sent the whole body + content_length = resp.headers.get("Content-Length") + if content_length is not None and content_length.isdigit(): + return int(content_length) + lgr.debug("Could not determine size of %s from probe response", url) + return None + + def open_url(self, url: str, mode: str = "rb", **kwargs: Any) -> IO: # noqa: U100 + if mode != "rb": + # Mirror can_handle()'s mode check: remfile is binary-only, so + # refuse text modes explicitly rather than silently returning a + # binary stream when called directly (bypassing can_handle()). + raise NotImplementedError( + f"RemfileBackend only supports mode='rb', got {mode!r}" + ) + size = self._probe_size(url) + # `_size` is private API; see the remfile pin in setup.cfg. + remfile_obj = self._remfile.File(url, _size=size, disk_cache=self._disk_cache) + return cast(IO, RemfileWrapper(remfile_obj, url)) + + def clear(self) -> None: + if self._disk_cache is not None: + shutil.rmtree(self._cache_dir, ignore_errors=True) + + +class RemfileWrapper: + """Wraps ``remfile.File`` to satisfy the contracts expected by datalad-fuse. + + Adds context manager protocol, line iteration (for ``fsspec_head``), and an + ``info()`` method compatible with ``file_getattr`` in *fuse_.py*. + """ + + _ITER_CHUNK = 8192 + # Hard cap on bytes returned per __next__ call. Without this, iterating + # a binary file (e.g. HDF5/NWB) — which has no '\n' — would download the + # entire remote file in a single iteration step. + _MAX_LINE_BYTES = 1 << 20 # 1 MiB + + def __init__(self, remfile_obj: Any, url: str) -> None: + self._f = remfile_obj + self._url = url + self.closed = False + + def read(self, size: Optional[int] = -1) -> bytes: + """Read up to *size* bytes, or to the end of the file by default""" + # ``remfile.RemFile.read()`` requires an explicit, non-negative size: + # it returns b"" (and rewinds by one) for -1, never clamps at EOF, and + # issues an unsatisfiable Range request when reading at EOF. Clamp + # here so the wrapper behaves like a regular file object. + remaining = max(self._f.length - self._f.tell(), 0) + if size is None or size < 0 or size > remaining: + size = remaining + if not size: + return b"" + return self._f.read(size) # type: ignore[no-any-return] + + def seek(self, offset: int, whence: int = 0) -> int: + """Set the read position, and return it""" + # ``remfile.RemFile.seek()`` returns None + self._f.seek(offset, whence) + return self.tell() + + def tell(self) -> int: + """Return the current read position""" + return self._f.tell() # type: ignore[no-any-return] + + def close(self) -> None: + """Close the file""" + self._f.close() + self.closed = True + + def readable(self) -> bool: + return True + + def seekable(self) -> bool: + return True + + def writable(self) -> bool: + return False + + def __enter__(self) -> RemfileWrapper: + return self + + def __exit__( + self, + _exc_type: Optional[type[BaseException]], + _exc_val: Optional[BaseException], + _exc_tb: Optional[TracebackType], + ) -> None: + self.close() + + def __iter__(self) -> RemfileWrapper: + return self + + def __next__(self) -> bytes: + chunks: list[bytes] = [] + total = 0 + while total < self._MAX_LINE_BYTES: + chunk = self.read(self._ITER_CHUNK) + if not chunk: + if chunks: + return b"".join(chunks) + raise StopIteration + idx = chunk.find(b"\n") + if idx != -1: + chunks.append(chunk[: idx + 1]) + # Seek back past the bytes we read beyond the newline + overshoot = len(chunk) - idx - 1 + if overshoot: + self._f.seek(-overshoot, 1) + return b"".join(chunks) + chunks.append(chunk) + total += len(chunk) + # Hit the cap without finding a newline — return what we have so the + # caller still makes progress instead of OOM-ing on a binary file. + return b"".join(chunks) + + def info(self) -> dict[str, Any]: + """Minimal info dict matching the fsspec convention. + + ``DataLadFUSE.getattr(path, fh)`` reaches this for every ``fstat()`` on + an open handle, so it must not hit the network: a HEAD request would + fail outright against presigned S3 GET URLs, which reject HEAD. + remfile already determined the size when the file was opened. + """ + return {"type": "file", "size": self._f.length} diff --git a/datalad_fuse/tests/test_backends.py b/datalad_fuse/tests/test_backends.py new file mode 100644 index 0000000..c242bbb --- /dev/null +++ b/datalad_fuse/tests/test_backends.py @@ -0,0 +1,1102 @@ +"""Tests for the pluggable backend system (remfile, fsspec) and adapters.""" + +from __future__ import annotations + +import asyncio +import builtins +from datetime import datetime, timezone +from io import BytesIO +import logging +from pathlib import Path +import sys +from types import SimpleNamespace +from typing import Iterator, Optional +from unittest.mock import MagicMock, patch + +from fsspec.exceptions import BlocksizeMismatchError +import pytest + +from datalad_fuse.adapter import ( + DatasetAdapter, + FileState, + RemoteFilesystemAdapter, + create_backends, + is_http_url, + resolve_backends, +) +from datalad_fuse.backends import DEFAULT_BACKENDS, Backend +from datalad_fuse.fsspec import FsspecBackend, on_request_start +from datalad_fuse.remfile import RemfileBackend, RemfileWrapper, _get_remfile +from datalad_fuse.utils import AnnexKey + +try: + import remfile # noqa: F401 + + _has_remfile = True +except ImportError: + _has_remfile = False + +requires_remfile = pytest.mark.skipif(not _has_remfile, reason="remfile not installed") + + +# -- RemfileBackend.can_handle tests ----------------------------------------- + + +@requires_remfile +class TestRemfileBackendCanHandle: + """Tests for RemfileBackend extension detection via can_handle.""" + + pytestmark = pytest.mark.ai_generated + + @pytest.fixture() + def remfile_backend(self, tmp_path) -> RemfileBackend: + return RemfileBackend(tmp_path, caching=False) + + @pytest.mark.parametrize("ext", sorted(RemfileBackend.EXTENSIONS)) + def test_hdf5_extensions_accepted( + self, remfile_backend: RemfileBackend, ext: str + ) -> None: + key = AnnexKey(backend="MD5E", name="abc123", size=100, suffix=ext) + assert remfile_backend.can_handle(key, "rb") is True + + @pytest.mark.parametrize("ext", [".txt", ".png", ".csv", ".json", ".zip"]) + def test_non_hdf5_extensions_rejected( + self, remfile_backend: RemfileBackend, ext: str + ) -> None: + key = AnnexKey(backend="MD5E", name="abc123", size=100, suffix=ext) + assert remfile_backend.can_handle(key, "rb") is False + + def test_suffixless_key_falls_back_to_relpath( + self, remfile_backend: RemfileBackend + ) -> None: + """URL/VURL keys carry no suffix; dispatch on the path the user sees. + + ``git annex addurl --fast/--relaxed`` (and ``datalad addurls --fast``) + produce such keys, so an .nwb file would otherwise never reach remfile. + """ + key = AnnexKey(backend="VURL", name="abc123", size=None, suffix=None) + assert remfile_backend.can_handle(key, "rb", "sub-01/sub-01.nwb") is True + assert remfile_backend.can_handle(key, "rb", "README.md") is False + assert remfile_backend.can_handle(None, "rb", "x.h5") is True + + def test_key_suffix_wins_over_relpath( + self, remfile_backend: RemfileBackend + ) -> None: + key = AnnexKey(backend="MD5E", name="abc123", size=100, suffix=".nwb") + assert remfile_backend.can_handle(key, "rb", "renamed.txt") is True + + def test_no_suffix_rejected(self, remfile_backend: RemfileBackend) -> None: + key = AnnexKey(backend="MD5", name="abc123", size=100, suffix=None) + assert remfile_backend.can_handle(key, "rb") is False + + def test_none_key_rejected(self, remfile_backend: RemfileBackend) -> None: + assert remfile_backend.can_handle(None, "rb") is False + + def test_case_insensitive(self, remfile_backend: RemfileBackend) -> None: + key = AnnexKey(backend="MD5E", name="abc123", size=100, suffix=".NWB") + assert remfile_backend.can_handle(key, "rb") is True + + def test_text_mode_rejected(self, remfile_backend: RemfileBackend) -> None: + key = AnnexKey(backend="MD5E", name="abc123", size=100, suffix=".nwb") + assert remfile_backend.can_handle(key, "r") is False + assert remfile_backend.can_handle(key, "rt") is False + + @pytest.mark.parametrize("mode", ["r", "rt", "w", "wb"]) + def test_open_url_rejects_non_binary_mode( + self, remfile_backend: RemfileBackend, mode: str + ) -> None: + """open_url() refuses non-'rb' modes even when called directly. + + ``can_handle()`` already returns False for these, so the backend chain + will never call open_url() with such a mode in practice — but if a + caller bypasses can_handle() we want a loud failure rather than a + silently-binary handle. + """ + with pytest.raises(NotImplementedError, match="mode='rb'"): + remfile_backend.open_url("http://example.com/x.h5", mode=mode) + + +def _mock_remfile_backend(tmp_path, caching: bool = False): + """Build a RemfileBackend against a stand-in ``remfile`` module. + + Lets the probe / disk-cache behaviour be tested without remfile installed. + """ + fake = MagicMock() + with patch("datalad_fuse.remfile._get_remfile", return_value=fake): + return RemfileBackend(tmp_path, caching=caching), fake + + +def _probe_response(headers: dict, status: int = 206) -> MagicMock: + resp = MagicMock() + resp.headers = headers + resp.status = status + resp.__enter__ = MagicMock(return_value=resp) + resp.__exit__ = MagicMock(return_value=None) + return resp + + +@pytest.mark.ai_generated +class TestRemfileBackendProbe: + """open_url() probes the size itself instead of letting remfile do it. + + remfile retries its own ``Content-Length`` probe 8x with backoff (~25 s), + so an unreachable URL would stall the backend chain — and, in ``fusefs``, + the whole mount, since ``open()`` holds the global rwlock. + """ + + def test_probe_result_passed_to_remfile(self, tmp_path, monkeypatch) -> None: + backend, fake = _mock_remfile_backend(tmp_path) + resp = _probe_response({"Content-Range": "bytes 0-0/177728"}) + urlopen = MagicMock(return_value=resp) + monkeypatch.setattr("urllib.request.urlopen", urlopen) + backend.open_url("http://example.com/x.h5") + assert fake.File.call_args.kwargs["_size"] == 177728 + # A short, retry-free timeout is the point of probing here. + assert urlopen.call_args.kwargs["timeout"] == RemfileBackend.PROBE_TIMEOUT + req = urlopen.call_args.args[0] + assert req.get_header("Range") == "bytes=0-0" + + def test_probe_falls_back_to_content_length(self, tmp_path, monkeypatch) -> None: + """A server that ignores Range answers 200 with the full length.""" + backend, fake = _mock_remfile_backend(tmp_path) + resp = _probe_response({"Content-Length": "4096"}, status=200) + monkeypatch.setattr("urllib.request.urlopen", MagicMock(return_value=resp)) + backend.open_url("http://example.com/x.h5") + assert fake.File.call_args.kwargs["_size"] == 4096 + + def test_unknown_size_lets_remfile_probe(self, tmp_path, monkeypatch) -> None: + """An unparsable size is not fatal — remfile falls back to its own.""" + backend, fake = _mock_remfile_backend(tmp_path) + resp = _probe_response({"Content-Range": "bytes 0-0/*"}) + monkeypatch.setattr("urllib.request.urlopen", MagicMock(return_value=resp)) + backend.open_url("http://example.com/x.h5") + assert fake.File.call_args.kwargs["_size"] is None + + def test_probe_failure_propagates_without_opening( + self, tmp_path, monkeypatch + ) -> None: + """A dead URL fails at the probe, so the chain moves on immediately.""" + backend, fake = _mock_remfile_backend(tmp_path) + monkeypatch.setattr( + "urllib.request.urlopen", MagicMock(side_effect=OSError("unreachable")) + ) + with pytest.raises(OSError, match="unreachable"): + backend.open_url("http://example.com/x.h5") + fake.File.assert_not_called() + + +@pytest.mark.ai_generated +class TestRemfileBackendCaching: + """``--caching=ondisk`` must apply to remfile-handled files too.""" + + def test_no_disk_cache_without_caching(self, tmp_path) -> None: + backend, fake = _mock_remfile_backend(tmp_path, caching=False) + fake.DiskCache.assert_not_called() + with patch("urllib.request.urlopen", MagicMock(side_effect=OSError("x"))): + pass + assert backend._disk_cache is None + + def test_disk_cache_under_dataset_git_dir(self, tmp_path) -> None: + backend, fake = _mock_remfile_backend(tmp_path, caching=True) + expected = str(tmp_path / ".git" / "datalad" / "cache" / "remfile") + fake.DiskCache.assert_called_once_with(expected) + assert backend._disk_cache is fake.DiskCache.return_value + + def test_disk_cache_passed_to_remfile(self, tmp_path, monkeypatch) -> None: + backend, fake = _mock_remfile_backend(tmp_path, caching=True) + resp = _probe_response({"Content-Range": "bytes 0-0/10"}) + monkeypatch.setattr("urllib.request.urlopen", MagicMock(return_value=resp)) + backend.open_url("http://example.com/x.h5") + assert fake.File.call_args.kwargs["disk_cache"] is fake.DiskCache.return_value + + def test_clear_removes_cache_dir(self, tmp_path) -> None: + """So ``fsspec-cache-clear`` / ``datalad.fusefs.cache-clear`` work.""" + backend, _ = _mock_remfile_backend(tmp_path, caching=True) + cache_dir = tmp_path / ".git" / "datalad" / "cache" / "remfile" + (cache_dir / "ab" / "cd").mkdir(parents=True) + (cache_dir / "ab" / "cd" / "chunk").write_bytes(b"x") + backend.clear() + assert not cache_dir.exists() + + def test_clear_without_caching_is_noop(self, tmp_path) -> None: + backend, _ = _mock_remfile_backend(tmp_path, caching=False) + cache_dir = tmp_path / ".git" / "datalad" / "cache" / "remfile" + cache_dir.mkdir(parents=True) + backend.clear() + assert cache_dir.exists() + + +class TestFsspecBackendCanHandle: + """FsspecBackend.can_handle always returns True.""" + + @pytest.mark.ai_generated + def test_always_true(self, tmp_path) -> None: + backend = FsspecBackend(tmp_path, caching=False) + key = AnnexKey(backend="MD5E", name="abc123", size=100, suffix=".txt") + assert backend.can_handle(key, "rb") is True + assert backend.can_handle(key, "r") is True + assert backend.can_handle(None, "rb") is True + assert backend.can_handle(None, "rb", "x.nwb") is True + + +# -- ABC compliance shared across backends ----------------------------------- + + +def _all_backends(tmp_path) -> list[Backend]: + out: list[Backend] = [FsspecBackend(tmp_path, caching=False)] + if _has_remfile: + out.append(RemfileBackend(tmp_path, caching=False)) + return out + + +@pytest.mark.ai_generated +def test_abc_compliance(tmp_path) -> None: + for backend in _all_backends(tmp_path): + assert isinstance(backend, Backend) + assert isinstance(backend.name, str) and backend.name + assert callable(backend.can_handle) + assert callable(backend.open_url) + assert callable(backend.clear) + + +# -- FsspecBackend retry on BlocksizeMismatchError --------------------------- + + +@pytest.mark.ai_generated +def test_fsspec_open_retries_on_blocksize_mismatch(tmp_path) -> None: + """First open() raises BlocksizeMismatchError; cache cleared; retry succeeds.""" + backend = FsspecBackend(tmp_path, caching=True) + fake_fs = MagicMock() + good_handle = BytesIO(b"ok") + fake_fs.open.side_effect = [BlocksizeMismatchError("mismatch"), good_handle] + backend.fs = fake_fs + result = backend.open_url("http://example.com/x.bin") + assert result is good_handle + fake_fs.pop_from_cache.assert_called_once_with("http://example.com/x.bin") + assert fake_fs.open.call_count == 2 + + +@pytest.mark.ai_generated +def test_fsspec_open_propagates_blocksize_mismatch_without_caching(tmp_path) -> None: + """With caching=False there is no cache to evict; the error must propagate + instead of attempting pop_from_cache (which would AttributeError on a + plain HTTPFileSystem).""" + backend = FsspecBackend(tmp_path, caching=False) + fake_fs = MagicMock(spec=["open"]) # no pop_from_cache attribute + fake_fs.open.side_effect = BlocksizeMismatchError("mismatch") + backend.fs = fake_fs + with pytest.raises(BlocksizeMismatchError): + backend.open_url("http://example.com/x.bin") + assert fake_fs.open.call_count == 1 + + +# -- resolve_backends / create_backends tests -------------------------------- + + +class TestResolveBackends: + """Verify argument > config > default precedence. + + Note: resolve_backends returns ``(spec, explicit)``. ``explicit`` is True + for the argument and config paths, False for the default fallback. + """ + + pytestmark = pytest.mark.ai_generated + + def test_explicit_wins(self) -> None: + spec, explicit = resolve_backends("fsspec") + assert spec == "fsspec" + assert explicit is True + + def test_default(self, monkeypatch) -> None: + # Ensure config is not consulted; this returns None for any key so + # resolve_backends falls through to DEFAULT_BACKENDS. + monkeypatch.setattr( + "datalad_fuse.adapter.cfg.get", + lambda _key, default=None: default, + ) + spec, explicit = resolve_backends(None) + assert spec == DEFAULT_BACKENDS + assert explicit is False + + def test_config_override(self, monkeypatch) -> None: + get = MagicMock(return_value="fsspec") + monkeypatch.setattr("datalad_fuse.adapter.cfg.get", get) + spec, explicit = resolve_backends(None) + assert spec == "fsspec" + assert explicit is True + get.assert_called_once_with("datalad.fusefs.backends", None) + + def test_explicit_beats_config(self, monkeypatch) -> None: + get = MagicMock(return_value="fsspec") + monkeypatch.setattr("datalad_fuse.adapter.cfg.get", get) + spec, explicit = resolve_backends("remfile,fsspec") + assert spec == "remfile,fsspec" + assert explicit is True + get.assert_not_called() + + def test_dataset_config_consulted(self, monkeypatch) -> None: + """A dataset's own config must be honored, not just the global one.""" + global_get = MagicMock(return_value=None) + monkeypatch.setattr("datalad_fuse.adapter.cfg.get", global_get) + ds_config = MagicMock() + ds_config.get.return_value = "fsspec" + spec, explicit = resolve_backends(None, config=ds_config) + assert spec == "fsspec" + assert explicit is True + ds_config.get.assert_called_once_with("datalad.fusefs.backends", None) + global_get.assert_not_called() + + def test_dataset_config_unset_falls_back_to_default(self, monkeypatch) -> None: + monkeypatch.setattr( + "datalad_fuse.adapter.cfg.get", + lambda _key, default=None: default, + ) + ds_config = MagicMock() + ds_config.get.return_value = None + spec, explicit = resolve_backends(None, config=ds_config) + assert spec == DEFAULT_BACKENDS + assert explicit is False + + +class TestCreateBackends: + pytestmark = pytest.mark.ai_generated + + def test_fsspec_only(self, tmp_path) -> None: + backends = create_backends("fsspec", tmp_path, caching=False) + assert len(backends) == 1 + assert backends[0].name == "fsspec" + + def test_unknown_backend_raises(self, tmp_path) -> None: + with pytest.raises(ValueError, match="Unknown backend"): + create_backends("nosuch", tmp_path, caching=False) + + @pytest.mark.parametrize( + "spec", + [ + "fsspec,,fsspec", # double comma + "fsspec, ,fsspec", # comma + whitespace + "fsspec,", # trailing comma + ",fsspec", # leading comma + " fsspec , fsspec ", # internal/external whitespace + ], + ) + def test_empty_entries_tolerated(self, tmp_path, spec) -> None: + """Stray commas and whitespace in the spec are silently skipped.""" + backends = create_backends(spec, tmp_path, caching=False) + assert all(b.name == "fsspec" for b in backends) + assert len(backends) >= 1 + + def test_unavailable_backend_skipped(self, tmp_path, caplog) -> None: + with patch("datalad_fuse.remfile._get_remfile", return_value=None): + backends = create_backends( + "remfile,fsspec", tmp_path, caching=False, explicit=False + ) + # remfile skipped, fsspec remains + assert len(backends) == 1 + assert backends[0].name == "fsspec" + # No warning when non-explicit + assert not any(r.levelname == "WARNING" for r in caplog.records) + + def test_unavailable_explicit_backend_warns(self, tmp_path, caplog) -> None: + import logging + + caplog.set_level(logging.WARNING, logger="datalad.fuse.adapter") + with patch("datalad_fuse.remfile._get_remfile", return_value=None): + backends = create_backends( + "remfile,fsspec", tmp_path, caching=False, explicit=True + ) + assert len(backends) == 1 and backends[0].name == "fsspec" + assert any("remfile" in r.message for r in caplog.records) + + def test_all_unavailable_raises(self, tmp_path) -> None: + with patch("datalad_fuse.remfile._get_remfile", return_value=None): + with pytest.raises(ValueError, match="No usable backends"): + create_backends("remfile", tmp_path, caching=False) + + @requires_remfile + def test_remfile_gets_path_and_caching(self, tmp_path) -> None: + """create_backends must wire --caching through to remfile too.""" + (backend,) = create_backends("remfile", tmp_path, caching=True) + assert backend._disk_cache is not None + + @requires_remfile + def test_order_preserved(self, tmp_path) -> None: + backends = create_backends("remfile,fsspec", tmp_path, caching=False) + assert [b.name for b in backends] == ["remfile", "fsspec"] + + def test_fsspec_with_caching(self, tmp_path) -> None: + backends = create_backends("fsspec", tmp_path, caching=True) + assert len(backends) == 1 + assert backends[0].name == "fsspec" + from fsspec.implementations.cached import CachingFileSystem + + assert isinstance(backends[0].fs, CachingFileSystem) + + +# -- is_http_url tests ------------------------------------------------------- + + +@pytest.mark.ai_generated +@pytest.mark.parametrize( + "url,expected", + [ + ("http://example.com/x", True), + ("https://example.com/x", True), + ("HTTP://example.com/x", True), + ("HTTPS://example.com/x", True), + ("ftp://example.com/x", False), + ("file:///tmp/x", False), + ("", False), + ("/local/path", False), + ], +) +def test_is_http_url(url: str, expected: bool) -> None: + assert is_http_url(url) is expected + + +# -- RemfileWrapper tests ---------------------------------------------------- + + +def _make_mock_remfile(data: bytes = b"hello world\nline two\n") -> MagicMock: + """Create a mock ``remfile.File`` object backed by *data*. + + Deliberately mimics ``remfile.RemFile``'s quirks rather than a well-behaved + file object, so that the wrapper is tested against what it actually wraps: + + - ``read()`` requires an explicit, non-negative size, + - it advances the position by that size even past EOF, and + - a read starting at or past EOF issues an unsatisfiable Range request, + - ``seek()`` returns ``None``. + """ + mock = MagicMock() + pos = [0] + + def read(size=None) -> bytes: + if size is None: + raise Exception("The size argument must be provided in remfile") + start = pos[0] + pos[0] = start + size # remfile advances unconditionally + if size < 0: + # remfile's chunk range comes out empty for a negative size + return b"" + if start >= len(data) and size: + raise OSError( + f"Error fetching bytes {start}-{start + size - 1}: " + "416 Range Not Satisfiable" + ) + return data[start : start + size] + + def seek(offset: int, whence: int = 0) -> None: + if whence == 0: + pos[0] = offset + elif whence == 1: + pos[0] += offset + elif whence == 2: + pos[0] = len(data) + offset + else: + raise ValueError("Invalid argument: 'whence' must be 0, 1, or 2.") + + def tell() -> int: + return pos[0] + + mock.length = len(data) + mock.read = read + mock.seek = seek + mock.tell = tell + mock.close = MagicMock() + return mock + + +class TestRemfileWrapper: + """Tests for the RemfileWrapper adapter class.""" + + pytestmark = pytest.mark.ai_generated + + def test_read(self) -> None: + mock_rf = _make_mock_remfile(b"abcdef") + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + assert w.read(3) == b"abc" + assert w.read(3) == b"def" + + def test_seek_and_tell(self) -> None: + mock_rf = _make_mock_remfile(b"abcdef") + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + w.seek(3) + assert w.tell() == 3 + assert w.read(2) == b"de" + + def test_read_without_size_reads_to_eof(self) -> None: + """``read()`` with no size must behave like a file object, not remfile. + + ``remfile.RemFile.read(-1)`` returns ``b""`` and moves the position + back by one. + """ + mock_rf = _make_mock_remfile(b"abcdef") + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + w.seek(2) + assert w.read() == b"cdef" + assert w.tell() == 6 + assert w.read() == b"" + + def test_read_none_size_reads_to_eof(self) -> None: + mock_rf = _make_mock_remfile(b"abcdef") + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + assert w.read(None) == b"abcdef" + + def test_read_clamps_at_eof(self) -> None: + """An over-long read returns what is there and leaves tell() at EOF.""" + mock_rf = _make_mock_remfile(b"abcdef") + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + w.seek(4) + assert w.read(100) == b"ef" + assert w.tell() == 6 + + def test_read_at_eof_returns_empty(self) -> None: + """Reading at EOF must not issue an unsatisfiable Range request. + + remfile raises ``416 Range Not Satisfiable`` for that, which breaks + line iteration (and hence ``fsspec-head -n``) whenever the file size is + an exact multiple of remfile's chunk size. + """ + mock_rf = _make_mock_remfile(b"abcdef") + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + w.seek(6) + assert w.read(10) == b"" + assert w.tell() == 6 + + def test_seek_returns_new_position(self) -> None: + """``remfile.RemFile.seek()`` returns None; the wrapper declares -> int.""" + mock_rf = _make_mock_remfile(b"abcdef") + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + assert w.seek(3) == 3 + assert w.seek(-1, 2) == 5 + + def test_context_manager(self) -> None: + mock_rf = _make_mock_remfile() + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + assert not w.closed + with w as f: + assert f is w + assert w.closed + mock_rf.close.assert_called_once() + + def test_iteration(self) -> None: + mock_rf = _make_mock_remfile(b"line one\nline two\nline three") + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + assert list(w) == [b"line one\n", b"line two\n", b"line three"] + + def test_iteration_empty(self) -> None: + mock_rf = _make_mock_remfile(b"") + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + assert list(w) == [] + + def test_iteration_caps_no_newline(self) -> None: + """Binary data without '\\n' must not cause runaway reads. + + The wrapper caps each ``__next__`` at ``_MAX_LINE_BYTES`` so iterating + an HDF5 file does not download the whole remote object. + """ + # 3 chunks worth of newline-less bytes (well below 1 MiB cap so the + # test stays fast); shrink the cap for the assertion. + data = b"x" * (3 * RemfileWrapper._ITER_CHUNK) + mock_rf = _make_mock_remfile(data) + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + # Patch the cap to 2 chunks so we can verify the bound takes effect. + w._MAX_LINE_BYTES = 2 * RemfileWrapper._ITER_CHUNK + first = next(iter(w)) + assert len(first) == 2 * RemfileWrapper._ITER_CHUNK + assert b"\n" not in first + + def test_close(self) -> None: + mock_rf = _make_mock_remfile() + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + w.close() + assert w.closed + mock_rf.close.assert_called_once() + + def test_io_protocol_methods(self) -> None: + mock_rf = _make_mock_remfile() + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + assert w.readable() is True + assert w.seekable() is True + assert w.writable() is False + + def test_info_uses_remfile_length(self) -> None: + """info() reports the size remfile already knows, with no HEAD request. + + A HEAD here would break ``fstat()`` on an open FUSE handle whenever the + server rejects HEAD, which presigned S3 GET URLs do. + """ + mock_rf = _make_mock_remfile(b"0123456789") + w = RemfileWrapper(mock_rf, "http://example.com/test.h5") + with patch("urllib.request.urlopen") as urlopen: + assert w.info() == {"type": "file", "size": 10} + urlopen.assert_not_called() + + +# -- DatasetAdapter backend-chain and RemoteFilesystemAdapter tests ---------- + + +class _StubBackend(Backend): + """In-memory backend for exercising DatasetAdapter.open fallback.""" + + def __init__( + self, + name: str, + can_handle_result: bool = True, + open_result: Optional[bytes] = None, + raises: Optional[Exception] = None, + ) -> None: + self.name = name + self._can = can_handle_result + self._open_result = open_result + self._raises = raises + self.open_calls: list[str] = [] + self.clear_calls = 0 + + def can_handle(self, key, mode: str, relpath=None) -> bool: # noqa: U100 + return self._can + + def open_url(self, url: str, mode: str = "rb", **kwargs): # noqa: U100 + self.open_calls.append(url) + if self._raises is not None: + raise self._raises + return BytesIO(self._open_result or b"") + + def clear(self) -> None: + self.clear_calls += 1 + + +def _make_dataset_adapter_with_backends( + tmp_path, backends: list[Backend], urls: list[str] +) -> DatasetAdapter: + """Build a DatasetAdapter whose file-state + URL generator are stubbed.""" + # Bypass __init__ (which requires a real datalad dataset on disk) + adapter = DatasetAdapter.__new__(DatasetAdapter) + adapter.path = tmp_path + adapter.mode_transparent = False + adapter.annex = None + adapter._backends = backends + + key = AnnexKey(backend="MD5E", name="abc123", size=100, suffix=".nwb") + + def fake_get_file_state(_relpath: str): + return (FileState.NO_CONTENT, key) + + def fake_get_urls(_key: str) -> Iterator[str]: + yield from urls + + adapter.get_file_state = fake_get_file_state # type: ignore[method-assign] + adapter.get_urls = fake_get_urls # type: ignore[method-assign] + return adapter + + +@pytest.mark.ai_generated +class TestDatasetAdapterOpen: + """Exercise the backend-chain fallback semantics in DatasetAdapter.open.""" + + def test_first_backend_succeeds(self, tmp_path) -> None: + a = _StubBackend("a", open_result=b"hello") + b = _StubBackend("b", open_result=b"world") + adapter = _make_dataset_adapter_with_backends(tmp_path, [a, b], ["http://u1"]) + handle = adapter.open("x.nwb") + assert handle.read() == b"hello" + assert a.open_calls == ["http://u1"] + assert b.open_calls == [] + + def test_falls_through_to_next_backend(self, tmp_path) -> None: + boom = _StubBackend("boom", raises=FileNotFoundError("nope")) + ok = _StubBackend("ok", open_result=b"good") + adapter = _make_dataset_adapter_with_backends( + tmp_path, [boom, ok], ["http://u1"] + ) + handle = adapter.open("x.nwb") + assert handle.read() == b"good" + assert boom.open_calls == ["http://u1"] + assert ok.open_calls == ["http://u1"] + + def test_skips_backend_that_cannot_handle(self, tmp_path) -> None: + skipper = _StubBackend( + "skip", can_handle_result=False, raises=RuntimeError("!") + ) + ok = _StubBackend("ok", open_result=b"good") + adapter = _make_dataset_adapter_with_backends( + tmp_path, [skipper, ok], ["http://u1"] + ) + handle = adapter.open("x.nwb") + assert handle.read() == b"good" + # Skipper must not have been asked to open anything. + assert skipper.open_calls == [] + assert ok.open_calls == ["http://u1"] + + def test_all_fail_raises_ioerror_with_cause(self, tmp_path) -> None: + last_exc = FileNotFoundError("last") + a = _StubBackend("a", raises=FileNotFoundError("first")) + b = _StubBackend("b", raises=last_exc) + adapter = _make_dataset_adapter_with_backends(tmp_path, [a, b], ["http://u1"]) + with pytest.raises(IOError) as exc_info: + adapter.open("x.nwb") + assert exc_info.value.__cause__ is last_exc + assert "a,b" in str(exc_info.value) + + def test_non_http_urls_still_tried_per_backend(self, tmp_path) -> None: + """All provided URLs are tried in order until one succeeds.""" + a = _StubBackend( + "a", + raises=FileNotFoundError("no"), + ) + # With multiple URLs, backend is retried for each. + adapter = _make_dataset_adapter_with_backends( + tmp_path, [a], ["http://u1", "http://u2"] + ) + with pytest.raises(IOError): + adapter.open("x.nwb") + assert a.open_calls == ["http://u1", "http://u2"] + + def test_urls_enumerated_once_across_backends(self, tmp_path) -> None: + """Falling through to the next backend must not re-enumerate URLs. + + ``get_urls()`` runs a ``git annex whereis`` subprocess plus two + ``examinekey`` calls, and the exporttree path adds a boto3 + ``list_object_versions`` call, so repeating it per backend is costly. + """ + a = _StubBackend("a", raises=FileNotFoundError("no")) + b = _StubBackend("b", raises=FileNotFoundError("no")) + adapter = _make_dataset_adapter_with_backends( + tmp_path, [a, b], ["http://u1", "http://u2"] + ) + calls: list[str] = [] + inner = adapter.get_urls + + def counting_get_urls(key: str) -> Iterator[str]: + calls.append(key) + yield from inner(key) + + adapter.get_urls = counting_get_urls # type: ignore[method-assign] + with pytest.raises(IOError): + adapter.open("x.nwb") + assert len(calls) == 1 + # Every backend still gets offered every URL. + assert a.open_calls == ["http://u1", "http://u2"] + assert b.open_calls == ["http://u1", "http://u2"] + + def test_urls_not_enumerated_when_no_backend_handles(self, tmp_path) -> None: + """URL enumeration stays lazy when nothing can handle the file.""" + skipper = _StubBackend("skip", can_handle_result=False) + adapter = _make_dataset_adapter_with_backends( + tmp_path, [skipper], ["http://u1"] + ) + calls: list[str] = [] + inner = adapter.get_urls + + def counting_get_urls(key: str) -> Iterator[str]: + calls.append(key) + yield from inner(key) + + adapter.get_urls = counting_get_urls # type: ignore[method-assign] + with pytest.raises(IOError): + adapter.open("x.nwb") + assert calls == [] + + def test_clear_delegates_to_backends(self, tmp_path) -> None: + a = _StubBackend("a") + b = _StubBackend("b") + adapter = _make_dataset_adapter_with_backends(tmp_path, [a, b], []) + adapter.clear() + assert a.clear_calls == 1 + assert b.clear_calls == 1 + + def test_unsupported_mode_raises(self, tmp_path) -> None: + adapter = _make_dataset_adapter_with_backends(tmp_path, [_StubBackend("a")], []) + with pytest.raises(NotImplementedError, match="modes"): + adapter.open("x.nwb", mode="wb") + + @pytest.mark.parametrize("mode", ["r", "rt"]) + def test_text_mode_forwards_encoding(self, tmp_path, mode) -> None: + """Mode 'r'/'rt' should pass encoding/errors through to backend.open_url.""" + captured: dict = {} + + class CapturingBackend(Backend): + name = "cap" + + def can_handle(self, key, mode_, relpath=None): # noqa: U100 + return True + + def open_url(self, url, mode_="rb", **kwargs): # noqa: U100 + captured.update(kwargs) + return BytesIO(b"") + + def clear(self): + pass + + adapter = _make_dataset_adapter_with_backends( + tmp_path, [CapturingBackend()], ["http://u1"] + ) + adapter.open("x.nwb", mode=mode, encoding="latin-1", errors="replace") + assert captured == {"encoding": "latin-1", "errors": "replace"} + + def test_generic_exception_falls_through(self, tmp_path) -> None: + """A non-FileNotFoundError (e.g. RuntimeError) should still fall through.""" + boom = _StubBackend("boom", raises=RuntimeError("wat")) + ok = _StubBackend("ok", open_result=b"fallback") + adapter = _make_dataset_adapter_with_backends( + tmp_path, [boom, ok], ["http://u1"] + ) + handle = adapter.open("x.nwb") + assert handle.read() == b"fallback" + assert boom.open_calls == ["http://u1"] + assert ok.open_calls == ["http://u1"] + + def test_not_annexed_opens_local(self, tmp_path) -> None: + """NOT_ANNEXED files are opened directly from disk.""" + (tmp_path / "local.txt").write_text("direct read") + adapter = DatasetAdapter.__new__(DatasetAdapter) + adapter.path = tmp_path + adapter.mode_transparent = False + adapter.annex = None + adapter._backends = [] + adapter.get_file_state = lambda _r: ( # type: ignore[method-assign] + FileState.NOT_ANNEXED, + None, + ) + with adapter.open("local.txt", mode="r") as f: + assert f.read() == "direct read" + + +@pytest.mark.ai_generated +class TestRemoteFilesystemAdapterLifecycle: + """Context-manager semantics: __exit__ closes datasets and empties the dict.""" + + def test_exit_closes_datasets(self, tmp_path) -> None: + rfs = RemoteFilesystemAdapter(tmp_path, caching=False) + ds_a = MagicMock() + ds_b = MagicMock() + rfs.datasets = {tmp_path / "a": ds_a, tmp_path / "b": ds_b} + with rfs: + pass + ds_a.close.assert_called_once() + ds_b.close.assert_called_once() + assert rfs.datasets == {} + + def test_exit_closes_even_on_exception(self, tmp_path) -> None: + rfs = RemoteFilesystemAdapter(tmp_path, caching=False) + ds_a = MagicMock() + rfs.datasets = {tmp_path / "a": ds_a} + with pytest.raises(RuntimeError): + with rfs: + raise RuntimeError("boom") + ds_a.close.assert_called_once() + assert rfs.datasets == {} + + def test_get_dataset_path_not_datalad(self, tmp_path, monkeypatch) -> None: + monkeypatch.setattr("datalad_fuse.adapter.get_dataset_root", lambda _p: None) + rfs = RemoteFilesystemAdapter(tmp_path, caching=False) + with pytest.raises(ValueError, match="not under DataLad"): + rfs.get_dataset_path("anything") + + def test_get_dataset_path_outside_root(self, tmp_path, monkeypatch) -> None: + other = tmp_path.parent / "elsewhere" + monkeypatch.setattr( + "datalad_fuse.adapter.get_dataset_root", lambda _p: str(other) + ) + rfs = RemoteFilesystemAdapter(tmp_path, caching=False) + with pytest.raises(ValueError, match="not under root dataset"): + rfs.get_dataset_path("x") + + def test_delegation_to_dataset_adapter(self, tmp_path) -> None: + """get_file_state, is_under_annex, get_commit_datetime delegate to DS.""" + rfs = RemoteFilesystemAdapter(tmp_path, caching=False) + fake_ds = MagicMock() + dt = datetime(2024, 1, 2, tzinfo=timezone.utc) + fake_ds.commit_dt = dt + fake_ds.get_file_state.return_value = (FileState.HAS_CONTENT, None) + rfs.resolve_dataset = lambda _p: (fake_ds, "rel") # type: ignore[method-assign] + + assert rfs.get_file_state("x") == (FileState.HAS_CONTENT, None) + assert rfs.is_under_annex("x") is True + assert rfs.get_commit_datetime("x") == dt + + fake_ds.get_file_state.return_value = (FileState.NOT_ANNEXED, None) + assert rfs.is_under_annex("x") is False + + def test_resolve_dataset_reuses_cached_adapter(self, tmp_path, monkeypatch) -> None: + """resolve_dataset caches DatasetAdapter instances by dataset path.""" + monkeypatch.setattr( + "datalad_fuse.adapter.get_dataset_root", lambda p: str(Path(p).parent) + ) + rfs = RemoteFilesystemAdapter(tmp_path, caching=False) + # Pre-populate so __init__ isn't invoked (no real dataset) + prebuilt = MagicMock() + dspath = tmp_path + rfs.datasets[dspath] = prebuilt + (tmp_path / "a.txt").write_text("x") + dsap1, rel1 = rfs.resolve_dataset(tmp_path / "a.txt") + dsap2, _ = rfs.resolve_dataset(tmp_path / "a.txt") + assert dsap1 is prebuilt is dsap2 + assert rel1 == "a.txt" + + +# -- FsspecBackend additional coverage --------------------------------------- + + +@pytest.mark.ai_generated +def test_fsspec_clear_with_caching(tmp_path) -> None: + """FsspecBackend.clear() calls fs.clear_cache() when caching=True.""" + backend = FsspecBackend(tmp_path, caching=True) + backend.fs = MagicMock() + backend._caching = True + backend.clear() + backend.fs.clear_cache.assert_called_once() + + +@pytest.mark.ai_generated +def test_fsspec_clear_without_caching_noop(tmp_path) -> None: + """FsspecBackend.clear() is a no-op when caching=False.""" + backend = FsspecBackend(tmp_path, caching=False) + backend.fs = MagicMock() + backend.clear() + backend.fs.clear_cache.assert_not_called() + + +@pytest.mark.ai_generated +def test_on_request_start_logs_retry(caplog) -> None: + """on_request_start emits a warning for attempts beyond the first.""" + caplog.set_level(logging.WARNING, logger="datalad.fuse.fsspec") + ctx = SimpleNamespace(trace_request_ctx={"current_attempt": 2}) + params = SimpleNamespace(url="http://example.com") + asyncio.run(on_request_start(MagicMock(), ctx, params)) + assert any("Retrying" in r.message for r in caplog.records) + + caplog.clear() + ctx.trace_request_ctx["current_attempt"] = 1 + asyncio.run(on_request_start(MagicMock(), ctx, params)) + assert not any("Retrying" in r.message for r in caplog.records) + + +# -- RemfileBackend availability detection ----------------------------------- + + +@pytest.mark.ai_generated +def test_remfile_backend_raises_when_unavailable(monkeypatch, tmp_path) -> None: + """RemfileBackend.__init__ raises ImportError when remfile is not installed.""" + monkeypatch.setattr("datalad_fuse.remfile._get_remfile", lambda: None) + with pytest.raises(ImportError, match="remfile"): + RemfileBackend(tmp_path, caching=False) + + +@pytest.mark.ai_generated +def test_get_remfile_returns_none_on_import_error(monkeypatch) -> None: + """_get_remfile() swallows ImportError and returns None.""" + # Block remfile from being imported: shadow sys.modules so the + # next `import remfile` fails, and shadow __import__ to guarantee + # ImportError is raised even if remfile was cached. + monkeypatch.setitem(sys.modules, "remfile", None) + real_import = builtins.__import__ + + def fake_import(name, *args, **kwargs): + if name == "remfile": + raise ImportError("not here") + return real_import(name, *args, **kwargs) + + monkeypatch.setattr(builtins, "__import__", fake_import) + assert _get_remfile() is None + + +# -- Integration tests with real S3 URLs ------------------------------------- + +# Pinned versions from dandiarchive S3 bucket +S3_HDF5_URL = ( + "https://dandiarchive.s3.amazonaws.com/ros3test.hdf5" + "?versionId=_8Zs6qF7E6vpc5BcPOhizlFBY2oHCN8T" +) +S3_NWB_URL = ( + "https://dandiarchive.s3.amazonaws.com/ros3test.nwb" + "?versionId=jRN4ejAcjOAaFDXTQrO99WCZqTxhZL32" +) +S3_README_URL = ( + "https://dandiarchive.s3.amazonaws.com/README.md" + "?versionId=mAywZ4KP9BCgGIERb3DtlPzTWYv.5sUi" +) + + +@requires_remfile +class TestIntegrationRemfileBackend: + """Integration tests using real S3 URLs via the RemfileBackend.""" + + pytestmark = [pytest.mark.ai_generated, pytest.mark.network] + + @pytest.fixture() + def remfile_backend(self, tmp_path) -> RemfileBackend: + return RemfileBackend(tmp_path, caching=False) + + def test_open_hdf5(self, remfile_backend: RemfileBackend) -> None: + with remfile_backend.open_url(S3_HDF5_URL) as f: + assert f.read(8) == b"\x89HDF\r\n\x1a\n" + + def test_open_nwb(self, remfile_backend: RemfileBackend) -> None: + with remfile_backend.open_url(S3_NWB_URL) as f: + assert f.read(8) == b"\x89HDF\r\n\x1a\n" + + def test_seek_and_reread(self, remfile_backend: RemfileBackend) -> None: + with remfile_backend.open_url(S3_HDF5_URL) as f: + first = f.read(8) + f.seek(0) + second = f.read(8) + assert first == second == b"\x89HDF\r\n\x1a\n" + + def test_wrapper_context_manager(self, remfile_backend: RemfileBackend) -> None: + f = remfile_backend.open_url(S3_HDF5_URL) + with f: + assert f.read(4) == b"\x89HDF" + assert f.closed + + +class TestIntegrationFsspecBackend: + """Integration tests using real S3 URLs via the FsspecBackend.""" + + pytestmark = [pytest.mark.ai_generated, pytest.mark.network] + + @pytest.fixture() + def fsspec_backend(self, tmp_path) -> FsspecBackend: + return FsspecBackend(tmp_path, caching=False) + + def test_open_readme(self, fsspec_backend: FsspecBackend) -> None: + with fsspec_backend.open_url(S3_README_URL) as f: + data = f.read(100) + assert len(data) > 0 + assert isinstance(data, bytes) + + def test_open_hdf5(self, fsspec_backend: FsspecBackend) -> None: + with fsspec_backend.open_url(S3_HDF5_URL) as f: + assert f.read(8) == b"\x89HDF\r\n\x1a\n" + + +@requires_remfile +class TestIntegrationBackendChain: + """Test backend chain selection logic.""" + + pytestmark = pytest.mark.ai_generated + + def test_remfile_selected_for_hdf5(self, tmp_path) -> None: + backends = create_backends("remfile,fsspec", tmp_path, caching=False) + key = AnnexKey(backend="MD5E", name="abc", size=100, suffix=".hdf5") + for b in backends: + if b.can_handle(key, "rb"): + assert b.name == "remfile" + break + + def test_fsspec_selected_for_md(self, tmp_path) -> None: + backends = create_backends("remfile,fsspec", tmp_path, caching=False) + key = AnnexKey(backend="MD5E", name="abc", size=100, suffix=".md") + for b in backends: + if b.can_handle(key, "rb"): + assert b.name == "fsspec" + break + + def test_fsspec_only_chain(self, tmp_path) -> None: + backends = create_backends("fsspec", tmp_path, caching=False) + key = AnnexKey(backend="MD5E", name="abc", size=100, suffix=".nwb") + for b in backends: + if b.can_handle(key, "rb"): + assert b.name == "fsspec" + break diff --git a/datalad_fuse/tests/test_deprecations.py b/datalad_fuse/tests/test_deprecations.py new file mode 100644 index 0000000..3b074dc --- /dev/null +++ b/datalad_fuse/tests/test_deprecations.py @@ -0,0 +1,55 @@ +"""Verify old datalad_fuse.fsspec import paths still work with warnings.""" + +from __future__ import annotations + +import importlib +import warnings + +import pytest + +from datalad_fuse.fsspec import _COMPAT_MAP + + +@pytest.mark.ai_generated +@pytest.mark.parametrize( + "old_name,module_path,canonical_name", + [(name, mod, canon) for name, (mod, canon) in _COMPAT_MAP.items()], + ids=list(_COMPAT_MAP), +) +def test_compat_import_warns(old_name, module_path, canonical_name): + """Moved names imported from datalad_fuse.fsspec emit DeprecationWarning.""" + import datalad_fuse.fsspec as fsspec_mod + + with pytest.warns(DeprecationWarning, match=old_name): + obj = getattr(fsspec_mod, old_name) + + canonical_mod = importlib.import_module(module_path) + assert obj is getattr(canonical_mod, canonical_name) + + +@pytest.mark.ai_generated +def test_unknown_attr_raises(): + import datalad_fuse.fsspec as fsspec_mod + + # Direct attribute access raises AttributeError (module-level __getattr__). + with warnings.catch_warnings(): + warnings.simplefilter("error", DeprecationWarning) + with pytest.raises(AttributeError, match="no_such_name"): + fsspec_mod.no_such_name # noqa: B018 + + +@pytest.mark.ai_generated +def test_unknown_attr_from_import_raises(): + # `from x import y` on a missing name converts AttributeError to ImportError. + with pytest.raises(ImportError): + from datalad_fuse.fsspec import ( # type: ignore[attr-defined] # noqa: F401 + no_such_name, + ) + + +@pytest.mark.ai_generated +def test_fsspec_backend_no_warning(): + """FsspecBackend still lives in fsspec.py — no deprecation.""" + with warnings.catch_warnings(): + warnings.simplefilter("error", DeprecationWarning) + from datalad_fuse.fsspec import FsspecBackend # noqa: F401 diff --git a/datalad_fuse/tests/test_exporttree.py b/datalad_fuse/tests/test_exporttree.py index f962094..75239b4 100644 --- a/datalad_fuse/tests/test_exporttree.py +++ b/datalad_fuse/tests/test_exporttree.py @@ -13,10 +13,9 @@ import pytest -from datalad_fuse.fsspec import DatasetAdapter +from datalad_fuse.adapter import DatasetAdapter from datalad_fuse.utils import AnnexKey - # --- remote.log parsing --- @@ -261,9 +260,7 @@ def test_get_exporttree_urls_construction(adapter): } ], ), - patch.object( - DatasetAdapter, "_list_s3_versions", return_value=versions - ), + patch.object(DatasetAdapter, "_list_s3_versions", return_value=versions), ): urls = list(adapter.get_exporttree_urls("sub-01/anat/sub-01_T1w.nii.gz", key)) @@ -338,9 +335,7 @@ def test_get_exporttree_urls_ambiguous_skips(adapter): } ], ), - patch.object( - DatasetAdapter, "_list_s3_versions", return_value=versions - ), + patch.object(DatasetAdapter, "_list_s3_versions", return_value=versions), ): urls = list(adapter.get_exporttree_urls("sub-01/anat/sub-01_T1w.nii.gz", key)) diff --git a/datalad_fuse/tests/test_forgejo.py b/datalad_fuse/tests/test_forgejo.py index c9a4b11..1b63d8d 100644 --- a/datalad_fuse/tests/test_forgejo.py +++ b/datalad_fuse/tests/test_forgejo.py @@ -12,7 +12,7 @@ import pytest import requests -from datalad_fuse.fsspec import DatasetAdapter, _is_aneksajo +from datalad_fuse.adapter import DatasetAdapter, _is_aneksajo from .conftest_forgejo import ( ForgejoInstance, diff --git a/docs/source/api.rst b/docs/source/api.rst index 8a63087..fee7e00 100644 --- a/docs/source/api.rst +++ b/docs/source/api.rst @@ -18,15 +18,15 @@ The commands added to DataLad by ``datalad-fuse``, available from fsspec_cache_clear -Opening files: ``datalad_fuse.fsspec`` -====================================== +Opening files: ``datalad_fuse.adapter`` +======================================= -.. currentmodule:: datalad_fuse.fsspec +.. currentmodule:: datalad_fuse.adapter .. autoclass:: DatasetAdapter :members: open, get_file_state, get_urls, clear, close -.. autoclass:: FsspecAdapter +.. autoclass:: RemoteFilesystemAdapter :members: open, get_file_state, is_under_annex, get_commit_datetime, resolve_dataset, get_dataset_path @@ -34,6 +34,43 @@ Opening files: ``datalad_fuse.fsspec`` :members: :undoc-members: +.. autofunction:: resolve_backends + +.. autofunction:: create_backends + +.. note:: + Up to 0.6.0 these lived in ``datalad_fuse.fsspec``, and + `RemoteFilesystemAdapter` was called ``FsspecAdapter``. Importing + ``FsspecAdapter``, ``DatasetAdapter``, ``FileState`` or ``is_http_url`` + from ``datalad_fuse.fsspec`` still works, with a + :exc:`DeprecationWarning`. + + +Backends: ``datalad_fuse.backends`` +=================================== + +Backends do the actual reading of remote files; see :ref:`concepts-backends`. + +.. currentmodule:: datalad_fuse.backends + +.. autoclass:: Backend + :members: can_handle, open_url, clear + +.. autodata:: DEFAULT_BACKENDS + +.. currentmodule:: datalad_fuse.fsspec + +.. autoclass:: FsspecBackend + :members: can_handle, open_url, clear + +.. currentmodule:: datalad_fuse.remfile + +.. autoclass:: RemfileBackend + :members: can_handle, open_url, clear + +.. autoclass:: RemfileWrapper + :members: read, seek, tell, info, close + git-annex helpers: ``datalad_fuse.utils`` ========================================= diff --git a/docs/source/cli.rst b/docs/source/cli.rst index d4acd14..1bbec6f 100644 --- a/docs/source/cli.rst +++ b/docs/source/cli.rst @@ -67,6 +67,11 @@ Options (see :ref:`caching`). The default, ``none``, only buffers data in memory while a file is open. +``--backends `` + Comma-separated, priority-ordered backends to read remote files with, e.g. + ``--backends fsspec``. The default is ``remfile,fsspec`` (see + :ref:`concepts-backends`). + ``--allow-other`` Let other users access the mount; by default, only the user who mounted it can. This requires the line ``user_allow_other`` in @@ -120,13 +125,23 @@ content of a file can be reached, or to look at the header of a file: The output is the raw content of the file, without any result rendering, so it can be piped into other tools. ``--caching ondisk`` stores the fetched -data in the dataset's cache. +data in the dataset's cache, and ``--backends`` chooses the backends to read +with, as for ``datalad fusefs`` above. + +Because it reports errors directly, ``datalad fsspec-head`` is also a handy +way to check which backend handles a file, with debug logging: + +.. code-block:: console + + $ datalad -l debug fsspec-head -d ds -c 8 sub-01/sub-01_ecephys.nwb 2>&1 | grep backend + [DEBUG] sub-01/sub-01_ecephys.nwb: opening via backend remfile Clearing the cache: ``datalad fsspec-cache-clear`` ================================================== -Removes the on-disk cache of a dataset (``.git/datalad/cache/fsspec/``): +Removes the on-disk caches of a dataset (``.git/datalad/cache/``, one +directory per backend): .. code-block:: console @@ -134,3 +149,25 @@ Removes the on-disk cache of a dataset (``.git/datalad/cache/fsspec/``): Add ``-r`` (``--recursive``) to also clear the caches of all installed subdatasets. + + +.. _cli-backends: + +Choosing the backends +===================== + +``datalad fusefs`` and ``datalad fsspec-head`` both take ``--backends``, a +comma-separated, priority-ordered list of the backends to read remote files +with (see :ref:`concepts-backends`). Without it, the configuration option +``datalad.fusefs.backends`` is used, and failing that the default +``remfile,fsspec``. Set the option like any DataLad or git configuration +option, for a dataset, globally, or for a single call: + +.. code-block:: console + + $ git config datalad.fusefs.backends fsspec # in a dataset + $ git config --global datalad.fusefs.backends fsspec # everywhere + $ datalad -c datalad.fusefs.backends=fsspec fsspec-head -d ds -c 8 file.nwb + +``datalad fsspec-cache-clear`` needs no such option: it clears the caches of +all backends. diff --git a/docs/source/concepts.rst b/docs/source/concepts.rst index 6ef5593..51f6d1e 100644 --- a/docs/source/concepts.rst +++ b/docs/source/concepts.rst @@ -2,9 +2,10 @@ How it works ************ This page explains what ``datalad-fuse`` does when a file is opened, where it -looks for content, and what that means for performance and caching. The -details apply equally to the FUSE mount, the Python adapters and ``datalad -fsspec-head``, which all share the same machinery. +looks for content, which *backend* fetches it, and what that means for +performance and caching. The details apply equally to the FUSE mount, the +Python adapters and ``datalad fsspec-head``, which all share the same +machinery. Annexed files and their content @@ -20,7 +21,7 @@ git-annex knows which *remotes* have a copy, and for some of them, URLs to download it from. For every file it opens, ``datalad-fuse`` determines one of three states -(`~datalad_fuse.fsspec.FileState`): +(`~datalad_fuse.adapter.FileState`): ``NOT_ANNEXED`` The file is in git (or untracked); it is read from disk. @@ -76,10 +77,10 @@ in order until one can be opened: .. note:: This fallback is newer than the 0.6.0 release. -If none of the candidates can be opened, opening the file fails with -``Could not find a usable URL for within ``. +If no backend can open any of the candidates, opening the file fails with +``Could not open within (backends=)``. -To see which URLs are tried, enable debug logging (see +To see which URLs are tried, and by which backend, enable debug logging (see :ref:`troubleshooting-logging`). Requests answered with a server error (HTTP 5xx) are retried up to four @@ -112,20 +113,89 @@ key's checksum. ``datalad-fuse`` reads parts of files and does not verify them, so it relies on the URLs serving the right content. +.. _concepts-backends: + +Backends +======== + +The candidate URLs above are opened by a *backend*. Two are available, and +they are tried in order; the first one that both handles the file and manages +to open one of its URLs wins: + +.. list-table:: + :header-rows: 1 + :stub-columns: 1 + :widths: 14 43 43 + + * - + - ``remfile`` + - ``fsspec`` + * - Handles + - HDF5-structured files, by extension: ``.nwb``, ``.h5``, ``.hdf5``, + ``.hdf``, ``.he5``, ``.nc``, ``.nc4`` (binary reads only) + - Every file + * - Fetches + - 100 KiB chunks, several per request, reading further ahead as long as + reads stay sequential + - Blocks of 5 MiB + * - Installed + - Optional, with the ``remfile`` extra (see :ref:`installation-backends`) + - Always, as a dependency of ``datalad-fuse`` + +The default is ``remfile,fsspec``: NWB/HDF5 files go to `remfile +`_, which is written for the +many small, scattered reads that HDF5 libraries make, and everything else +goes to fsspec. If remfile is not installed, it is skipped silently and +fsspec handles everything, as before. + +Reading a 73 MB NWB file of a DANDI dandiset from end to end through a FUSE +mount took 4.6 s with remfile and 22.6 s with fsspec; for reads of a few +arrays the two are comparable. The gain depends on the access pattern, so +measure your own if it matters. + +Choosing the backends +--------------------- + +Give a comma-separated list, in priority order, by any of: + +- the ``--backends`` option of ``datalad fusefs`` and ``datalad fsspec-head``; +- the ``datalad.fusefs.backends`` configuration option, for example + ``git config datalad.fusefs.backends fsspec`` in a dataset, or + ``datalad -c datalad.fusefs.backends=fsspec ...`` for a single call; +- the ``backends`` argument of the Python adapters (see :doc:`python`). + +So ``--backends fsspec`` reads everything with fsspec, and ``--backends +remfile`` reads *only* HDF5-structured files — other files then have no +backend that handles them, and opening them fails. A name that is requested +explicitly but not installed is skipped with a warning; an unknown name is an +error. + +Which file goes to which backend is decided from the extension of the file's +annex key (``SHA256E``, ``MD5E`` and other ``*E`` backends keep it), falling +back to the extension of the file's path for keys that carry none, such as +the ``URL``/``VURL`` keys that ``git annex addurl --fast`` produces. + +.. note:: + Backends are newer than the 0.6.0 release. Up to 0.6.0, all files were + read with fsspec, which remains the behaviour when remfile is not + installed. + + Reading only what is needed =========================== -Remote files are read with HTTP range requests, through fsspec's -`HTTPFileSystem -`_. -Data are fetched in blocks of 5 MiB, so even reading a few bytes transfers up -to 5 MiB, while reading the metadata and a few arrays of a multi-gigabyte -NWB/HDF5 file transfers only a small fraction of it. +Remote files are read with HTTP range requests: by fsspec's `HTTPFileSystem +`_ +in blocks of 5 MiB, or by remfile in chunks of 100 KiB, several of them per +request. Even reading a few bytes transfers a whole block or chunk, while +reading the metadata and a few arrays of a multi-gigabyte NWB/HDF5 file +transfers only a small fraction of it. This works best for file formats designed for partial access, such as HDF5 -and NWB, and with tools that only read the parts they need. Reading a whole file, e.g. to decompress a ``.nii.gz`` file or compute -a checksum, transfers all of it, one 5 MiB request after the other; if you -need entire files, ``datalad get`` is usually faster. +and NWB, and with tools that only read the parts they need. Reading a whole +file, e.g. to decompress a ``.nii.gz`` file or compute a checksum, transfers +all of it, one request after the other; if you need entire files, ``datalad +get`` is usually faster. The first access to a remote file takes a little time, as git-annex has to be queried and a connection established; subsequent reads of the same open file @@ -139,16 +209,18 @@ Caching By default, fetched data are kept in memory, and only while a file is open. With caching enabled (``caching=True`` for the Python adapters, ``--caching -ondisk`` for ``datalad fusefs`` and ``datalad fsspec-head``), fsspec's -`CachingFileSystem -`_ -stores the fetched blocks on disk and reuses them when the same file is read -again, also in later sessions, for up to a week; after that, fsspec considers -them expired and fetches the data again. - -- The cache of a dataset is located at ``.git/datalad/cache/fsspec/`` inside - that dataset, so each (sub)dataset has its own. -- Files are cached *sparsely*: only the blocks that were read are stored. +ondisk`` for ``datalad fusefs`` and ``datalad fsspec-head``), the backends +store what they fetch on disk and reuse it when the same file is read again, +also in later sessions. fsspec uses its `CachingFileSystem +`_, +which keeps blocks for up to a week and fetches them again after that; +remfile keeps its chunks without an expiry. + +- The caches of a dataset are located inside it, under + ``.git/datalad/cache/``, one directory per backend (``fsspec/`` and + ``remfile/``), so each (sub)dataset has its own. +- Files are cached *sparsely*: only the blocks or chunks that were read are + stored. - The cache is separate from the git-annex object store: cached files do not count as present content for ``git annex`` or ``datalad``. - The cache does not shrink by itself. Remove it with @@ -161,7 +233,7 @@ them expired and fetches the data again. Datasets with subdatasets ========================= -The FUSE mount, `~datalad_fuse.fsspec.FsspecAdapter` and ``datalad +The FUSE mount, `~datalad_fuse.adapter.RemoteFilesystemAdapter` and ``datalad fsspec-head`` work across dataset boundaries: for each file, they determine the (sub)dataset it belongs to and query that dataset's git-annex. diff --git a/docs/source/conf.py b/docs/source/conf.py index a55b059..9d4abea 100644 --- a/docs/source/conf.py +++ b/docs/source/conf.py @@ -9,7 +9,7 @@ import sys import datalad_fuse -import datalad_fuse.fsspec +import datalad_fuse.adapter docs_source = dirname(abspath(__file__)) repo_root = dirname(dirname(docs_source)) @@ -37,7 +37,10 @@ # Methods decorated with methodtools.lru_cache are hidden behind a descriptor # that autodoc does not recognize as a method. For the documentation only, # replace them with the functions they wrap, so that they are documented. -for _cls in (datalad_fuse.fsspec.DatasetAdapter, datalad_fuse.fsspec.FsspecAdapter): +for _cls in ( + datalad_fuse.adapter.DatasetAdapter, + datalad_fuse.adapter.RemoteFilesystemAdapter, +): for _name, _attr in list(vars(_cls).items()): if type(_attr).__module__.startswith("wirerope") and hasattr( _attr, "__wrapped__" diff --git a/docs/source/index.rst b/docs/source/index.rst index 1c0c7da..2b44fab 100644 --- a/docs/source/index.rst +++ b/docs/source/index.rst @@ -13,9 +13,11 @@ web server. ``datalad get`` downloads such files in full before you can open them. ``datalad-fuse`` lets you open them right away instead. It looks up the URLs -that git-annex has recorded for a file and, using -`fsspec `_, fetches only the parts of -the file that are actually read. +that git-annex has recorded for a file and fetches only the parts of the file +that are actually read, with +`fsspec `_ or, for NWB/HDF5 files and +if installed, `remfile `_ (see +:ref:`concepts-backends`). Is it a good fit? ================= @@ -44,8 +46,8 @@ There are two ways to use it: * - What - ``datalad fusefs`` presents a dataset as a read-only directory tree in which annexed files can be opened like local files. - - ``FsspecAdapter`` and ``DatasetAdapter`` return Python file objects - for files of a dataset. + - ``DatasetAdapter`` and ``RemoteFilesystemAdapter`` return Python file + objects for files of a dataset. * - Works with - Any program: ``h5ls``, MATLAB, ``nwbinspector``, shell tools, ... - Python libraries that accept file objects: h5py, pynwb, pandas, @@ -102,7 +104,7 @@ Or open files directly from Python: import h5py import pynwb - from datalad_fuse.fsspec import DatasetAdapter + from datalad_fuse.adapter import DatasetAdapter # caching=False: keep fetched data in memory only (see "Caching") with closing(DatasetAdapter("000582", caching=False)) as dsa: diff --git a/docs/source/installation.rst b/docs/source/installation.rst index 138bac1..38c4baa 100644 --- a/docs/source/installation.rst +++ b/docs/source/installation.rst @@ -57,6 +57,34 @@ released yet: $ python3 -m pip install git+https://github.com/datalad/datalad-fuse.git +.. _installation-backends: + +Optional backends +================= + +Remote files are read by a *backend* (see :ref:`concepts-backends`). fsspec, +which handles every kind of file, is always installed. `remfile +`_, which is faster for +NWB/HDF5 files, is optional: + +.. code-block:: console + + $ python3 -m pip install "datalad-fuse[remfile]" + +The ``full`` extra installs every optional backend, currently the same one: + +.. code-block:: console + + $ python3 -m pip install "datalad-fuse[full]" + +Installing remfile is enough to start using it: NWB/HDF5 files are read with +it by default from then on. Without it, those files are read with fsspec, as +before. + +.. note:: + The optional backends are newer than the 0.6.0 release. + + Checking the installation ========================= @@ -77,3 +105,9 @@ and that the Python package can be imported: .. code-block:: console $ python3 -c "import datalad_fuse; print(datalad_fuse.__version__)" + +To check whether the optional remfile backend is available: + +.. code-block:: console + + $ python3 -c "import remfile; print(remfile.__file__)" diff --git a/docs/source/python.rst b/docs/source/python.rst index bddc1e9..29d287c 100644 --- a/docs/source/python.rst +++ b/docs/source/python.rst @@ -5,12 +5,21 @@ From Python, files of a dataset can be opened without any FUSE mount, as file objects that many libraries accept in place of a file name. This page covers: -- `~datalad_fuse.fsspec.DatasetAdapter`: open files of a single dataset; -- `~datalad_fuse.fsspec.FsspecAdapter`: open files across a dataset and its - subdatasets; +- `~datalad_fuse.adapter.DatasetAdapter`: open files of a single dataset; +- `~datalad_fuse.adapter.RemoteFilesystemAdapter`: open files across a dataset + and its subdatasets; - the ``datalad`` commands of ``datalad-fuse``, called from Python; - mounting a dataset with FUSE from Python. +.. note:: + Up to 0.6.0, the adapters lived in ``datalad_fuse.fsspec``, and + `~datalad_fuse.adapter.RemoteFilesystemAdapter` was called + ``FsspecAdapter``. Both old names still work, with a + :exc:`DeprecationWarning`:: + + from datalad_fuse.fsspec import FsspecAdapter # deprecated + from datalad_fuse.adapter import RemoteFilesystemAdapter # use this + The examples use the dandiset cloned in the :doc:`tutorial`, and are run from the directory containing it: @@ -22,7 +31,7 @@ the directory containing it: Opening files of a dataset ========================== -`~datalad_fuse.fsspec.DatasetAdapter` takes the path to a dataset (or any +`~datalad_fuse.adapter.DatasetAdapter` takes the path to a dataset (or any git-annex repository), and opens files by their path relative to the dataset's top directory: @@ -30,7 +39,7 @@ dataset's top directory: from contextlib import closing - from datalad_fuse.fsspec import DatasetAdapter + from datalad_fuse.adapter import DatasetAdapter with closing(DatasetAdapter("000582", caching=False)) as dsa: with dsa.open(nwb_path) as f: # binary mode, like open(..., "rb") @@ -50,8 +59,10 @@ dataset's top directory: - for files read from disk (not annexed, or with content present), a regular Python file object; -- for files read from a URL, an fsspec file object, which fetches data as it - is read. +- for files read from a URL, an object from the backend that opened it, which + fetches data as it is read: an fsspec file object from the ``fsspec`` + backend, a `~datalad_fuse.remfile.RemfileWrapper` from ``remfile`` (see + :ref:`concepts-backends`). Text mode (``"r"`` or ``"rt"``) accepts ``encoding`` (default ``"utf-8"``) and ``errors`` arguments, as the built-in :func:`open` does. Writing is not @@ -62,6 +73,22 @@ on-disk cache inside the dataset, ``False`` only buffers them in memory while a file is open (see :ref:`caching`). With ``caching=True``, ``dsa.clear()`` removes the dataset's cache. +The optional ``backends`` argument chooses which backends to read remote +files with, as a comma-separated, priority-ordered string: + +.. code-block:: python + + # read every file with fsspec, even if remfile is installed + DatasetAdapter("000582", caching=False, backends="fsspec") + +Without it, the ``datalad.fusefs.backends`` configuration option is used, and +failing that the default ``"remfile,fsspec"`` +(`~datalad_fuse.backends.DEFAULT_BACKENDS`). See :ref:`concepts-backends` +for what the backends do. + +.. note:: + The ``backends`` argument is newer than the 0.6.0 release. + The adapter starts ``git annex`` processes to answer its queries; ``close()``, called by :func:`contextlib.closing` above, stops them. @@ -95,7 +122,7 @@ objects; use a FUSE mount for them (see :ref:`python-mount`). Inspecting files ---------------- -`~datalad_fuse.fsspec.DatasetAdapter.get_file_state` tells whether a file is +`~datalad_fuse.adapter.DatasetAdapter.get_file_state` tells whether a file is annexed and whether its content is present, and returns its git-annex key as an `~datalad_fuse.utils.AnnexKey`: @@ -132,19 +159,19 @@ The possible states are described in :doc:`concepts`. Datasets with subdatasets ========================= -`~datalad_fuse.fsspec.FsspecAdapter` works on a dataset together with its -installed subdatasets: for each path, it finds the (sub)dataset that contains -it and uses a `~datalad_fuse.fsspec.DatasetAdapter` for that dataset. It is a -context manager: +`~datalad_fuse.adapter.RemoteFilesystemAdapter` works on a dataset together +with its installed subdatasets: for each path, it finds the (sub)dataset that +contains it and uses a `~datalad_fuse.adapter.DatasetAdapter` for that +dataset. It is a context manager: .. code-block:: python from pathlib import Path - from datalad_fuse.fsspec import FsspecAdapter + from datalad_fuse.adapter import RemoteFilesystemAdapter root = Path("path/to/superdataset").resolve() - with FsspecAdapter(root, caching=False) as fsa: + with RemoteFilesystemAdapter(root, caching=False) as fsa: path = root / "subdataset" / "data" / "file.nwb" print(fsa.get_file_state(path)) print(fsa.is_under_annex(path)) @@ -158,7 +185,7 @@ context manager: Besides ``open()``, ``get_file_state()`` and ``is_under_annex()``, it offers ``get_commit_datetime()`` (the date of the last commit of the dataset containing a path) and ``resolve_dataset()`` (the -`~datalad_fuse.fsspec.DatasetAdapter` and relative path used for a path). +`~datalad_fuse.adapter.DatasetAdapter` and relative path used for a path). .. _python-commands: diff --git a/docs/source/troubleshooting.rst b/docs/source/troubleshooting.rst index 47e1302..553725b 100644 --- a/docs/source/troubleshooting.rst +++ b/docs/source/troubleshooting.rst @@ -6,14 +6,15 @@ Troubleshooting Seeing what happens =================== -Debug logging shows which files are opened, their state, and every URL that -is tried: +Debug logging shows which files are opened, their state, which backend takes +them and every URL that is tried: .. code-block:: console $ datalad -l debug fsspec-head -d ds -c 8 path/to/file [DEBUG] path/to/file: under annex, does not have content - [DEBUG] path/to/file: Attempting to open via URL https://... + [DEBUG] path/to/file: opening via backend remfile + [DEBUG] path/to/file: trying URL https://... (backend=remfile) ... ``datalad -l debug fusefs ...`` works the same way, but logs every file @@ -36,11 +37,11 @@ reading from a FUSE mount often only report a generic error. Common problems =============== -"Could not find a usable URL for within " ---------------------------------------------------------- +"Could not open within (backends=...)" +------------------------------------------------------- -The content of the file is not present locally, and none of the candidate -URLs (see :doc:`concepts`) could be opened. Check what git-annex knows about +The content of the file is not present locally, and no backend could open any +of the candidate URLs (see :doc:`concepts`). Check what git-annex knows about the file: .. code-block:: console @@ -53,6 +54,13 @@ the file: - If URLs are listed, try one of them with e.g. ``curl -I ``: the server may be down, or require authentication, which ``datalad-fuse`` does not support. +- If ``backends=`` in the message lists only backends that do not handle this + kind of file — ``backends=remfile`` for anything that is not HDF5-structured + — no backend was even tried. Add ``fsspec`` to the list (see + :ref:`concepts-backends`). + + Up to 0.6.0 this message read ``Could not find a usable URL for + within ``. Errors mentioning "identifier is not of specified type" ------------------------------------------------------- @@ -66,7 +74,7 @@ while the file is open, inside the ``with`` blocks (see :doc:`python`). ``AttributeError: 'NoneType' object has no attribute 'get_commit_date'`` ------------------------------------------------------------------------ -The path given to `~datalad_fuse.fsspec.DatasetAdapter` is not a dataset (or +The path given to `~datalad_fuse.adapter.DatasetAdapter` is not a dataset (or git repository). Check the path, and the current directory if the path is relative. @@ -151,7 +159,7 @@ get -n path/to/subdataset``, and remount if you are using a FUSE mount. ``ValueError`` "Path not under root dataset" or "is not in the subpath of" -------------------------------------------------------------------------- -A path passed to `~datalad_fuse.fsspec.FsspecAdapter` was relative. Use an +A path passed to `~datalad_fuse.adapter.RemoteFilesystemAdapter` was relative. Use an absolute ``root`` and absolute paths (see :doc:`python`). Changes to the dataset are not reflected @@ -162,13 +170,45 @@ remembered by an adapter or a mount once determined. After ``datalad get``, ``datalad drop``, ``git checkout`` etc. in the dataset, create a new adapter or remount. +Warning "Backend 'remfile' requested but not available" +------------------------------------------------------- + +A backend named in ``--backends`` or in ``datalad.fusefs.backends`` is not +installed, so it is skipped and the remaining ones are used. Install it (see +:ref:`installation-backends`) or drop it from the list. When the default +``remfile,fsspec`` is in effect and remfile is missing, nothing is said: that +is the normal case, and fsspec handles everything. + +``ValueError`` "No usable backends from spec ..." +-------------------------------------------------- + +None of the requested backends could be created: either every name is an +uninstalled backend, or the list is empty. ``ValueError: Unknown backend: +'...'`` instead means a name that does not exist; the names are ``remfile`` +and ``fsspec``. + +A file is read by the wrong backend +----------------------------------- + +Which backend takes a file is decided from the extension of its annex key, +falling back to the extension of its path (see :ref:`concepts-backends`). Run +``datalad -l debug fsspec-head ...`` to see the decision; the ``cannot handle +(suffix=..., mode=...)`` lines show what each backend was offered. Force a +single backend with ``--backends fsspec`` or ``--backends remfile``, for +example to check whether a problem is specific to one of them. + Reading is slow --------------- - The first access to a remote file takes a moment, to query git-annex and to connect to the server. -- Data are fetched in blocks of 5 MiB, so reading many small, scattered - pieces of a file is slow. +- Data are fetched in blocks of 5 MiB (fsspec) or chunks of 100 KiB + (remfile), so reading many small, scattered pieces of a file is slow. For + NWB/HDF5 files, installing remfile (see :ref:`installation-backends`) + usually helps. +- Unreachable URLs delay the fall-through to the next one. ``git annex + whereis`` shows which URLs are recorded; ``datalad -l debug`` shows which + one is slow. - Reading entire files, or decompressing them, fetches everything, and is faster with ``datalad get``. - Use ``--caching ondisk`` (``caching=True`` in Python) if the same files are diff --git a/docs/source/tutorial.rst b/docs/source/tutorial.rst index b209ddc..9d92e56 100644 --- a/docs/source/tutorial.rst +++ b/docs/source/tutorial.rst @@ -22,7 +22,14 @@ the packages used to read and plot the data: .. code-block:: console - $ python3 -m pip install datalad-fuse h5py pynwb matplotlib + $ python3 -m pip install "datalad-fuse[remfile]" h5py pynwb matplotlib + +The ``[remfile]`` extra pulls in `remfile +`_, which ``datalad-fuse`` then +uses for NWB and other HDF5 files; it is written for the access pattern of +HDF5 libraries and is noticeably faster for them (see +:ref:`concepts-backends`). Everything below works without it too, just more +slowly. The FUSE part of the tutorial also uses ``h5ls`` from the HDF5 command-line tools (``sudo apt-get install hdf5-tools`` on Debian/Ubuntu, or ``conda install @@ -91,7 +98,7 @@ prints the first bytes of the file, here the signature of an HDF5 file: Read the data from Python ========================= -`DatasetAdapter ` gives access to the +`DatasetAdapter ` gives access to the files of one dataset, addressed by paths relative to the dataset's top directory. Its ``open()`` method returns a file object that h5py, and hence PyNWB, can read from: @@ -105,7 +112,7 @@ PyNWB, can read from: import numpy as np import pynwb - from datalad_fuse.fsspec import DatasetAdapter + from datalad_fuse.adapter import DatasetAdapter nwb_path = "sub-10073/sub-10073_ses-17010302_behavior+ecephys.nwb" @@ -262,13 +269,13 @@ or, for a mount: $ datalad fusefs -d 000582 --foreground --caching ondisk mnt & -The cache is kept inside the dataset, under ``.git/datalad/cache/fsspec/``, -and holds only the parts of files that were read: +The cache is kept inside the dataset, under ``.git/datalad/cache/``, one +directory per backend, and holds only the parts of files that were read: .. code-block:: console - $ du -sh 000582/.git/datalad/cache/fsspec - 10M 000582/.git/datalad/cache/fsspec + $ du -sh 000582/.git/datalad/cache/* + 12M 000582/.git/datalad/cache/remfile Remove it when it is no longer needed: diff --git a/setup.cfg b/setup.cfg index eaf9a43..887b5c8 100644 --- a/setup.cfg +++ b/setup.cfg @@ -28,6 +28,12 @@ include_package_data = True include = datalad_fuse* [options.extras_require] +full = + datalad-fuse[remfile] +remfile = + # Pinned to 0.1.x: RemfileBackend passes the private ``_size=`` kwarg to + # remfile.File to skip its retrying size probe (see datalad_fuse/remfile.py) + remfile ~= 0.1.15 test = coverage~=6.0 linesep~=0.2 diff --git a/tox.ini b/tox.ini index 29fe157..5263e1f 100644 --- a/tox.ini +++ b/tox.ini @@ -19,7 +19,13 @@ python = passenv = HOME DATALAD_* -extras = test +extras = + test + # `full` pulls in optional backends (e.g. remfile). Installed only in + # envs whose name carries `full` or `libfuse`, so the default + # `py3{10..14}-nonetwork` rows stay lean and remfile remains optional. + # Tests that require remfile are gated by a `requires_remfile` skipif. + full,libfuse: full # Factors: # nonetwork — disable network tests (sets bogus proxies to catch leaks) # full — also enable --libfuse tests @@ -58,8 +64,11 @@ deps = extras = test commands = mypy --follow-imports skip \ + datalad_fuse/adapter.py \ + datalad_fuse/backends.py \ datalad_fuse/fsspec.py \ datalad_fuse/fuse_.py \ + datalad_fuse/remfile.py \ datalad_fuse/utils.py [pytest]