Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions .github/workflows/run-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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: |
Expand Down
160 changes: 160 additions & 0 deletions dandi/support/datalad_fuse.py
Original file line number Diff line number Diff line change
@@ -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)
84 changes: 84 additions & 0 deletions dandi/support/tests/test_datalad_fuse.py
Original file line number Diff line number Diff line change
@@ -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
41 changes: 41 additions & 0 deletions dandi/tests/fixtures.py
Original file line number Diff line number Diff line change
Expand Up @@ -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():
Expand Down
5 changes: 5 additions & 0 deletions dandi/tests/skip.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down Expand Up @@ -162,6 +166,7 @@ def on_windows():
no_docker_commands,
no_docker_engine,
no_git,
no_git_annex,
no_network,
# no_singularity,
no_ssh,
Expand Down
Loading
Loading