From 2ca1c25ee7dfd120b1965f620955f39d1e49117e Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 20:46:50 +0000 Subject: [PATCH 1/5] Read missing annexed content with datalad-fuse, cached by annex key Add `dandi.support.datalad_fuse.AnnexedReadableFile`, a `Readable` for a locked annexed file whose content is not present locally, streamed with the adapter of datalad-fuse (no FUSE mount), which asks git-annex where the content is. It is available with the new `datalad` extra. Cache NWB metadata of locked annexed files, streamed or local, under their git-annex key paired with their path, with fscacher 0.5.0's `annex_key_fingerprint`. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01Qn5WBSiQgoZoL4fytF6nEr --- .github/workflows/run-tests.yml | 7 ++ dandi/metadata/nwb.py | 3 +- dandi/pynwb_utils.py | 25 +++- dandi/support/datalad_fuse.py | 147 +++++++++++++++++++++++ dandi/support/tests/test_datalad_fuse.py | 97 +++++++++++++++ dandi/tests/fixtures.py | 123 +++++++++++++++++++ dandi/tests/skip.py | 5 + pyproject.toml | 8 +- 8 files changed, 410 insertions(+), 5 deletions(-) create mode 100644 dandi/support/datalad_fuse.py create mode 100644 dandi/support/tests/test_datalad_fuse.py diff --git a/.github/workflows/run-tests.yml b/.github/workflows/run-tests.yml index fec7c9104..d33f99a4f 100644 --- a/.github/workflows/run-tests.yml +++ b/.github/workflows/run-tests.yml @@ -73,6 +73,9 @@ jobs: - os: ubuntu-latest python: '3.11' mode: lowest-deps + - os: ubuntu-latest + python: '3.12' + mode: datalad steps: - name: Set up environment @@ -123,6 +126,10 @@ jobs: git+https://github.com/hdmf-dev/hdmf \ git+https://github.com/hdmf-dev/hdmf-zarr + - name: Install git-annex and datalad-fuse (for streaming annexed content) + if: matrix.mode == 'datalad' + run: pip install git-annex ".[datalad]" + - name: Create NFS filesystem if: matrix.mode == 'nfs' run: | diff --git a/dandi/metadata/nwb.py b/dandi/metadata/nwb.py index bf43dbcb2..140f2e1f7 100644 --- a/dandi/metadata/nwb.py +++ b/dandi/metadata/nwb.py @@ -16,6 +16,7 @@ from ..misctypes import DUMMY_DANDI_ETAG, Digest, LocalReadableFile, Readable from ..pynwb_utils import ( _get_pynwb_metadata, + annex_fingerprint, get_neurodata_types, get_nwb_version, ignore_benign_pynwb_warnings, @@ -28,7 +29,7 @@ # Disable this for clean hacking -@metadata_cache.memoize_path +@metadata_cache.memoize_path(custom_fingerprint=annex_fingerprint) def get_metadata( path: str | Path | Readable, digest: Digest | None = None ) -> dict[str, Any]: diff --git a/dandi/pynwb_utils.py b/dandi/pynwb_utils.py index 297df05f3..da6e4d835 100644 --- a/dandi/pynwb_utils.py +++ b/dandi/pynwb_utils.py @@ -23,7 +23,7 @@ import warnings import dandischema -from fscacher import PersistentCache +from fscacher import PersistentCache, annex_key_fingerprint import h5py import hdmf import numpy as np @@ -73,6 +73,25 @@ ) +def annex_fingerprint(source: Any) -> tuple[str, str] | None: + """ + Fingerprint of ``source`` for ``PersistentCache.memoize_path`` + + Pass it as ``custom_fingerprint`` to cache the results of a function of a + local path or a `Readable` under the git-annex key of a locked annexed + file, paired with its path (see `fscacher.annex_key_fingerprint`), rather + than under ``stat()``: the content of an `AnnexedReadableFile`, which is + not present locally, cannot be ``stat()``-ed. Anything else is handled as + without this. + """ + # Avoid circular import: + from .support.datalad_fuse import AnnexedReadableFile + + if isinstance(source, AnnexedReadableFile): + source = source.filepath + return cast("tuple[str, str] | None", annex_key_fingerprint(source)) + + def _sanitize_nwb_version( v: Any, filename: str | Path | None = None, @@ -194,7 +213,7 @@ def get_neurodata_types_to_modalities_map() -> dict[str, str]: return ndtypes -@metadata_cache.memoize_path +@metadata_cache.memoize_path(custom_fingerprint=annex_fingerprint) def get_neurodata_types(filepath: str | Path | Readable) -> list[str]: with open_readable(filepath) as fp, h5py.File(fp, "r") as h5file: all_pairs = _scan_neurodata_types(h5file) @@ -808,7 +827,7 @@ def copy_nwb_file(src: str | Path, dest: str | Path) -> str: return str(dest) -@metadata_cache.memoize_path +@metadata_cache.memoize_path(custom_fingerprint=annex_fingerprint) def nwb_has_external_links(filepath: str | Path | Readable) -> bool: with open_readable(filepath) as f, h5py.File(f, "r") as fp: visited = set() diff --git a/dandi/support/datalad_fuse.py b/dandi/support/datalad_fuse.py new file mode 100644 index 000000000..8460039f3 --- /dev/null +++ b/dandi/support/datalad_fuse.py @@ -0,0 +1,147 @@ +""" +Reading the content of annexed files of git-annex_ repositories (such as +DataLad_ datasets) that is not present locally, with datalad-fuse_. + +In a git-annex repository, a (locked) annexed file is a symbolic link into +:file:`.git/annex/objects/`, which is broken until the content is fetched. +datalad-fuse can mount such a repository with FUSE so that the content is +streamed on demand, but its adapter, which the mount is built on, also opens +the files directly with no mount. It asks git-annex where the content is +(``git annex whereis``) and streams it over HTTP(S) with fsspec_. That adapter +is used here, so that the files of a DataLad Dandiset can be read (e.g., for +their metadata, or to be validated) without downloading the (possibly +terabytes of) data. + +This requires datalad-fuse (``pip install "dandi[datalad]"``), which requires +DataLad and git-annex. + +.. _git-annex: https://git-annex.branchable.com +.. _DataLad: https://www.datalad.org +.. _datalad-fuse: https://github.com/datalad/datalad-fuse +.. _fsspec: https://github.com/fsspec/filesystem_spec +""" + +from __future__ import annotations + +import atexit +from dataclasses import dataclass, field +from datetime import datetime +from functools import cache +import logging +import os +from pathlib import Path +import subprocess +from typing import IO, Any, cast + +from ..misctypes import Readable + +lgr = logging.getLogger("dandi.support.datalad_fuse") + + +@cache +def get_adapter() -> Any: + """ + The datalad-fuse adapter shared by all files, rooted at the filesystem + root so that it serves files of any dataset, and closed at exit + """ + # Optional dependency: + from datalad_fuse.fsspec import FsspecAdapter + + # caching=False: do not keep the streamed blocks on disk, in the dataset + adapter = FsspecAdapter(Path(os.path.abspath(os.sep)), caching=False) + atexit.register(adapter.__exit__, None, None, None) + return adapter + + +@cache +def annex_initialized(directory: Path) -> bool: + """ + Whether ``directory`` is in a Git repository in which git-annex is + initialized. (DataLad would otherwise initialize it, modifying the + repository.) + """ + try: + r = subprocess.run( + ["git", "-C", str(directory), "config", "--get", "annex.uuid"], + capture_output=True, + text=True, + ) + except FileNotFoundError: + return False + return r.returncode == 0 and bool(r.stdout.strip()) + + +@dataclass +class AnnexedReadableFile(Readable): + """ + A `Readable` for a (locked) annexed file whose content is not present + locally, which is streamed with datalad-fuse (see the module docstring). + + Instances are obtained by calling `get_annexed_readable()`. + """ + + #: The absolute path of the (broken) symbolic link to the file's content + filepath: Path + + #: The git-annex key of the file + key: str + + #: The size of the content, as recorded in the key + size: int + + #: The first URL that the content may be streamed from (for reporting; + #: the adapter tries the others if it cannot be read from that one) + url: str + + adapter: Any = field(repr=False, compare=False) + + def open(self) -> IO[bytes]: + return cast("IO[bytes]", self.adapter.open(self.filepath, "rb")) + + def get_size(self) -> int: + return self.size + + def get_mtime(self) -> datetime | None: + return None + + def get_filename(self) -> str: + return self.filepath.name + + def __str__(self) -> str: + return str(self.filepath) + + +def get_annexed_readable(path: str | Path) -> AnnexedReadableFile | None: + """ + Return an `AnnexedReadableFile` for streaming the content of the annexed + file at ``path``, or `None` if datalad-fuse is not installed, ``path`` is + not a (locked) annexed file whose content is missing in a repository in + which git-annex is initialized, its key does not record the size of the + content, or git-annex knows no URL to stream it from + """ + filepath = Path(os.path.abspath(path)) + if not filepath.is_symlink() or filepath.exists(): + # Not a broken symbolic link + return None + try: + from datalad_fuse.fsspec import FileState + except ImportError: + return None + if not annex_initialized(filepath.parent): + return None + try: + dsap, relpath = get_adapter().resolve_dataset(filepath) + state, key = dsap.get_file_state(relpath) + except Exception as e: + # E.g., a repository without any commit + lgr.debug("%s: Cannot be read with datalad-fuse: %s", filepath, e) + return None + if state is not FileState.NO_CONTENT or key is None or key.size is None: + return None + url = next(dsap.get_urls(str(key)), None) + if url is None: + lgr.debug("%s: git-annex knows no URL for key %s", filepath, key) + return None + return AnnexedReadableFile( + filepath=filepath, key=str(key), size=key.size, url=url, adapter=get_adapter() + ) diff --git a/dandi/support/tests/test_datalad_fuse.py b/dandi/support/tests/test_datalad_fuse.py new file mode 100644 index 000000000..6dc31495a --- /dev/null +++ b/dandi/support/tests/test_datalad_fuse.py @@ -0,0 +1,97 @@ +from __future__ import annotations + +from pathlib import Path +import shutil + +from fscacher import PersistentCache +import h5py +import pytest + +from ..datalad_fuse import AnnexedReadableFile, get_annexed_readable +from ...pynwb_utils import annex_fingerprint, get_neurodata_types +from ...tests.fixtures import RangeHTTPServer, make_git_annex_dandiset +from ...tests.skip import mark + +pytestmark = mark.skipif_no_git_annex + + +@pytest.fixture +def streamed_nwb( + range_http_server: tuple[RangeHTTPServer, str], simple2_nwb: Path, tmp_path: Path +) -> tuple[Path, RangeHTTPServer]: + """ + An annexed copy of ``simple2_nwb`` in a git-annex repository, whose content + is not present but registered at a URL of the HTTP server + """ + pytest.importorskip("datalad_fuse.fsspec") + server, base_url = range_http_server + shutil.copy(simple2_nwb, tmp_path / "served" / "content.nwb") + make_git_annex_dandiset( + tmp_path / "ds", + {"sub-01/sub-01.nwb": (simple2_nwb, f"{base_url}/content.nwb")}, + ) + nwb = tmp_path / "ds" / "sub-01" / "sub-01.nwb" + assert not nwb.exists() + return nwb, server + + +@pytest.mark.ai_generated +def test_annexed_readable( + streamed_nwb: tuple[Path, RangeHTTPServer], simple2_nwb: Path +) -> None: + nwb, server = streamed_nwb + r = get_annexed_readable(nwb) + assert isinstance(r, AnnexedReadableFile) + assert r.key.startswith(f"SHA256E-s{simple2_nwb.stat().st_size}--") + assert r.url.endswith("/content.nwb") + assert r.get_size() == simple2_nwb.stat().st_size + assert r.get_filename() == "sub-01.nwb" + fp = r.open() + try: + # Random access over HTTP is enough for h5py to read the file + with h5py.File(fp, "r") as h5: + assert h5.attrs["nwb_version"] + finally: + fp.close() + assert server.ranged_requests > 0 + + +@pytest.mark.ai_generated +def test_annexed_readable_none( + streamed_nwb: tuple[Path, RangeHTTPServer], simple2_nwb: Path, tmp_path: Path +) -> None: + nwb, _ = streamed_nwb + # The content is present + (nwb.parent / "regular.nwb").write_bytes(b"content") + assert get_annexed_readable(nwb.parent / "regular.nwb") is None + # Not in a Git repository + assert get_annexed_readable(simple2_nwb) is None + # A broken symbolic link in a repository without git-annex initialized + plain = tmp_path / "plain" + plain.mkdir() + shutil.copytree(nwb.parent, plain / "sub-01", symlinks=True) + shutil.copytree(nwb.parent.parent / ".git", plain / ".git") + (plain / ".git" / "config").write_text("[core]\n\tbare = false\n") + assert get_annexed_readable(plain / "sub-01" / "sub-01.nwb") is None + + +@pytest.mark.ai_generated +def test_annexed_readable_cached_by_key( + streamed_nwb: tuple[Path, RangeHTTPServer], simple2_nwb: Path, tmp_path: Path +) -> None: + nwb, server = streamed_nwb + cache = PersistentCache(path=tmp_path / "cache") + neurodata_types = cache.memoize_path(custom_fingerprint=annex_fingerprint)( + get_neurodata_types.__wrapped__ + ) + r = get_annexed_readable(nwb) + assert r is not None + expected = get_neurodata_types.__wrapped__(simple2_nwb) + assert neurodata_types(r) == expected + requests = server.ranged_requests + assert requests > 0 + # Served from the cache, without the content being read again: by the + # readable, and by the path of the (broken) link itself + assert neurodata_types(get_annexed_readable(nwb)) == expected + assert neurodata_types(nwb) == expected + assert server.ranged_requests == requests diff --git a/dandi/tests/fixtures.py b/dandi/tests/fixtures.py index 9b128a23d..438219c99 100644 --- a/dandi/tests/fixtures.py +++ b/dandi/tests/fixtures.py @@ -3,12 +3,16 @@ from collections.abc import Callable, Iterator from dataclasses import dataclass, field, replace from datetime import datetime, timezone +from functools import partial +from http import HTTPStatus +from http.server import SimpleHTTPRequestHandler, ThreadingHTTPServer import logging import os from pathlib import Path import re import shutil from subprocess import DEVNULL, check_output, run +import threading from time import sleep from typing import Any, Literal from uuid import uuid4 @@ -352,6 +356,125 @@ def fixture( return fixture +def make_git_annex_dandiset(dandiset: Path, files: dict[str, tuple[Path, str]]) -> None: + """ + Make the directory ``dandiset`` (created if needed, with a minimal + :file:`dandiset.yaml`) a git-annex repository looking like a clone of a + DataLad Dandiset whose annexed content has not been fetched: for each + ``relpath: (source, url)`` item of ``files``, a copy of ``source`` is + annexed at ``relpath`` with ``url`` registered for it, and then dropped. + This needs git-annex. + """ + env = { + **os.environ, + "GIT_AUTHOR_NAME": "DANDI tests", + "GIT_AUTHOR_EMAIL": "tests@dandiarchive.org", + "GIT_COMMITTER_NAME": "DANDI tests", + "GIT_COMMITTER_EMAIL": "tests@dandiarchive.org", + } + + def git(*args: str) -> str: + return check_output( + ["git", "-C", str(dandiset), *args], env=env, text=True + ).strip() + + dandiset.mkdir(parents=True, exist_ok=True) + (dandiset / dandiset_metadata_file).write_text( + "identifier: '000027'\nname: Test\ndescription: Test dandiset\n" + ) + git("init", "-q") + git("annex", "init", "-q") + git("config", "annex.backend", "SHA256E") + git("add", dandiset_metadata_file) + for relpath, (source, _) in files.items(): + (dandiset / relpath).parent.mkdir(parents=True, exist_ok=True) + shutil.copy(source, dandiset / relpath) + git("annex", "add", "-q", *files) + for relpath, (_, url) in files.items(): + key = git("annex", "lookupkey", relpath) + git("annex", "registerurl", key, url) + git("commit", "-q", "-m", "Add files") + git("annex", "drop", "-q", "--force", *files) + + +class RangeHTTPServer(ThreadingHTTPServer): + """An HTTP server counting the range requests it served""" + + ranged_requests: int = 0 + + +class RangeHTTPRequestHandler(SimpleHTTPRequestHandler): + """ + A handler serving files with support for ``HEAD`` and single-range ``GET`` + requests, like S3 does + """ + + def log_message(self, format: str, *args: Any) -> None: + pass + + def _send_headers( + self, status: HTTPStatus, length: int, content_range: str | None = None + ) -> None: + self.send_response(status) + self.send_header("Content-Type", "application/octet-stream") + self.send_header("Content-Length", str(length)) + self.send_header("Accept-Ranges", "bytes") + if content_range is not None: + self.send_header("Content-Range", content_range) + self.end_headers() + + def do_HEAD(self) -> None: + path = self.translate_path(self.path) + if not os.path.isfile(path): + self.send_error(HTTPStatus.NOT_FOUND) + return + self._send_headers(HTTPStatus.OK, os.path.getsize(path)) + + def do_GET(self) -> None: + path = self.translate_path(self.path) + if not os.path.isfile(path): + self.send_error(HTTPStatus.NOT_FOUND) + return + size = os.path.getsize(path) + start, end = 0, size - 1 + if m := re.fullmatch(r"bytes=(\d+)-(\d*)", self.headers.get("Range", "")): + start = int(m[1]) + end = min(int(m[2]) if m[2] else size - 1, size - 1) + assert isinstance(self.server, RangeHTTPServer) + self.server.ranged_requests += 1 + self._send_headers( + HTTPStatus.PARTIAL_CONTENT, + end - start + 1, + f"bytes {start}-{end}/{size}", + ) + else: + self._send_headers(HTTPStatus.OK, size) + with open(path, "rb") as fp: + fp.seek(start) + self.wfile.write(fp.read(end - start + 1)) + + +@pytest.fixture() +def range_http_server(tmp_path: Path) -> Iterator[tuple[RangeHTTPServer, str]]: + """ + Serve the files in ``tmp_path / "served"`` over HTTP with support for range + requests; yields the server and its base URL + """ + served = tmp_path / "served" + served.mkdir() + server = RangeHTTPServer( + ("127.0.0.1", 0), partial(RangeHTTPRequestHandler, directory=str(served)) + ) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + yield server, f"http://127.0.0.1:{server.server_address[1]}" + finally: + server.shutdown() + server.server_close() + thread.join() + + def _make_subdirs_dandisets(path: Path) -> None: for bids_dataset_path in path.iterdir(): if bids_dataset_path.is_dir(): diff --git a/dandi/tests/skip.py b/dandi/tests/skip.py index 54d864b76..9e30b195d 100644 --- a/dandi/tests/skip.py +++ b/dandi/tests/skip.py @@ -125,6 +125,10 @@ def no_git(): return "Git not installed", shutil.which("git") is None +def no_git_annex(): + return "git-annex not installed", shutil.which("git-annex") is None + + # ### END MODIFIED CODE @@ -162,6 +166,7 @@ def on_windows(): no_docker_commands, no_docker_engine, no_git, + no_git_annex, no_network, # no_singularity, no_ssh, diff --git a/pyproject.toml b/pyproject.toml index b98d4f7a0..e7d3637fb 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -53,7 +53,7 @@ dependencies = [ "dandischema >= 0.12.0, != 0.13.0, < 0.15.0", "etelemetry >= 0.2.2", "fasteners >= 0.19", - "fscacher >= 0.3.0", + "fscacher >= 0.5.0", # Floor raised to what `py3-lowest` actually resolves: modern pynwb pins # hdmf >= 4, so declared floor is now the tested floor. "hdmf >= 4.1.0", @@ -119,6 +119,11 @@ test = [ # 5.0 switched to urllib3.connection.VerifiedHTTPSConnection (urllib3 2.x compat) "vcrpy >= 5.0.0", ] +# For reading the content of annexed files of DataLad datasets without +# fetching it (requires git-annex) +datalad = [ + "datalad-fuse >= 0.5.1", +] tools = [ "boto3", ] @@ -214,6 +219,7 @@ follow_imports = "skip" module = [ "bidsschematools.*", "click_didyoumean.*", + "datalad_fuse.*", "etelemetry.*", "fasteners.*", "fscacher.*", From 52b675ee05b2eee9b2c1a44abafe8e8183338510 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 21:17:17 +0000 Subject: [PATCH 2/5] Move annex_fingerprint next to AnnexedReadableFile; trim field comments annex_fingerprint is not specific to NWB files, and living in dandi.support.datalad_fuse it needs no deferred import. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01Qn5WBSiQgoZoL4fytF6nEr --- dandi/metadata/nwb.py | 2 +- dandi/pynwb_utils.py | 22 ++------------- dandi/support/datalad_fuse.py | 34 ++++++++++++++++-------- dandi/support/tests/test_datalad_fuse.py | 4 +-- 4 files changed, 28 insertions(+), 34 deletions(-) diff --git a/dandi/metadata/nwb.py b/dandi/metadata/nwb.py index 140f2e1f7..d0468d607 100644 --- a/dandi/metadata/nwb.py +++ b/dandi/metadata/nwb.py @@ -16,13 +16,13 @@ from ..misctypes import DUMMY_DANDI_ETAG, Digest, LocalReadableFile, Readable from ..pynwb_utils import ( _get_pynwb_metadata, - annex_fingerprint, get_neurodata_types, get_nwb_version, ignore_benign_pynwb_warnings, metadata_cache, nwb_has_external_links, ) +from ..support.datalad_fuse import annex_fingerprint from ..utils import find_parent_directory_containing lgr = get_logger() diff --git a/dandi/pynwb_utils.py b/dandi/pynwb_utils.py index da6e4d835..c49293f57 100644 --- a/dandi/pynwb_utils.py +++ b/dandi/pynwb_utils.py @@ -23,7 +23,7 @@ import warnings import dandischema -from fscacher import PersistentCache, annex_key_fingerprint +from fscacher import PersistentCache import h5py import hdmf import numpy as np @@ -42,6 +42,7 @@ metadata_nwb_subject_fields, ) from .misctypes import Readable +from .support.datalad_fuse import annex_fingerprint from .utils import get_module_version, is_url from .validate._types import ( Origin, @@ -73,25 +74,6 @@ ) -def annex_fingerprint(source: Any) -> tuple[str, str] | None: - """ - Fingerprint of ``source`` for ``PersistentCache.memoize_path`` - - Pass it as ``custom_fingerprint`` to cache the results of a function of a - local path or a `Readable` under the git-annex key of a locked annexed - file, paired with its path (see `fscacher.annex_key_fingerprint`), rather - than under ``stat()``: the content of an `AnnexedReadableFile`, which is - not present locally, cannot be ``stat()``-ed. Anything else is handled as - without this. - """ - # Avoid circular import: - from .support.datalad_fuse import AnnexedReadableFile - - if isinstance(source, AnnexedReadableFile): - source = source.filepath - return cast("tuple[str, str] | None", annex_key_fingerprint(source)) - - def _sanitize_nwb_version( v: Any, filename: str | Path | None = None, diff --git a/dandi/support/datalad_fuse.py b/dandi/support/datalad_fuse.py index 8460039f3..a33c84fe9 100644 --- a/dandi/support/datalad_fuse.py +++ b/dandi/support/datalad_fuse.py @@ -33,6 +33,8 @@ import subprocess from typing import IO, Any, cast +from fscacher import annex_key_fingerprint + from ..misctypes import Readable lgr = logging.getLogger("dandi.support.datalad_fuse") @@ -74,25 +76,19 @@ def annex_initialized(directory: Path) -> bool: @dataclass class AnnexedReadableFile(Readable): """ - A `Readable` for a (locked) annexed file whose content is not present - locally, which is streamed with datalad-fuse (see the module docstring). + A `Readable` for a (locked) annexed file at ``filepath`` (a broken symbolic + link) whose content is not present locally, which is streamed with + datalad-fuse (see the module docstring). ``url`` is only the first URL + known for the content, for reporting: if it cannot be read, the adapter + tries the others. Instances are obtained by calling `get_annexed_readable()`. """ - #: The absolute path of the (broken) symbolic link to the file's content filepath: Path - - #: The git-annex key of the file key: str - - #: The size of the content, as recorded in the key size: int - - #: The first URL that the content may be streamed from (for reporting; - #: the adapter tries the others if it cannot be read from that one) url: str - adapter: Any = field(repr=False, compare=False) def open(self) -> IO[bytes]: @@ -111,6 +107,22 @@ def __str__(self) -> str: return str(self.filepath) +def annex_fingerprint(source: Any) -> tuple[str, str] | None: + """ + Fingerprint of ``source`` for ``PersistentCache.memoize_path`` + + Pass it as ``custom_fingerprint`` to cache the results of a function of a + local path or a `Readable` under the git-annex key of a locked annexed + file, paired with its path (see `fscacher.annex_key_fingerprint`), rather + than under ``stat()``: the content of an `AnnexedReadableFile`, which is + not present locally, cannot be ``stat()``-ed. Anything else is handled as + without this. + """ + if isinstance(source, AnnexedReadableFile): + source = source.filepath + return cast("tuple[str, str] | None", annex_key_fingerprint(source)) + + def get_annexed_readable(path: str | Path) -> AnnexedReadableFile | None: """ Return an `AnnexedReadableFile` for streaming the content of the annexed diff --git a/dandi/support/tests/test_datalad_fuse.py b/dandi/support/tests/test_datalad_fuse.py index 6dc31495a..6f594eb07 100644 --- a/dandi/support/tests/test_datalad_fuse.py +++ b/dandi/support/tests/test_datalad_fuse.py @@ -7,8 +7,8 @@ import h5py import pytest -from ..datalad_fuse import AnnexedReadableFile, get_annexed_readable -from ...pynwb_utils import annex_fingerprint, get_neurodata_types +from ..datalad_fuse import AnnexedReadableFile, annex_fingerprint, get_annexed_readable +from ...pynwb_utils import get_neurodata_types from ...tests.fixtures import RangeHTTPServer, make_git_annex_dandiset from ...tests.skip import mark From 22a54594e2dbee28fce0d0867d25fdfac593492f Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 21:18:07 +0000 Subject: [PATCH 3/5] Note URL order and proxy behavior of datalad-fuse As documented in datalad/datalad-fuse#138. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01Qn5WBSiQgoZoL4fytF6nEr --- dandi/support/datalad_fuse.py | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/dandi/support/datalad_fuse.py b/dandi/support/datalad_fuse.py index a33c84fe9..cfe71d895 100644 --- a/dandi/support/datalad_fuse.py +++ b/dandi/support/datalad_fuse.py @@ -12,6 +12,10 @@ their metadata, or to be validated) without downloading the (possibly terabytes of) data. +The URLs are tried in the order git-annex lists them; for a DataLad Dandiset, +that is the DANDI Archive's API download URL (which redirects to S3) before the +direct S3 URL. datalad-fuse does not use HTTP proxy environment variables. + This requires datalad-fuse (``pip install "dandi[datalad]"``), which requires DataLad and git-annex. @@ -49,7 +53,8 @@ def get_adapter() -> Any: # Optional dependency: from datalad_fuse.fsspec import FsspecAdapter - # caching=False: do not keep the streamed blocks on disk, in the dataset + # caching=False: do not keep the streamed blocks on disk, in the dataset. + # The root and all paths passed to the adapter must be absolute. adapter = FsspecAdapter(Path(os.path.abspath(os.sep)), caching=False) atexit.register(adapter.__exit__, None, None, None) return adapter From 1ad6495fccfd121ddc6f0070a5a1ac0b3d063f0c Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 6 Oct 2026 21:26:15 +0000 Subject: [PATCH 4/5] Import fscacher lazily in dandi.support.datalad_fuse fscacher imports joblib, which imports numpy, so importing it with the module made the CLI (via dandi.validate in the stacked PR) import numpy, failing test_no_heavy_imports. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01Qn5WBSiQgoZoL4fytF6nEr --- dandi/support/datalad_fuse.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/dandi/support/datalad_fuse.py b/dandi/support/datalad_fuse.py index cfe71d895..54c27a18c 100644 --- a/dandi/support/datalad_fuse.py +++ b/dandi/support/datalad_fuse.py @@ -37,8 +37,6 @@ import subprocess from typing import IO, Any, cast -from fscacher import annex_key_fingerprint - from ..misctypes import Readable lgr = logging.getLogger("dandi.support.datalad_fuse") @@ -123,6 +121,10 @@ def annex_fingerprint(source: Any) -> tuple[str, str] | None: not present locally, cannot be ``stat()``-ed. Anything else is handled as without this. """ + # Avoid heavy import (fscacher imports joblib, which imports numpy) when + # this module is imported, e.g., by the CLI: + from fscacher import annex_key_fingerprint + if isinstance(source, AnnexedReadableFile): source = source.filepath return cast("tuple[str, str] | None", annex_key_fingerprint(source)) From 796463300563e9ee347cdf38540ffdbd30552a0c Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 7 Oct 2026 02:40:19 +0000 Subject: [PATCH 5/5] Disable automatic GC while a streamed annexed file is open h5py holds its global lock while reading from a Python file object, and the read waits for another thread to fetch the data over HTTP (fsspec's I/O thread, or an HTTP server in the same process as in the tests). A garbage collection in that thread could then run the finalizer of an unreferenced h5py-backed object (e.g., hdmf's HDF5IO.__del__), which waits for the same lock: a deadlock, as seen in CI in test_validate_stream, ended by pytest-timeout with a segfault. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01Qn5WBSiQgoZoL4fytF6nEr --- dandi/support/datalad_fuse.py | 60 +++++++++++++++++++++++- dandi/support/tests/test_datalad_fuse.py | 21 +++++++++ 2 files changed, 80 insertions(+), 1 deletion(-) diff --git a/dandi/support/datalad_fuse.py b/dandi/support/datalad_fuse.py index 54c27a18c..1dfaf117c 100644 --- a/dandi/support/datalad_fuse.py +++ b/dandi/support/datalad_fuse.py @@ -31,10 +31,12 @@ from dataclasses import dataclass, field from datetime import datetime from functools import cache +import gc import logging import os from pathlib import Path import subprocess +import threading from typing import IO, Any, cast from ..misctypes import Readable @@ -76,6 +78,62 @@ def annex_initialized(directory: Path) -> bool: return r.returncode == 0 and bool(r.stdout.strip()) +_gc_lock = threading.Lock() +_gc_pauses = 0 +_gc_was_enabled = False + + +class _GCPausedFile: + """ + A file object streamed by datalad-fuse, during whose lifetime automatic + garbage collection is disabled + + h5py holds a global lock while it reads from a Python file object, and the + read waits for another thread (fsspec's I/O thread, or an HTTP server in + the same process) to fetch the data. If a garbage collection ran in that + thread meanwhile, finalizers of h5py objects (e.g., of an unreferenced + ``NWBHDF5IO``) would wait for the same lock, deadlocking the process. + """ + + def __init__(self, fp: Any) -> None: + global _gc_pauses, _gc_was_enabled + with _gc_lock: + if _gc_pauses == 0: + # Collect pending garbage now, in this thread, which does not + # hold h5py's lock yet + gc.collect() + _gc_was_enabled = gc.isenabled() + gc.disable() + _gc_pauses += 1 + self._fp = fp + self._paused = True + + def __getattr__(self, name: str) -> Any: + return getattr(self._fp, name) + + def __enter__(self) -> _GCPausedFile: + return self + + def __exit__(self, *_exc: Any) -> None: + self.close() + + def __del__(self) -> None: + if self._paused: + self.close() + + def close(self) -> None: + global _gc_pauses + try: + self._fp.close() + finally: + with _gc_lock: + if self._paused: + self._paused = False + _gc_pauses -= 1 + if _gc_pauses == 0 and _gc_was_enabled: + gc.enable() + + @dataclass class AnnexedReadableFile(Readable): """ @@ -95,7 +153,7 @@ class AnnexedReadableFile(Readable): adapter: Any = field(repr=False, compare=False) def open(self) -> IO[bytes]: - return cast("IO[bytes]", self.adapter.open(self.filepath, "rb")) + return cast("IO[bytes]", _GCPausedFile(self.adapter.open(self.filepath, "rb"))) def get_size(self) -> int: return self.size diff --git a/dandi/support/tests/test_datalad_fuse.py b/dandi/support/tests/test_datalad_fuse.py index 6f594eb07..d75b20c32 100644 --- a/dandi/support/tests/test_datalad_fuse.py +++ b/dandi/support/tests/test_datalad_fuse.py @@ -1,5 +1,6 @@ from __future__ import annotations +import gc from pathlib import Path import shutil @@ -56,6 +57,26 @@ def test_annexed_readable( assert server.ranged_requests > 0 +@pytest.mark.ai_generated +def test_annexed_readable_pauses_gc( + streamed_nwb: tuple[Path, RangeHTTPServer], +) -> None: + # Automatic garbage collection is disabled while any streamed file is open + # (to avoid deadlocks with h5py's lock; see `_GCPausedFile`) + nwb, _ = streamed_nwb + r = get_annexed_readable(nwb) + assert r is not None + assert gc.isenabled() + with r.open() as fp: + assert not gc.isenabled() + with r.open() as fp2: + assert fp2.read(8) == b"\x89HDF\r\n\x1a\n" + assert not gc.isenabled() + fp.seek(1) + assert fp.read(3) == b"HDF" + assert gc.isenabled() + + @pytest.mark.ai_generated def test_annexed_readable_none( streamed_nwb: tuple[Path, RangeHTTPServer], simple2_nwb: Path, tmp_path: Path