diff --git a/dandi/metadata/nwb.py b/dandi/metadata/nwb.py index 38c2ac4e7..785320687 100644 --- a/dandi/metadata/nwb.py +++ b/dandi/metadata/nwb.py @@ -21,6 +21,7 @@ ignore_benign_pynwb_warnings, metadata_cache, nwb_has_external_links, + readable_fingerprint, ) from ..utils import find_parent_directory_containing @@ -28,10 +29,10 @@ # Disable this for clean hacking -@metadata_cache.memoize_path +@metadata_cache.memoize_path(custom_fingerprint=readable_fingerprint) def get_metadata( path: str | Path | Readable, digest: Digest | None = None -) -> dict | None: +) -> dict[str, Any]: """ Get "flatdata" from a .nwb file diff --git a/dandi/misctypes.py b/dandi/misctypes.py index b50cffc42..9631444e4 100644 --- a/dandi/misctypes.py +++ b/dandi/misctypes.py @@ -284,6 +284,25 @@ def get_filename(self) -> str: """ ... + def get_fingerprint(self) -> str | None: + """ + .. versionadded:: 0.81.0 + + Returns a fingerprint of the resource's own content, such as a content + digest (e.g., a git-annex key), or `None` if none is known + + Two resources with equal fingerprints must have identical bytes. The + fingerprint covers that one file only, not its location nor any other + file that may affect how it is interpreted, such as the BIDS sidecar + ``.json`` files (possibly inherited from parent directories) describing + a ``.nii.gz``. Results computed from the file alone may thus be + shared by resources with equal fingerprints (see + `dandi.pynwb_utils.readable_fingerprint`, which also pairs it with the + file name); results that depend on other files must not be keyed by + it alone. + """ + return None + class LocalReadableFile(Readable): """ @@ -343,9 +362,8 @@ class RemoteReadableAsset(Readable): def open(self) -> IO[bytes]: # Optional dependency: - import fsspec - from aiohttp import ClientTimeout + import fsspec # We need to call open() on the return value of fsspec.open() because # otherwise the filehandle will only be opened when used to enter a diff --git a/dandi/pynwb_utils.py b/dandi/pynwb_utils.py index 297df05f3..3732123f3 100644 --- a/dandi/pynwb_utils.py +++ b/dandi/pynwb_utils.py @@ -73,6 +73,27 @@ ) +def readable_fingerprint(source: Any) -> tuple[str, str] | None: + """ + Content 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`. Paths keep being cached under their location + and ``stat()``; a `Readable` cannot be fingerprinted that way, so one whose + `~Readable.get_fingerprint` returns a value is cached under its file name + and that fingerprint instead. A `Readable` without a fingerprint is handled + as without this: cached by its path if it is path-like (as + `LocalReadableFile` is), not cached at all otherwise. + + Only for functions whose result depends on the file's own content and name + (see `~Readable.get_fingerprint`): those it is applied to here each read a + single NWB file. + """ + if isinstance(source, Readable) and (fp := source.get_fingerprint()) is not None: + return (source.get_filename(), fp) + return None + + def _sanitize_nwb_version( v: Any, filename: str | Path | None = None, @@ -194,7 +215,7 @@ def get_neurodata_types_to_modalities_map() -> dict[str, str]: return ndtypes -@metadata_cache.memoize_path +@metadata_cache.memoize_path(custom_fingerprint=readable_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 +829,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=readable_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/annex.py b/dandi/support/annex.py new file mode 100644 index 000000000..06d9225a8 --- /dev/null +++ b/dandi/support/annex.py @@ -0,0 +1,359 @@ +""" +Support for reading the content of annexed files in git-annex_ repositories +(such as DataLad_ datasets) when that content is not present locally. + +In a git-annex repository, an annexed file is a symbolic link into +:file:`.git/annex/objects/`. Until the content is fetched (e.g., with ``git +annex get`` or ``datalad get``), the link is broken. The name of the link's +target is the file's *key*, which for the ``SHA256E`` and ``MD5E`` backends used +for DANDI Dandisets records the size of the content. The URLs from which the +content can be retrieved are recorded in the ``git-annex`` branch of the +repository (the "URL log" of the key). + +This module reads that information using only ``git`` (git-annex itself is not +required) and exposes the content of such files as a `Readable` that streams it +over HTTP(S) with fsspec_, so that the files of a DataLad Dandiset can be read +(e.g., for their metadata) without downloading the (possibly terabytes of) +data. Its fingerprint is the key, so that results computed from the content +are cached under it. + +.. _git-annex: https://git-annex.branchable.com +.. _DataLad: https://www.datalad.org +.. _fsspec: https://github.com/fsspec/filesystem_spec +""" + +from __future__ import annotations + +from dataclasses import dataclass +from datetime import datetime +import hashlib +import logging +import os +from pathlib import Path +import re +import subprocess +from typing import IO, ClassVar, cast + +from fscacher import annex_key_fingerprint + +from ..misctypes import Readable + +lgr = logging.getLogger("dandi.support.annex") + +#: Regular expression matching DANDI Archive API "download" URLs, which +#: redirect to the actual storage location on every request and are thus less +#: efficient for streaming than the direct URLs also registered in git-annex +DANDI_API_DOWNLOAD_URL_RE = re.compile( + r"/api/(?:dandisets/[^/]+/versions/[^/]+/)?assets/[^/]+/download/?(?:[?#].*)?$" +) + +#: Block size (in bytes) for streaming reads of remote content +STREAM_BLOCK_SIZE = 4 * 1024 * 1024 + +#: Maximum number of blocks of remote content kept in memory per open file +STREAM_MAX_BLOCKS = 64 + + +@dataclass(frozen=True) +class AnnexKey: + """ + A parsed git-annex key, e.g., ``SHA256E-s1234--.nwb`` (see + https://git-annex.branchable.com/internals/key_format/) + """ + + #: The key as a whole + key: str + + #: The key backend, e.g., ``SHA256E`` or ``MD5E`` + backend: str + + #: The size of the content in bytes, if recorded in the key + size: int | None + + #: The backend-specific part of the key following ``--``; for the checksum + #: backends, this is the checksum followed (for the ``*E`` variants) by the + #: file extension + name: str + + @classmethod + def parse(cls, key: str) -> AnnexKey: + """Parse a git-annex key string""" + fields, sep, name = key.partition("--") + backend, *extra = fields.split("-") + if not sep or not name or not backend: + raise ValueError(f"Not a git-annex key: {key!r}") + size = None + for f in extra: + if f[:1] == "s" and f[1:].isdigit(): + size = int(f[1:]) + return cls(key=key, backend=backend, size=size, name=name) + + def __str__(self) -> str: + return self.key + + @property + def hashdir_lower(self) -> str: + """ + The two-level hash directory under which git-annex stores the key's + log files in the ``git-annex`` branch + """ + h = hashlib.md5(self.key.encode("utf-8"), usedforsecurity=False).hexdigest() + return f"{h[:3]}/{h[3:6]}" + + +def get_annex_key(path: str | Path) -> AnnexKey | None: + """ + Return the git-annex key of the locked annexed file at ``path``, or `None` + if ``path`` is not a symbolic link into a git-annex object store or the + key's backend does not hash the content (e.g., ``WORM`` or ``URL`` keys), + in which case the key cannot fingerprint the content (see + `fscacher.annex_key_fingerprint`, which this uses) + """ + key = annex_key_fingerprint(path, pair_with_path=False) + return None if key is None else AnnexKey.parse(key) + + +def parse_url_log(text: str) -> list[str]: + """ + Parse the content of a git-annex URL log (a ``.log.web`` file in the + ``git-annex`` branch), whose lines are of the form ``s + ``, and return the URLs whose most recent status is ``1`` (present). + + Direct URLs are sorted before DANDI Archive API download URLs; otherwise, + the order is that of first appearance in the log. + """ + latest: dict[str, tuple[float, str]] = {} + for line in text.splitlines(): + try: + ts, status, url = line.split(" ", 2) + t = float(ts.rstrip("s")) + except ValueError: + lgr.debug("Ignoring unparsable git-annex URL log line: %r", line) + continue + if url not in latest or t > latest[url][0]: + latest[url] = (t, status) + urls = [url for url, (_, status) in latest.items() if status == "1"] + return sorted(urls, key=lambda u: DANDI_API_DOWNLOAD_URL_RE.search(u) is not None) + + +class AnnexRepo: + """ + Read-only access to the git-annex metadata (the ``git-annex`` branch) of a + Git repository, using only ``git``. Instances are obtained with `find()` + and are cached per repository root. + """ + + _instances: ClassVar[dict[Path, AnnexRepo]] = {} + + def __init__(self, root: Path) -> None: + #: The root (top-level directory) of the repository + self.root = root + self._annex_ref: str | None = None + self._annex_ref_resolved = False + self._urls: dict[str, list[str]] = {} + + def __repr__(self) -> str: + return f"{type(self).__name__}({str(self.root)!r})" + + @classmethod + def find(cls, path: str | Path) -> AnnexRepo | None: + """ + Return the `AnnexRepo` for the repository containing ``path`` (the + closest parent directory containing a ``.git`` directory or file), or + `None` if ``path`` is not inside a Git repository + """ + p = Path(os.path.abspath(path)) + for d in (p, *p.parents): + if os.path.lexists(d / ".git"): + return cls._instances.setdefault(d, cls(d)) + return None + + def _git(self, *args: str) -> str | None: + """ + Run a ``git`` command in the repository and return its standard output, + or `None` if it failed + """ + try: + r = subprocess.run( + ["git", "-C", str(self.root), *args], + capture_output=True, + text=True, + encoding="utf-8", + errors="replace", + ) + except FileNotFoundError: + lgr.warning( + "git is not installed; cannot read git-annex metadata of %s", self.root + ) + return None + if r.returncode != 0: + lgr.debug( + "git %s failed in %s: %s", " ".join(args), self.root, r.stderr.strip() + ) + return None + return r.stdout + + @property + def annex_ref(self) -> str | None: + """ + The ref of the ``git-annex`` branch: the local branch if there is one, + otherwise a remote-tracking branch (preferring that of ``origin``); or + `None` if the repository has no git-annex metadata + """ + if not self._annex_ref_resolved: + self._annex_ref_resolved = True + out = self._git( + "for-each-ref", + "--format=%(refname)", + "refs/heads/git-annex", + "refs/remotes/*/git-annex", + ) + if out is not None and (refs := out.split()): + for preferred in ( + "refs/heads/git-annex", + "refs/remotes/origin/git-annex", + ): + if preferred in refs: + self._annex_ref = preferred + break + else: + self._annex_ref = sorted(refs)[0] + return self._annex_ref + + def get_urls(self, key: AnnexKey | str) -> list[str]: + """ + Return the URLs registered in git-annex for the given key (see + `parse_url_log()` for their order), or an empty list if there are none + or the git-annex metadata cannot be read + """ + if isinstance(key, str): + key = AnnexKey.parse(key) + if (urls := self._urls.get(key.key)) is None: + urls = [] + if (ref := self.annex_ref) is not None: + out = self._git( + "cat-file", "-p", f"{ref}:{key.hashdir_lower}/{key.key}.log.web" + ) + if out is not None: + urls = parse_url_log(out) + self._urls[key.key] = urls + return urls + + +@dataclass +class AnnexReadableFile(Readable): + """ + A `Readable` for an annexed file whose content is not present locally and + is instead streamed from the URL(s) registered for it in git-annex. The + fsspec_ library must be installed with the ``http`` extra (e.g., ``pip + install "dandi[extras]"``) in order for `.open()` to be usable. + + Instances are obtained by calling `get_annex_readable()`. + + .. _fsspec: http://github.com/fsspec/filesystem_spec + """ + + #: The path of the (broken) symbolic link to the file's content + filepath: Path + + #: The git-annex key of the file + key: AnnexKey + + #: The URLs from which the content can be retrieved, in order of preference + urls: list[str] + + def open(self) -> IO[bytes]: + """ + Open the content for random-access reading from the first URL that can + be opened, streaming it in blocks on demand + """ + # Optional dependency: + from aiohttp import ClientTimeout + import fsspec + from fsspec.caching import caches as fsspec_caches + + if not self.urls: + raise RuntimeError(f"{self.filepath}: No URLs registered in git-annex") + # fsspec's LRU block cache (which suits h5py's random access) was + # registered as "block" before it was renamed to "blockcache" in 2023 + cache_type = "blockcache" if "blockcache" in fsspec_caches else "block" + # fsspec logs every block it fetches at INFO level, which is too noisy + # for the (INFO-level by default) dandi CLI output; quiet that unless + # the user configured that logger themselves + if (fsspec_lgr := logging.getLogger("fsspec.caching")).level == logging.NOTSET: + fsspec_lgr.setLevel(logging.WARNING) + error: Exception | None = None + for url in self.urls: + lgr.debug("%s: Opening %s for streaming", self.filepath, url) + try: + # We need to call open() on the return value of fsspec.open() + # because otherwise the filehandle will only be opened when + # used to enter a context manager. + return cast( + IO[bytes], + fsspec.open( + url, + mode="rb", + block_size=STREAM_BLOCK_SIZE, + cache_type=cache_type, + cache_options={"maxblocks": STREAM_MAX_BLOCKS}, + client_kwargs={ + # Explicit timeouts prevent indefinite hangs in + # fsspec's sync() wrapper on a stalled connection; + # see the same in `RemoteReadableAsset.open()`. + "timeout": ClientTimeout( + total=120, sock_read=60, sock_connect=30 + ), + # Honor HTTP(S)_PROXY etc. environment variables + "trust_env": True, + }, + ).open(), + ) + except Exception as e: + lgr.warning( + "%s: Could not open %s for streaming: %s: %s", + self.filepath, + url, + type(e).__name__, + e, + ) + error = e + assert error is not None + raise error + + def get_size(self) -> int: + if self.key.size is None: + raise ValueError(f"git-annex key {self.key} does not record a size") + 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_annex_readable(path: str | Path) -> AnnexReadableFile | None: + """ + Return an `AnnexReadableFile` for streaming the content of the annexed file + at ``path``, or `None` if ``path`` is not a symbolic link into a git-annex + object store, or its key does not pin the content (see `get_annex_key()`) + or does not record its size. The + returned object's ``urls`` list is empty if no URLs are registered for the + file (or the repository's git-annex metadata cannot be read). + """ + filepath = Path(path) + key = get_annex_key(filepath) + if key is None or key.size is None: + return None + repo = AnnexRepo.find(filepath.parent) + urls = repo.get_urls(key) if repo is not None else [] + return AnnexReadableFile(filepath=filepath, key=key, urls=urls) diff --git a/dandi/support/tests/test_annex.py b/dandi/support/tests/test_annex.py new file mode 100644 index 000000000..a65dccaac --- /dev/null +++ b/dandi/support/tests/test_annex.py @@ -0,0 +1,382 @@ +from __future__ import annotations + +from collections.abc import Iterator +from functools import partial +from http import HTTPStatus +from http.server import SimpleHTTPRequestHandler, ThreadingHTTPServer +import os +from pathlib import Path +import re +import shutil +from subprocess import run +import threading +from typing import Any + +from fscacher import PersistentCache +import h5py +import pytest + +from ..annex import ( + AnnexKey, + AnnexReadableFile, + AnnexRepo, + get_annex_key, + get_annex_readable, + parse_url_log, +) +from ...pynwb_utils import get_neurodata_types, readable_fingerprint +from ...tests.fixtures import ( + annex_key_for_file, + create_git_annex_branch, + make_annexed_dandiset, +) +from ...tests.skip import skipif + +# The key of an actual asset of https://github.com/dandisets/000003 +REAL_KEY = ( + "SHA256E-s61510864725--" + "e063aec023141b5a11fb95a7605fb31057aab05944a81f37c5c6211739ed7f62.nwb" +) + + +@pytest.mark.ai_generated +@pytest.mark.parametrize( + "key,backend,size,name", + [ + ( + REAL_KEY, + "SHA256E", + 61510864725, + "e063aec023141b5a11fb95a7605fb31057aab05944a81f37c5c6211739ed7f62.nwb", + ), + ( + "MD5E-s123--0123456789abcdef0123456789abcdef.nwb", + "MD5E", + 123, + "0123456789abcdef0123456789abcdef.nwb", + ), + ("WORM-s512-m1700000000--foo.txt", "WORM", 512, "foo.txt"), + ("URL--http&c%%example.com%file", "URL", None, "http&c%%example.com%file"), + ], +) +def test_annex_key_parse(key: str, backend: str, size: int | None, name: str) -> None: + k = AnnexKey.parse(key) + assert (k.backend, k.size, k.name) == (backend, size, name) + assert str(k) == key + + +@pytest.mark.ai_generated +@pytest.mark.parametrize( + "key", ["", "SHA256E-s123", "--abc", "SHA256E-s1--", "-s1--abc"] +) +def test_annex_key_parse_invalid(key: str) -> None: + with pytest.raises(ValueError): + AnnexKey.parse(key) + + +@pytest.mark.ai_generated +def test_annex_key_hashdir_lower() -> None: + # This is where the key's log files are located in the git-annex branch of + # https://github.com/dandisets/000003 + assert AnnexKey.parse(REAL_KEY).hashdir_lower == "595/864" + + +@pytest.mark.ai_generated +def test_get_annex_key(tmp_path: Path) -> None: + (tmp_path / "regular.nwb").write_bytes(b"data") + assert get_annex_key(tmp_path / "regular.nwb") is None + assert get_annex_key(tmp_path / "nonexistent.nwb") is None + (tmp_path / "other.nwb").symlink_to("somewhere/else.nwb") + assert get_annex_key(tmp_path / "other.nwb") is None + (tmp_path / "notakey.nwb").symlink_to(".git/annex/objects/Xx/Yy/notakey/notakey") + assert get_annex_key(tmp_path / "notakey.nwb") is None + (tmp_path / "annexed.nwb").symlink_to( + f"../.git/annex/objects/Xx/Yy/{REAL_KEY}/{REAL_KEY}" + ) + key = get_annex_key(tmp_path / "annexed.nwb") + assert key is not None + assert key.key == REAL_KEY + assert key.size == 61510864725 + # Keys that do not pin the content cannot be used to fingerprint it + for i, other_key in enumerate( + ["WORM-s4-m1700000000--file.nwb", "URL-s4--https&c%%example.com%file.nwb"] + ): + (tmp_path / f"unpinned{i}.nwb").symlink_to( + f"../.git/annex/objects/Xx/Yy/{other_key}/{other_key}" + ) + assert get_annex_key(tmp_path / f"unpinned{i}.nwb") is None + + +@pytest.mark.ai_generated +def test_parse_url_log() -> None: + s3_url = "https://dandiarchive.s3.amazonaws.com/blobs/9d7/5f6/uuid?versionId=abc" + api_url = "https://api.dandiarchive.org/api/assets/25564f6b/download/" + old_api_url = ( + "https://api.dandiarchive.org/api/dandisets/000003/versions/draft" + "/assets/25564f6b/download/" + ) + log = ( + f"1620050417.649816s 1 {s3_url}\n" + f"1630339841.045183s 1 {api_url}\n" + f"1630339840.983847s 0 {old_api_url}\n" + "1620050418.205986s 0 https://dandiarchive.s3.amazonaws.com/girder/71/54/old\n" + # A URL that was present and then removed, with the lines out of order + "1700000002s 0 https://example.com/removed\n" + "1700000001s 1 https://example.com/removed\n" + # A URL that was removed and then restored + "1700000001s 0 https://example.com/restored\n" + "1700000002s 1 https://example.com/restored\n" + "garbage line\n" + ) + assert parse_url_log(log) == [ + s3_url, + "https://example.com/restored", + # DANDI API download URLs come last + api_url, + ] + assert parse_url_log("") == [] + + +def _git(repo: Path, *args: str) -> None: + 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", + } + run(["git", "-C", str(repo), *args], check=True, env=env) + + +@pytest.mark.ai_generated +def test_annex_repo_get_urls(tmp_path: Path) -> None: + skipif.no_git() + repo = tmp_path / "repo" + repo.mkdir() + _git(repo, "init", "-q") + _git(repo, "commit", "-q", "--allow-empty", "-m", "Initial commit") + key = AnnexKey.parse(REAL_KEY) + + found = AnnexRepo.find(repo / "sub-01" / "sub-01.nwb") + assert found is not None + assert found.root == repo + assert AnnexRepo.find(repo) is found + assert AnnexRepo.find(tmp_path / "norepo" / "file.nwb") is None + + # Without a git-annex branch there are no URLs + annex_repo = AnnexRepo(repo) + assert annex_repo.annex_ref is None + assert annex_repo.get_urls(key) == [] + + create_git_annex_branch( + repo, + { + f"{key.hashdir_lower}/{REAL_KEY}.log.web": ( + "1s 1 https://api.dandiarchive.org/api/assets/x/download/\n" + "2s 1 https://example.com/direct\n" + ) + }, + ) + annex_repo = AnnexRepo(repo) + assert annex_repo.annex_ref == "refs/heads/git-annex" + urls = [ + "https://example.com/direct", + "https://api.dandiarchive.org/api/assets/x/download/", + ] + assert annex_repo.get_urls(key) == urls + assert annex_repo.get_urls(REAL_KEY) == urls + assert annex_repo.get_urls("MD5E-s1--00000000000000000000000000000000.nwb") == [] + + # A plain `git clone` only has the remote-tracking git-annex branch + clone = tmp_path / "clone" + _git(tmp_path, "clone", "-q", str(repo), str(clone)) + cloned = AnnexRepo(clone) + assert cloned.annex_ref == "refs/remotes/origin/git-annex" + assert cloned.get_urls(key) == urls + + +@pytest.mark.ai_generated +def test_get_annex_readable(tmp_path: Path, simple2_nwb: Path) -> None: + skipif.no_git() + ds = tmp_path / "ds" + key = annex_key_for_file(simple2_nwb) + make_annexed_dandiset( + ds, + { + "sub-01/sub-01.nwb": (key, [simple2_nwb.as_uri()]), + "sub-02/sub-02.nwb": ("SHA256E-s42--" + "0" * 64 + ".nwb", []), + "sub-03/sub-03.nwb": ("URL--http&c%%example.com%sub-03.nwb", ["x"]), + }, + ) + r = get_annex_readable(ds / "sub-01" / "sub-01.nwb") + assert r is not None + assert r.key.key == key + assert r.urls == [simple2_nwb.as_uri()] + assert r.get_size() == simple2_nwb.stat().st_size + assert r.get_filename() == "sub-01.nwb" + assert r.get_fingerprint() == r.key.key + assert str(r) == str(ds / "sub-01" / "sub-01.nwb") + # No URLs registered + r = get_annex_readable(ds / "sub-02" / "sub-02.nwb") + assert r is not None + assert r.urls == [] + # Key without a size + assert get_annex_readable(ds / "sub-03" / "sub-03.nwb") is None + # Not an annexed file + assert get_annex_readable(ds / "dandiset.yaml") is None + + +@pytest.mark.ai_generated +def test_annex_readable_file_open_file_url(tmp_path: Path) -> None: + pytest.importorskip("fsspec") + content = tmp_path / "content.bin" + content.write_bytes(b"0123456789" * 100) + key = AnnexKey.parse(annex_key_for_file(content)) + assert key.size == 1000 + filepath = tmp_path / "ds" / "content.bin" + missing = (tmp_path / "missing.bin").as_uri() + r = AnnexReadableFile(filepath=filepath, key=key, urls=[missing, content.as_uri()]) + assert r.get_size() == 1000 + assert r.get_mtime() is None + assert r.get_filename() == "content.bin" + assert r.get_fingerprint() == key.key + # The first URL cannot be opened, so the second one is used + fp = r.open() + try: + assert fp.read(10) == b"0123456789" + fp.seek(995) + assert fp.read() == b"56789" + finally: + fp.close() + with pytest.raises(FileNotFoundError): + AnnexReadableFile(filepath=filepath, key=key, urls=[missing]).open() + with pytest.raises(RuntimeError): + AnnexReadableFile(filepath=filepath, key=key, urls=[]).open() + with pytest.raises(ValueError): + AnnexReadableFile( + filepath=filepath, key=AnnexKey.parse("URL--x"), urls=[] + ).get_size() + + +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`` over HTTP with support for range requests""" + server = RangeHTTPServer( + ("127.0.0.1", 0), partial(RangeHTTPRequestHandler, directory=str(tmp_path)) + ) + 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() + + +@pytest.mark.ai_generated +def test_annex_readable_file_cached_by_key(tmp_path: Path, simple2_nwb: Path) -> None: + pytest.importorskip("fsspec") + cache = PersistentCache(path=tmp_path / "cache") + neurodata_types = cache.memoize_path(custom_fingerprint=readable_fingerprint)( + get_neurodata_types.__wrapped__ + ) + content = tmp_path / "content.nwb" + shutil.copy(simple2_nwb, content) + key = AnnexKey.parse(annex_key_for_file(content)) + expected = get_neurodata_types.__wrapped__(content) + streamed = AnnexReadableFile( + filepath=tmp_path / "ds" / "sub-01.nwb", key=key, urls=[content.as_uri()] + ) + assert neurodata_types(streamed) == expected + # A file with the same key and name elsewhere (e.g., in another clone) is + # served from the cache, without its content being read (it could not be) + twin = AnnexReadableFile( + filepath=tmp_path / "clone" / "sub-01.nwb", key=key, urls=[] + ) + with pytest.raises(RuntimeError): + twin.open() + assert neurodata_types(twin) == expected + + +@pytest.mark.ai_generated +def test_annex_readable_file_open_http( + range_http_server: tuple[RangeHTTPServer, str], simple2_nwb: Path, tmp_path: Path +) -> None: + pytest.importorskip("fsspec") + pytest.importorskip("aiohttp") + server, base_url = range_http_server + shutil.copy(simple2_nwb, tmp_path / "content.nwb") + key = AnnexKey.parse(annex_key_for_file(tmp_path / "content.nwb")) + r = AnnexReadableFile( + filepath=tmp_path / "ds" / "sub-01" / "sub-01.nwb", + key=key, + # The first URL does not exist, so the second one is used + urls=[f"{base_url}/missing.nwb", f"{base_url}/content.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"] + assert "session_description" in h5 + finally: + fp.close() + assert server.ranged_requests > 0 diff --git a/dandi/tests/fixtures.py b/dandi/tests/fixtures.py index 9b128a23d..0a0926b43 100644 --- a/dandi/tests/fixtures.py +++ b/dandi/tests/fixtures.py @@ -3,6 +3,7 @@ from collections.abc import Callable, Iterator from dataclasses import dataclass, field, replace from datetime import datetime, timezone +import hashlib import logging import os from pathlib import Path @@ -10,7 +11,7 @@ import shutil from subprocess import DEVNULL, check_output, run from time import sleep -from typing import Any, Literal +from typing import IO, Any, Literal, cast from uuid import uuid4 from click.testing import CliRunner @@ -38,7 +39,9 @@ metadata_nwb_file_fields, ) from ..dandiapi import DandiAPIClient, RemoteDandiset +from ..misctypes import LocalReadableFile from ..pynwb_utils import make_nwb_file +from ..support.annex import AnnexKey from ..upload import upload lgr = get_logger() @@ -94,6 +97,27 @@ def simple1_nwb_metadata() -> dict[str, Any]: return metadata +class FingerprintedReadable(LocalReadableFile): + """ + A local file posing as a `Readable` with a content fingerprint of its own + + Opening it is counted, to tell results served from a cache from those + computed from the content. + """ + + def __init__(self, filepath: str | Path, fingerprint: str | None) -> None: + super().__init__(filepath) + self.fingerprint = fingerprint + self.opened = 0 + + def open(self) -> IO[bytes]: + self.opened += 1 + return super().open() + + def get_fingerprint(self) -> str | None: + return self.fingerprint + + @pytest.fixture(scope="session") def simple1_nwb( simple1_nwb_metadata: dict[str, Any], tmp_path_factory: pytest.TempPathFactory @@ -352,6 +376,88 @@ def fixture( return fixture +def annex_key_for_file(path: Path) -> str: + """ + Return the git-annex key (using the ``SHA256E`` backend, as DANDI Dandisets + do) that git-annex would assign to the file at ``path`` + """ + digest = hashlib.sha256(path.read_bytes()).hexdigest() + return f"SHA256E-s{path.stat().st_size}--{digest}{path.suffix}" + + +def create_git_annex_branch( + repo: Path, files: dict[str, str], ref: str = "refs/heads/git-annex" +) -> None: + """ + Create (or replace) the ``git-annex`` branch of the Git repository at + ``repo`` so that its tree consists of the given files (a mapping from paths + to contents), using only ``git`` plumbing commands, i.e., without needing + git-annex. This is for testing code that reads git-annex metadata. + """ + index = repo / ".git" / "dandi-test-index" + env = { + **os.environ, + "GIT_INDEX_FILE": str(index), + "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, **kwargs: Any) -> str: + out = check_output( + ["git", "-C", str(repo), *args], env=env, text=True, **kwargs + ) + return cast(str, out).strip() + + for path, content in files.items(): + blob = git("hash-object", "-w", "--stdin", input=content) + git("update-index", "--add", "--cacheinfo", f"100644,{blob},{path}") + tree = git("write-tree") + commit = git("commit-tree", tree, "-m", "git-annex metadata") + git("update-ref", ref, commit) + index.unlink(missing_ok=True) + + +def make_annexed_dandiset( + dandiset: Path, files: dict[str, tuple[str, list[str]]] +) -> None: + """ + Make the directory ``dandiset`` (created if needed, with a minimal + :file:`dandiset.yaml` if it does not have one) look like a DataLad Dandiset + whose annexed content has not been fetched: for each ``relpath: (key, + urls)`` item of ``files``, ``relpath`` is created as a broken symbolic link + into the git-annex object store for ``key``, and the ``urls`` are + registered for ``key`` in the ``git-annex`` branch of the repository. + """ + dandiset.mkdir(parents=True, exist_ok=True) + metadata_file = dandiset / dandiset_metadata_file + if not metadata_file.exists(): + metadata_file.write_text( + "identifier: '000027'\nname: Test\ndescription: Test dandiset\n" + ) + if not (dandiset / ".git").exists(): + run(["git", "init", "-q", str(dandiset)], check=True) + logs: dict[str, str] = {} + for relpath, (key, urls) in files.items(): + link = dandiset / relpath + link.parent.mkdir(parents=True, exist_ok=True) + link.symlink_to( + Path(os.path.relpath(dandiset, link.parent)) + / ".git" + / "annex" + / "objects" + / "Xx" + / "Yy" + / key + / key + ) + logs[f"{AnnexKey.parse(key).hashdir_lower}/{key}.log.web"] = "".join( + f"{1700000000 + i}s 1 {url}\n" for i, url in enumerate(urls) + ) + create_git_annex_branch(dandiset, logs) + + 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/test_metadata.py b/dandi/tests/test_metadata.py index f6b40ddaa..020fa32f1 100644 --- a/dandi/tests/test_metadata.py +++ b/dandi/tests/test_metadata.py @@ -4,9 +4,11 @@ from datetime import datetime, timedelta from itertools import chain import json +import os from pathlib import Path import shutil from typing import Any +from uuid import uuid4 from anys import ANY_AWARE_DATETIME, ANY_INT, AnyFullmatch, AnyIn from dandischema.consts import DANDI_SCHEMA_VERSION @@ -39,7 +41,7 @@ import requests from semantic_version import Version -from .fixtures import SampleDandiset +from .fixtures import FingerprintedReadable, SampleDandiset from .skip import mark from .. import __version__ from ..consts import metadata_nwb_subject_fields @@ -1468,3 +1470,21 @@ def test_nwb2asset_remote_asset(nwb_dandiset: SampleDandiset) -> None: approach=[], relatedResource=[], ) + + +@pytest.mark.ai_generated +@pytest.mark.skipif( + os.environ.get("DANDI_CACHE") == "ignore", reason="the metadata cache is disabled" +) +def test_get_metadata_fingerprinted_readable(simple1_nwb: Path, tmp_path: Path) -> None: + # Unique per run: the module-level cache persists across test runs + fingerprint = f"test-{uuid4()}" + first = FingerprintedReadable(simple1_nwb, fingerprint) + metadata = get_metadata(first) + assert metadata == get_metadata(simple1_nwb) + assert first.opened > 0 + # Its twin is served from the cache by the fingerprint: never even opened, + # which it could not be + twin = FingerprintedReadable(tmp_path / "gone" / simple1_nwb.name, fingerprint) + assert get_metadata(twin) == metadata + assert twin.opened == 0 diff --git a/dandi/tests/test_pynwb_utils.py b/dandi/tests/test_pynwb_utils.py index 52b49367e..d72742e0e 100644 --- a/dandi/tests/test_pynwb_utils.py +++ b/dandi/tests/test_pynwb_utils.py @@ -2,21 +2,29 @@ from collections.abc import Callable from datetime import datetime, timezone +import os from pathlib import Path import re +import shutil +import time from types import SimpleNamespace from typing import Any, NoReturn +from fscacher import PersistentCache import h5py import numpy as np -import pytest from pynwb import NWBHDF5IO, NWBFile, TimeSeries +import pytest from pytest_mock import MockerFixture +from .fixtures import FingerprintedReadable +from ..misctypes import Readable from ..pynwb_utils import ( _rename_pose_estimation_original_videos, _sanitize_nwb_version, nwb_has_external_links, + open_readable, + readable_fingerprint, rename_nwb_external_files, ) @@ -66,7 +74,11 @@ def search(v: str) -> None: def test_rename_pose_estimation_original_videos() -> None: pose = SimpleNamespace( neurodata_type="PoseEstimation", - original_videos=[b"camera\\raw.mp4", "https://example.com/remote.mp4", "other.mp4"], + original_videos=[ + b"camera\\raw.mp4", + "https://example.com/remote.mp4", + "other.mp4", + ], ) unrelated = SimpleNamespace( neurodata_type="OtherContainer", original_videos=["camera/raw.mp4"] @@ -116,9 +128,7 @@ def test_rename_pose_estimation_original_videos_persists_hdf5( with h5py.File(filepath, "w") as f: f.create_dataset( "original_videos", - data=np.asarray( - ["camera/raw.mp4", "camera/other.mp4"], dtype=string_type - ), + data=np.asarray(["camera/raw.mp4", "camera/other.mp4"], dtype=string_type), ) with h5py.File(filepath, "r+") as f: @@ -228,3 +238,62 @@ def test_nwb_has_external_links(tmp_path): assert not nwb_has_external_links(filename1) assert nwb_has_external_links(filename4) + + +@pytest.mark.ai_generated +def test_readable_fingerprint(tmp_path: Path, simple1_nwb: Path) -> None: + assert readable_fingerprint(simple1_nwb) is None + assert readable_fingerprint(str(simple1_nwb)) is None + assert readable_fingerprint(FingerprintedReadable(simple1_nwb, None)) is None + assert readable_fingerprint(FingerprintedReadable(simple1_nwb, "A")) == ( + simple1_nwb.name, + "A", + ) + + +@pytest.mark.ai_generated +def test_memoize_path_readable_fingerprint(tmp_path: Path, simple1_nwb: Path) -> None: + cache = PersistentCache(path=tmp_path / "cache", tokens=["t1"]) + calls: list[Any] = [] + + @cache.memoize_path(custom_fingerprint=readable_fingerprint) + def size(source: str | Path | Readable, flag: bool = False) -> str: + calls.append(source) + with open_readable(source) as fp: + return f"{len(fp.read())}:{flag}" + + expected = f"{simple1_nwb.stat().st_size}:False" + + # A path is still cached by its stat() (which skips a file modified "just + # now", as the session-wide fixture may well have been: age a copy) + nwb = tmp_path / simple1_nwb.name + shutil.copyfile(simple1_nwb, nwb) + hour_ago = time.time() - 3600 + os.utime(nwb, (hour_ago, hour_ago)) + assert size(nwb) == expected + assert size(nwb) == expected, ( + "a repeated call on the unchanged file must return the same result," + " served from the cache" + ) + assert len(calls) == 1 + + # A Readable without a fingerprint is handled as before: a local one is + # path-like, so it is cached by the stat() of its path, sharing the entry + local = FingerprintedReadable(nwb, None) + assert size(local) == expected + assert local.opened == 0 + assert len(calls) == 1 + + # A Readable with a fingerprint is cached by it: its twin is served from + # the cache without being read (it could not be) + first = FingerprintedReadable(simple1_nwb, "A") + assert size(first) == expected + assert first.opened == 1 + twin = FingerprintedReadable(tmp_path / "gone" / simple1_nwb.name, "A") + assert size(twin) == expected + assert twin.opened == 0 + assert len(calls) == 2 + + # ... but the file name is part of the key + with pytest.raises(FileNotFoundError): + size(FingerprintedReadable(tmp_path / "gone" / "other.nwb", "A")) diff --git a/pyproject.toml b/pyproject.toml index b98d4f7a0..512697d71 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",