diff --git a/.github/workflows/run-tests.yml b/.github/workflows/run-tests.yml index fec7c9104..90b812978 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-fuse steps: - name: Set up environment @@ -123,6 +126,15 @@ jobs: git+https://github.com/hdmf-dev/hdmf \ git+https://github.com/hdmf-dev/hdmf-zarr + # TODO: install a release once datalad/datalad-fuse#131 is released + - name: Install git-annex and datalad-fuse (for streaming annexed content) + if: matrix.mode == 'datalad-fuse' + run: | + pip install git-annex \ + "datalad-fuse @ git+https://github.com/datalad/datalad-fuse@refs/pull/131/head" + git config --global user.name "DANDI tests" + git config --global user.email "tests@dandiarchive.org" + - name: Create NFS filesystem if: matrix.mode == 'nfs' run: | diff --git a/dandi/support/datalad_fuse.py b/dandi/support/datalad_fuse.py new file mode 100644 index 000000000..ab6324037 --- /dev/null +++ b/dandi/support/datalad_fuse.py @@ -0,0 +1,160 @@ +""" +Streaming the content of annexed files with datalad-fuse_, as an alternative +to the git-only reader of `dandi.support.annex`. + +datalad-fuse can mount a DataLad dataset with FUSE, but its remote filesystem +adapter, which the mount is built on, also opens files directly, with no mount. +That adapter is what is used here. Unlike `dandi.support.annex`, it asks +git-annex itself where the content is (``git annex whereis``), so it also finds +content available from remotes (e.g., a Forgejo-aneksajo instance or an S3 +export) that has no URL registered. It needs git-annex and DataLad installed, +and is used only in a git-annex repository, i.e., one with git-annex +initialized. + +.. _datalad-fuse: https://github.com/datalad/datalad-fuse +""" + +from __future__ import annotations + +import atexit +from dataclasses import dataclass, field +from datetime import datetime +from functools import cache +from itertools import chain +import logging +import os +from pathlib import Path +import subprocess +from typing import IO, TYPE_CHECKING, cast + +from .annex import AnnexKey +from ..misctypes import Readable + +if TYPE_CHECKING: + from datalad_fuse.adapter import RemoteFilesystemAdapter + +lgr = logging.getLogger("dandi.support.datalad_fuse") + + +@cache +def get_adapter() -> RemoteFilesystemAdapter: + """ + 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.adapter import RemoteFilesystemAdapter + + # caching=False: do not keep the streamed blocks on disk, in the dataset + adapter = RemoteFilesystemAdapter(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, e.g., in a repository with only a ``git-annex`` branch.) + """ + 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 DataladFuseReadableFile(Readable): + """ + A `Readable` for an annexed file whose content is not present locally and + is instead streamed with datalad-fuse's adapter (see the module + docstring). + + Instances are obtained by calling `get_datalad_fuse_readable()`. + """ + + #: The absolute path of the (broken) symbolic link to the file's content + filepath: Path + + #: The git-annex key of the file + key: AnnexKey + + #: 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: RemoteFilesystemAdapter = 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: + assert self.key.size is not None + return self.key.size + + def get_mtime(self) -> datetime | None: + return None + + def get_filename(self) -> str: + return self.filepath.name + + def get_fingerprint(self) -> str: + # The key is a digest of the content (plus its size), so results + # computed from it (validation, metadata) can be cached under it + return self.key.key + + def __str__(self) -> str: + return str(self.filepath) + + +def get_datalad_fuse_readable(path: str | Path) -> DataladFuseReadableFile | None: + """ + Return a `DataladFuseReadableFile` for streaming the content of the + annexed file at ``path``, or `None` if datalad-fuse is not installed, + ``path`` is not in a git-annex repository, is not an annexed file whose + content is missing, its key does not record the size of the content, or + git-annex knows no URL to stream it from. Only locked annexed files + (symbolic links) are handled. + """ + try: + from datalad_fuse.adapter import FileState + except ImportError: + return None + filepath = Path(os.path.abspath(path)) + if not filepath.is_symlink() or filepath.exists(): + # Not a broken symbolic link, i.e., not a locked annexed file whose + # content is missing + return None + if not annex_initialized(filepath.parent): + return None + adapter = get_adapter() + try: + dsap, relpath = adapter.resolve_dataset(filepath) + except Exception as e: + # Not in a Git repository, or one datalad-fuse cannot handle (e.g., + # without any commit) + lgr.debug("%s: Not streamable with datalad-fuse: %s", filepath, e) + return None + if dsap.annex is None: + # git-annex not initialized + return None + state, fkey = dsap.get_file_state(relpath) + if state is not FileState.NO_CONTENT or fkey is None: + return None + key = AnnexKey.parse(str(fkey)) + if key.size is None: + return None + # The URLs tried by the adapter, in its order + url = next( + chain(dsap.get_urls(key.key), dsap.get_exporttree_urls(relpath, fkey)), None + ) + if url is None: + lgr.debug("%s: git-annex knows no URL for key %s", filepath, key) + return None + return DataladFuseReadableFile(filepath=filepath, key=key, url=url, adapter=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..f4fb09d38 --- /dev/null +++ b/dandi/support/tests/test_datalad_fuse.py @@ -0,0 +1,84 @@ +from __future__ import annotations + +from pathlib import Path +import shutil + +import h5py +import pytest + +from .test_annex import RangeHTTPServer, range_http_server # noqa: F401 +from ..annex import AnnexKey +from ..datalad_fuse import DataladFuseReadableFile, get_datalad_fuse_readable +from ...pynwb_utils import get_neurodata_types +from ...tests.fixtures import ( + annex_key_for_file, + make_annexed_dandiset, + make_git_annex_dandiset, +) +from ...tests.skip import mark + +pytestmark = mark.skipif_no_git_annex + + +@pytest.fixture(autouse=True) +def _no_proxy(monkeypatch: pytest.MonkeyPatch) -> None: + # The test HTTP server is local + monkeypatch.setenv("NO_PROXY", "127.0.0.1") + monkeypatch.setenv("no_proxy", "127.0.0.1") + + +@pytest.mark.ai_generated +def test_datalad_fuse_readable( + range_http_server: tuple[RangeHTTPServer, str], # noqa: F811 + simple2_nwb: Path, + tmp_path: Path, +) -> None: + pytest.importorskip("datalad_fuse.adapter") + server, base_url = range_http_server + shutil.copy(simple2_nwb, tmp_path / "content.nwb") + ds = tmp_path / "ds" + make_git_annex_dandiset( + ds, {"sub-01/sub-01.nwb": (simple2_nwb, f"{base_url}/content.nwb")} + ) + nwb = ds / "sub-01" / "sub-01.nwb" + assert not nwb.exists() + + r = get_datalad_fuse_readable(nwb) + assert isinstance(r, DataladFuseReadableFile) + assert r.key == AnnexKey.parse(annex_key_for_file(simple2_nwb)) + assert r.url == f"{base_url}/content.nwb" + assert r.get_size() == simple2_nwb.stat().st_size + assert r.get_filename() == "sub-01.nwb" + assert r.get_fingerprint() == r.key.key + 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 + assert get_neurodata_types.__wrapped__(r) == get_neurodata_types.__wrapped__( + simple2_nwb + ) + + +@pytest.mark.ai_generated +def test_datalad_fuse_readable_none(tmp_path: Path, simple2_nwb: Path) -> None: + pytest.importorskip("datalad_fuse.adapter") + ds = tmp_path / "ds" + make_git_annex_dandiset( + ds, {"sub-01/sub-01.nwb": (simple2_nwb, "http://127.0.0.1:1/content.nwb")} + ) + # The content is present: no need to stream it + (ds / "regular.nwb").write_bytes(b"content") + assert get_datalad_fuse_readable(ds / "regular.nwb") is None + # Not in a Git repository + assert get_datalad_fuse_readable(simple2_nwb) is None + # git-annex not initialized: left to `dandi.support.annex` + fake = tmp_path / "fake" + make_annexed_dandiset( + fake, + {"sub-01/sub-01.nwb": (annex_key_for_file(simple2_nwb), ["http://x/y.nwb"])}, + ) + assert get_datalad_fuse_readable(fake / "sub-01" / "sub-01.nwb") is None diff --git a/dandi/tests/fixtures.py b/dandi/tests/fixtures.py index 0a0926b43..4411ad989 100644 --- a/dandi/tests/fixtures.py +++ b/dandi/tests/fixtures.py @@ -458,6 +458,47 @@ def make_annexed_dandiset( create_git_annex_branch(dandiset, logs) +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. + Unlike `make_annexed_dandiset()`, 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) + + 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/dandi/validate/_core.py b/dandi/validate/_core.py index 6fe3e6699..f268e601a 100644 --- a/dandi/validate/_core.py +++ b/dandi/validate/_core.py @@ -28,6 +28,7 @@ from ..consts import dandiset_metadata_file from ..files import DandiFile, LocalFileAsset, find_dandi_files from ..support.annex import AnnexReadableFile, get_annex_readable +from ..support.datalad_fuse import DataladFuseReadableFile, get_datalad_fuse_readable from ..utils import find_parent_directory_containing BIDS_TO_DANDI = { @@ -281,18 +282,28 @@ def _check_streaming_requirements() -> None: ) -def _prepare_streaming(df: DandiFile) -> AnnexReadableFile | None: +def _prepare_streaming( + df: DandiFile, +) -> AnnexReadableFile | DataladFuseReadableFile | None: """ - Set up streaming of the content of the annexed file represented by ``df`` - from the URLs registered for it in git-annex, and return the `Readable` - that will be used for it, or `None` if the content cannot be streamed - (``df`` is not a file asset, not an annexed file, or has no URLs registered) + Set up streaming of the content of the annexed file represented by ``df``, + and return the `Readable` that will be used for it, or `None` if the + content cannot be streamed (``df`` is not a file asset, not an annexed + file, or no URL is known for it). + + datalad-fuse is used if it is installed and git-annex is initialized in + the repository (see `dandi.support.datalad_fuse`); otherwise, the content + is streamed from the URLs registered in git-annex, read with ``git`` only + (see `dandi.support.annex`). """ if not isinstance(df, LocalFileAsset): return None - readable = get_annex_readable(df.filepath) - if readable is None or not readable.urls: - return None + readable: AnnexReadableFile | DataladFuseReadableFile | None + readable = get_datalad_fuse_readable(df.filepath) + if readable is None: + readable = get_annex_readable(df.filepath) + if readable is None or not readable.urls: + return None df.content_source = readable return readable @@ -300,7 +311,7 @@ def _prepare_streaming(df: DandiFile) -> AnnexReadableFile | None: def _handle_missing_content( df: DandiFile, policy: MissingFileContent, - readable: AnnexReadableFile | None = None, + readable: AnnexReadableFile | DataladFuseReadableFile | None = None, ) -> ValidationResult: """Produce a single :class:`ValidationResult` for a file with missing content. @@ -311,6 +322,10 @@ def _handle_missing_content( if policy == MissingFileContent.stream: if readable is not None: + if isinstance(readable, DataladFuseReadableFile): + source = f"with datalad-fuse, e.g., from {readable.url}" + else: + source = f"from {readable.urls[0]}" return ValidationResult( id="DANDI.FILE_CONTENT_STREAMED", origin=ORIGIN_VALIDATION_DANDI_LAYOUT, @@ -321,7 +336,7 @@ def _handle_missing_content( message=( f"File content is not present locally (git-annex key " f"{readable.key}); content-dependent validation streams " - f"it from {readable.urls[0]}" + f"it {source}" ), ) return ValidationResult( diff --git a/dandi/validate/tests/test_core.py b/dandi/validate/tests/test_core.py index 73d2b5374..64a190873 100644 --- a/dandi/validate/tests/test_core.py +++ b/dandi/validate/tests/test_core.py @@ -21,12 +21,14 @@ from ...consts import dandiset_metadata_file from ...pynwb_utils import validate as pynwb_validate from ...support.annex import AnnexKey, AnnexReadableFile, get_annex_readable +from ...support.tests.test_annex import RangeHTTPServer, range_http_server # noqa: F401 from ...tests.fixtures import ( BIDS_TESTDATA_SELECTION, annex_key_for_file, make_annexed_dandiset, + make_git_annex_dandiset, ) -from ...tests.skip import skipif +from ...tests.skip import mark, skipif def test_validate_nwb_error(simple3_nwb: Path) -> None: @@ -468,6 +470,48 @@ def test_validate_stream_unreadable_url(tmp_path: Path, simple3_nwb: Path) -> No } +@pytest.mark.ai_generated +@mark.skipif_no_git_annex +def test_validate_stream_datalad_fuse( + range_http_server: tuple[RangeHTTPServer, str], # noqa: F811 + simple3_nwb: Path, + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + """In a git-annex repository, stream policy streams content with datalad-fuse.""" + pytest.importorskip("datalad_fuse.adapter") + # The test HTTP server is local + monkeypatch.setenv("NO_PROXY", "127.0.0.1") + monkeypatch.setenv("no_proxy", "127.0.0.1") + server, base_url = range_http_server + shutil.copy(simple3_nwb, tmp_path / "content.nwb") + ds = tmp_path / "ds" + make_git_annex_dandiset( + ds, {"sub-001/sub-001.nwb": (simple3_nwb, f"{base_url}/content.nwb")} + ) + results = list(validate(ds, missing_file_content=MissingFileContent.stream)) + + streamed = [r for r in results if r.id == "DANDI.FILE_CONTENT_STREAMED"] + assert len(streamed) == 1 + assert streamed[0].message is not None + assert f"with datalad-fuse, e.g., from {base_url}/content.nwb" in ( + streamed[0].message + ) + assert server.ranged_requests > 0 + + # Content-dependent validation gives the same results as for a regular + # dandiset containing the file itself + local = tmp_path / "local" + (local / "sub-001").mkdir(parents=True) + shutil.copy(ds / dandiset_metadata_file, local / dandiset_metadata_file) + shutil.copy(simple3_nwb, local / "sub-001" / "sub-001.nwb") + expected = list(validate(local)) + assert any(r.origin.validator == Validator.nwbinspector for r in expected) + assert _content_results(results, "sub-001.nwb") == _content_results( + expected, "sub-001.nwb" + ) + + @pytest.mark.ai_generated def test_validate_stream_requires_fsspec( tmp_path: Path, monkeypatch: pytest.MonkeyPatch diff --git a/docs/source/cmdline/validate.rst b/docs/source/cmdline/validate.rst index 2ac46da9a..1e03cd375 100644 --- a/docs/source/cmdline/validate.rst +++ b/docs/source/cmdline/validate.rst @@ -124,6 +124,12 @@ Notes: remote-tracking counterpart, e.g., ``origin/git-annex``), so the clone must include that branch: do not clone with ``--single-branch`` or a ``--depth`` that excludes it. +- If datalad-fuse_ is installed (with git-annex and DataLad) and git-annex is + initialized in the clone (``git annex init``, which ``datalad clone`` does), + the content is streamed with datalad-fuse's adapter instead (no FUSE mount + is involved). It asks git-annex where the content is, so it also finds + content available from remotes with no URL registered. This currently + requires the development version from datalad/datalad-fuse#131. - Some nwbinspector checks read data arrays (e.g., timestamps), so the amount of data streamed for a file depends on its content; it is nevertheless usually a small fraction of the file. @@ -145,6 +151,7 @@ Notes: the file and folder names of the annexed files themselves. .. _fsspec: https://github.com/fsspec/filesystem_spec +.. _datalad-fuse: https://github.com/datalad/datalad-fuse Development Options diff --git a/pyproject.toml b/pyproject.toml index 512697d71..98b5fa15d 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -214,6 +214,7 @@ follow_imports = "skip" module = [ "bidsschematools.*", "click_didyoumean.*", + "datalad_fuse.*", "etelemetry.*", "fasteners.*", "fscacher.*",