From b8dbe4c3aeebcf15e2142f998861f943c4939222 Mon Sep 17 00:00:00 2001 From: Alex Rubinsteyn Date: Wed, 30 Sep 2026 15:47:14 -0400 Subject: [PATCH 1/3] Add transactional archive directory installs --- .github/workflows/tests.yml | 16 + CHANGELOG.md | 22 + README.md | 12 +- datacache/__init__.py | 6 + datacache/archives.py | 920 ++++++++++++++++++++++++++++++++++++ datacache/file_registry.py | 2 + datacache/version.py | 2 +- docs/api.md | 144 +++++- docs/archives.md | 197 ++++++++ docs/bundles.md | 5 + docs/file_registry.md | 2 +- docs/integration.md | 33 +- tests/test_archives.py | 563 ++++++++++++++++++++++ tests/test_file_registry.py | 8 +- 14 files changed, 1914 insertions(+), 18 deletions(-) create mode 100644 datacache/archives.py create mode 100644 docs/archives.md create mode 100644 tests/test_archives.py diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 6714edc..9f3ad89 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -57,3 +57,19 @@ jobs: # error. Don't let that flake fail the whole test run. continue-on-error: true uses: coverallsapp/github-action@v2.2.3 + + archive-windows: + runs-on: windows-latest + steps: + - name: Checkout repository + uses: actions/checkout@v4 + - name: Set up Python 3.12 + uses: actions/setup-python@v5 + with: + python-version: "3.12" + - name: Install package and test dependencies + run: | + python -m pip install --upgrade pip + python -m pip install -e ".[test]" + - name: Test cross-platform archive publication + run: python -m pytest -q tests/test_archives.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 9dc509f..e5d4dff 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,27 @@ # Changelog +## 1.16.0 + +- Add `install_archive` and `inspect_archive` for safely installing a complete + tar archive tree as one immutable generation. Ordered split archives are + concatenated byte-for-byte; URL and local-path sources share the same API. + Extraction rejects traversal, links, special files, path collisions and + reserved metadata names, and ignores archive ownership and permission bits. +- Add `VersionedArchiveRegistry`, a catalogue facade with pinned defaults, + explicit versions, reusable status rows, fast local path resolution, consumer + layout callbacks, local predownloaded-part overrides, and the same progress, + retry, verification and recovery behavior as `install_archive`. +- Include the recorded URL and observed SHA-256 in `VersionedFileRegistry` + status rows so fixed-file consumers can report the same download provenance. +- Publish consumer metadata such as MHCflurry's `DOWNLOAD_INFO.csv` inside the + same transaction with `extra_files`. Archive receipts record the ordered + source identities, assembled and per-part observed hashes, and every extracted + file. Trusted expectations remain distinct from observed consistency checks. +- Serialize archive writers with a cross-platform file lock, retain immutable + generations across refreshes, and recover a complete generation after an + interrupted pointer publication. Legacy or otherwise foreign directories are + never claimed, including with `force=True` (#85). + ## 1.15.0 - Add `VersionedFileRegistry` for established single-file caches with fixed diff --git a/README.md b/README.md index 3d94b46..f7884fe 100644 --- a/README.md +++ b/README.md @@ -5,8 +5,9 @@ Download, verify, transform, and cache datasets for Python applications, including OpenVax libraries such as pyensembl. DataCache provides streaming -downloads, gzip/ZIP decompression, reusable local paths, offline inspection, -and SQLite caches built from pandas DataFrames. +downloads, gzip/ZIP decompression, transactional tar-tree installation, +reusable local paths, offline inspection, and SQLite caches built from pandas +DataFrames. ## Install @@ -162,8 +163,10 @@ lists every public signature, default, return value, exception, and example. | Task | API | Result | | --- | --- | --- | | Install or reuse a versioned dataset | `VersionedDatasetRegistry`, `install_bundle(...)` | Mapping of asset names to snapshot paths | +| Install versioned archive trees | `VersionedArchiveRegistry`, `install_archive(...)` | Immutable extracted-generation `Path` | | Reuse an established fixed-path versioned file cache | `VersionedFileRegistry` | One Path and a legacy-compatible root receipt | | Inspect a complete dataset generation | `inspect_bundle(...)` | `BundleInspection` | +| Inspect an installed archive tree | `inspect_archive(...)` | `ArchiveInspection` | | Discard retained partial download bytes | `discard_partial(destination)` | Installed file unchanged | | Download or reuse one file | `fetch_file(...)`, `Cache.fetch(...)` | Local path string | | Compute a path without filesystem access | `expected_path(...)`, `Cache.local_path(...)` | Path string | @@ -186,6 +189,8 @@ FASTA files with `fetch_file`, then parse them in the consuming library. - [Versioned datasets and bundles](docs/bundles.md): pinned versions, atomic installation, offline recovery, path lifetime, and downstream adoption. +- [Archive directory installation](docs/archives.md): safe tar extraction, split archives, + transactional consumer receipts, and legacy-directory migration. - [API reference](https://github.com/openvax/datacache/blob/master/docs/api.md): all public functions, `Cache` methods, inspection results, and exceptions. - [Downloads and cache inspection](https://github.com/openvax/datacache/blob/master/docs/downloads.md): destinations, naming, @@ -202,7 +207,8 @@ FASTA files with `fetch_file`, then parse them in the consuming library. ## Guarantees and limits -Downloads are staged privately and published atomically after validation. +Downloads and complete archive trees are staged privately and published +atomically after validation. New files respect the process umask; replacements preserve existing access permissions. This includes pyensembl's private download helpers. diff --git a/datacache/__init__.py b/datacache/__init__.py index 300f5c0..679192a 100644 --- a/datacache/__init__.py +++ b/datacache/__init__.py @@ -33,6 +33,8 @@ from .cache import Cache from .resume import discard_partial from .bundles import BundleInspection, VersionedDatasetRegistry, inspect_bundle, install_bundle +from .archives import ( + ArchiveInspection, VersionedArchiveRegistry, inspect_archive, install_archive) from .file_registry import VersionedFileRegistry from .version import __version__ @@ -41,10 +43,14 @@ 'fetch_file', 'discard_partial', 'BundleInspection', + 'ArchiveInspection', + 'VersionedArchiveRegistry', 'VersionedDatasetRegistry', 'VersionedFileRegistry', 'inspect_bundle', 'install_bundle', + 'inspect_archive', + 'install_archive', 'expected_path', 'file_exists', 'resolve_path', diff --git a/datacache/archives.py b/datacache/archives.py new file mode 100644 index 0000000..8a18813 --- /dev/null +++ b/datacache/archives.py @@ -0,0 +1,920 @@ +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Safe, transactional installation of complete directory trees from tar archives.""" + +from dataclasses import dataclass, field +from datetime import datetime, timezone +import hashlib +import os +from pathlib import Path +import shutil +import stat +import tarfile +import tempfile +from uuid import uuid4 + +from filelock import FileLock + +from ._filesystem import open_regular, path_present, read_json, write_json +from .bundles import _relative_name +from .inspection import FileInspection +from .integrity import FileValidationError, _validate_expectations +from .provenance import redact_url + + +STORE = ".datacache-archive-store.json" +MANIFEST = ".datacache-archive-manifest.json" +CURRENT = "current.json" +FORMAT = 1 +CHUNK_SIZE = 2 ** 20 + + +@dataclass(frozen=True) +class ArchiveInspection: + """Read-only status for one installed archive-tree generation. + + ``verified`` means caller-supplied trusted archive or part hashes identified + the installed bytes. Receipt-only inspection still hashes every installed + file, but an observed receipt is not an independent source of trust. + """ + + path: str + status: str + verified: bool = False + generation: object = None + files: dict = field(default_factory=dict) + error: object = None + source_urls: tuple = () + fetched_at: object = None + archive_size: object = None + recorded_sha256: object = None + + +def _fingerprint(value): + return hashlib.sha256(value.encode("utf-8")).hexdigest() + + +def _local_identity(path): + return Path(path).absolute().as_uri() + + +def _normalize_sources(sources): + """Normalize one source or an explicitly ordered sequence of archive parts.""" + if isinstance(sources, (str, os.PathLike, dict)): + sources = [sources] + elif not isinstance(sources, (list, tuple)): + raise ValueError("sources must be a URL, path, mapping, or ordered sequence") + if not sources: + raise ValueError("sources must contain at least one archive source") + + result = [] + for source in sources: + if isinstance(source, os.PathLike): + source = {"path": source} + elif isinstance(source, str): + source = {"url": source} + elif isinstance(source, dict): + source = dict(source) + else: + raise ValueError("each archive source must be a URL, path, or mapping") + unknown = set(source) - {"url", "path", "sha256", "size"} + if unknown: + raise ValueError("unknown archive source options: %s" % sorted(unknown)) + url, path = source.get("url"), source.get("path") + if url is not None and (not isinstance(url, str) or not url): + raise ValueError("archive source url must be a nonempty string") + if path is not None: + try: + path = Path(path) + except TypeError as error: + raise ValueError("archive source path must be path-like") from error + if url is None and path is None: + raise ValueError("each archive source requires url or path") + identity = url if url is not None else _local_identity(path) + digest, size = source.get("sha256"), source.get("size") + _validate_expectations(digest, size) + result.append({ + "url": url, + "path": path, + "identity": identity, + "fingerprint": _fingerprint(identity), + "sha256": digest.lower() if digest else None, + "size": size, + }) + return tuple(result) + + +def _normalize_extra_files(extra_files): + if extra_files is None: + return {} + if not isinstance(extra_files, dict): + raise ValueError("extra_files must be a mapping of relative paths to text or bytes") + result = {} + nodes = {} + for name, value in extra_files.items(): + _relative_name(name) + if isinstance(value, str): + value = value.encode("utf-8") + elif not isinstance(value, bytes): + raise ValueError("extra file contents must be text or bytes") + parts = name.split("/") + for index in range(1, len(parts)): + parent = "/".join(parts[:index]) + key = parent.casefold() + previous = nodes.get(key) + if previous is None: + nodes[key] = (parent, "directory") + elif previous != (parent, "directory"): + raise ValueError("extra file paths collide as files and directories or ignoring case") + key = name.casefold() + if key in nodes: + raise ValueError("extra file paths collide as files and directories or ignoring case") + nodes[key] = (name, "file") + result[name] = value + return result + + +def _definition(sources, expected_sha256, expected_size, extra_files, require_verified=False): + normalized = _normalize_sources(sources) + _validate_expectations(expected_sha256, expected_size) + expected_sha256 = expected_sha256.lower() if expected_sha256 else None + extras = _normalize_extra_files(extra_files) + trusted = bool(expected_sha256 and expected_size is not None) or all( + source["sha256"] and source["size"] is not None for source in normalized) + if require_verified and not trusted: + raise ValueError( + "verified archive installs require expected_sha256 and expected_size " + "for the assembled archive, or sha256 and size for every part") + return { + "sources": normalized, + "expected_sha256": expected_sha256, + "expected_size": expected_size, + "extra_files": extras, + "trusted": trusted, + } + + +def _directory(path): + if not stat.S_ISDIR(Path(path).lstat().st_mode): + raise FileValidationError(path, "expected a directory, not a link") + + +def _file_mode(path): + from .download import _normal_creation_mode + return _normal_creation_mode(path) + + +def _store(path): + _directory(path) + if read_json(path / STORE) != {"format": FORMAT, "kind": "archive-tree"}: + raise FileValidationError(path, "unrecognized archive store") + _directory(path / "generations") + + +def _generation(store, generation): + if (not isinstance(generation, str) or len(generation) != 32 or + any(char not in "0123456789abcdef" for char in generation)): + raise FileValidationError(store, "invalid generation pointer") + result = store / "generations" / generation + _directory(result) + return result + + +def _manifest(receipt): + if (not isinstance(receipt, dict) or receipt.get("format") != FORMAT or + receipt.get("kind") != "archive-tree"): + raise ValueError("unrecognized archive manifest") + archive = receipt.get("archive") + sources = receipt.get("sources") + files = receipt.get("files") + directories = receipt.get("directories") + fetched_at = receipt.get("fetched_at") + if (not isinstance(archive, dict) or not isinstance(sources, list) or + not sources or not isinstance(files, dict) or not files or + not isinstance(directories, list) or not isinstance(fetched_at, str)): + raise ValueError("invalid archive manifest") + _validate_expectations(archive.get("sha256"), archive.get("size")) + if archive.get("sha256") is None or archive.get("size") is None: + raise ValueError("archive manifest does not identify assembled bytes") + normalized_sources = [] + for source in sources: + if not isinstance(source, dict): + raise ValueError("invalid archive source receipt") + if set(source) != {"url", "fingerprint", "sha256", "size"}: + raise ValueError("invalid archive source receipt") + if (not isinstance(source["url"], str) or + not isinstance(source["fingerprint"], str) or + len(source["fingerprint"]) != 64 or + any(char not in "0123456789abcdef" for char in source["fingerprint"])): + raise ValueError("invalid archive source receipt") + _validate_expectations(source["sha256"], source["size"]) + if source["sha256"] is None or source["size"] is None: + raise ValueError("archive source receipt does not identify observed bytes") + normalized_sources.append(source) + normalized_files = {} + for name, spec in files.items(): + _relative_name(name) + if not isinstance(spec, dict) or set(spec) != {"sha256", "size"}: + raise ValueError("invalid extracted file receipt") + _validate_expectations(spec["sha256"], spec["size"]) + if spec["sha256"] is None or spec["size"] is None: + raise ValueError("extracted file receipt does not identify observed bytes") + normalized_files[name] = spec + normalized_directories = [] + folded = set() + for name in directories: + _relative_name(name) + key = name.casefold() + if key in folded: + raise ValueError("duplicate directory in archive manifest") + folded.add(key) + normalized_directories.append(name) + file_keys = {name.casefold() for name in normalized_files} + if len(file_keys) != len(normalized_files) or file_keys & folded: + raise ValueError("colliding paths in archive manifest") + return archive, normalized_sources, normalized_files, normalized_directories, fetched_at + + +def _hash_handle(handle): + digest, size = hashlib.sha256(), 0 + for chunk in iter(lambda: handle.read(CHUNK_SIZE), b""): + digest.update(chunk) + size += len(chunk) + return digest.hexdigest(), size + + +def _open_local_source(path): + flags = os.O_RDONLY | getattr(os, "O_BINARY", 0) | getattr(os, "O_NOFOLLOW", 0) + descriptor = os.open(path, flags) + try: + info = os.fstat(descriptor) + if not stat.S_ISREG(info.st_mode): + raise FileValidationError(path, "expected a regular archive part") + return os.fdopen(descriptor, "rb") + except BaseException: + os.close(descriptor) + raise + + +def _walk_tree(root, receipt=None): + """Return observed file metadata, directories, and FileInspection values.""" + observed, directories, inspections = {}, [], {} + stack = [(Path(root), "")] + while stack: + directory, prefix = stack.pop() + with os.scandir(directory) as scanned: + entries = sorted( + scanned, key=lambda entry: entry.name.casefold(), reverse=True) + for entry in entries: + relative = entry.name if not prefix else prefix + "/" + entry.name + if not prefix and relative == MANIFEST: + continue + _relative_name(relative) + info = entry.stat(follow_symlinks=False) + if stat.S_ISDIR(info.st_mode): + directories.append(relative) + stack.append((Path(entry.path), relative)) + continue + if not stat.S_ISREG(info.st_mode): + raise FileValidationError(entry.path, "archive trees may contain only regular files and directories") + descriptor = open_regular(entry.path) + with os.fdopen(descriptor, "rb") as handle: + current = os.fstat(handle.fileno()) + digest, size = _hash_handle(handle) + observed[relative] = {"sha256": digest, "size": size} + expected = receipt.get(relative) if receipt is not None else None + if expected is not None and expected != observed[relative]: + raise FileValidationError(entry.path, "installed file disagrees with archive manifest") + inspections[relative] = FileInspection( + entry.path, "available", verified=expected is not None, + size=current.st_size, mtime=current.st_mtime) + directories.sort() + if receipt is not None and set(observed) != set(receipt): + missing = sorted(set(receipt) - set(observed)) + extra = sorted(set(observed) - set(receipt)) + raise FileValidationError(root, "archive tree file set differs from manifest; missing=%r extra=%r" % ( + missing, extra)) + return observed, directories, inspections + + +def _matches_definition(receipt_archive, receipt_sources, receipt_files, definition): + expected_sha256 = definition["expected_sha256"] + expected_size = definition["expected_size"] + if expected_sha256 is not None and receipt_archive["sha256"] != expected_sha256: + raise FileValidationError("archive", "assembled archive hash disagrees with requested source") + if expected_size is not None and receipt_archive["size"] != expected_size: + raise FileValidationError("archive", "assembled archive size disagrees with requested source") + if len(receipt_sources) != len(definition["sources"]): + raise FileValidationError("archive", "archive part count disagrees with requested source") + for recorded, expected in zip(receipt_sources, definition["sources"]): + if expected["sha256"] is not None and recorded["sha256"] != expected["sha256"]: + raise FileValidationError("archive", "archive part hash disagrees with requested source") + if expected["size"] is not None and recorded["size"] != expected["size"]: + raise FileValidationError("archive", "archive part size disagrees with requested source") + if not definition["trusted"]: + fingerprints = [source["fingerprint"] for source in receipt_sources] + expected = [source["fingerprint"] for source in definition["sources"]] + if fingerprints != expected: + raise FileValidationError("archive", "ordered archive source identities disagree") + for name, value in definition["extra_files"].items(): + expected = {"sha256": hashlib.sha256(value).hexdigest(), "size": len(value)} + if receipt_files.get(name) != expected: + raise FileValidationError(name, "extra file disagrees with requested content") + + +def _inspect_generation(store, generation, definition, verify_files=True): + directory = _generation(store, generation) + receipt = read_json(directory / MANIFEST) + archive, sources, files, directories, fetched_at = _manifest(receipt) + if definition is not None: + _matches_definition(archive, sources, files, definition) + metadata = { + "source_urls": tuple(source["url"] for source in sources), + "fetched_at": fetched_at, + "archive_size": archive["size"], + "recorded_sha256": archive["sha256"], + } + if not verify_files: + return ArchiveInspection( + str(store), "available", False, str(directory), **metadata) + _, observed_directories, inspections = _walk_tree(directory, files) + if observed_directories != sorted(directories): + raise FileValidationError(directory, "archive tree directories disagree with manifest") + return ArchiveInspection( + str(store), "available", bool(definition and definition["trusted"]), + str(directory), inspections, **metadata) + + +def inspect_archive( + destination, sources=None, *, expected_sha256=None, expected_size=None, + extra_files=None, verify_files=True): + """Inspect one installed archive tree without writes, locks, or network. + + Omit ``sources`` for receipt-only consistency checking. Supplying sources + checks the requested ordered archive identity or trusted content hashes. + ``verify_files=False`` validates publication and source metadata without + hashing the extracted tree; its result is never marked verified. + """ + if not isinstance(verify_files, bool): + raise ValueError("verify_files must be a boolean") + definition = None + if sources is not None: + definition = _definition( + sources, expected_sha256, expected_size, extra_files, + require_verified=False) + elif any(value is not None for value in (expected_sha256, expected_size, extra_files)): + raise ValueError("sources are required with archive expectations or extra_files") + path = Path(destination) + try: + if not path_present(path): + return ArchiveInspection(str(path), "missing") + _store(path) + try: + pointer = read_json(path / CURRENT) + except FileNotFoundError: + if any((path / "generations").iterdir()): + return ArchiveInspection(str(path), "recovery-required") + return ArchiveInspection(str(path), "missing") + return _inspect_generation( + path, pointer["generation"], definition, verify_files=verify_files) + except PermissionError as error: + return ArchiveInspection(str(path), "inaccessible", error=error) + except (OSError, ValueError, KeyError, TypeError, RecursionError) as error: + return ArchiveInspection(str(path), "invalid", error=error) + + +def _initialize(path): + if path_present(path): + _store(path) # Never claim a legacy, empty, or otherwise foreign path. + return + staging = Path(tempfile.mkdtemp(prefix=".datacache-archive-store-", dir=path.parent)) + try: + write_json( + staging / STORE, {"format": FORMAT, "kind": "archive-tree"}, + mode=_file_mode(staging)) + (staging / "generations").mkdir() + os.chmod(staging, stat.S_IMODE((staging / "generations").stat().st_mode)) + os.replace(staging, path) + finally: + if staging.exists(): + shutil.rmtree(staging) + + +def _recover(path, definition): + for entry in sorted((path / "generations").iterdir(), reverse=True): + try: + candidate = _inspect_generation(path, entry.name, definition) + except (OSError, ValueError, KeyError, TypeError, RecursionError): + continue + write_json(path / CURRENT, {"generation": entry.name}, mode=_file_mode(path)) + return candidate + return None + + +def _archive_members( + archive, max_members=None, max_extracted_size=None, extra_names=()): + members = archive.getmembers() + if not members: + raise FileValidationError(archive.name, "archive is empty") + if max_members is not None and len(members) > max_members: + raise FileValidationError(archive.name, "archive has more than %d members" % max_members) + nodes = {} + planned = [] + total = 0 + regular_files = 0 + for member in members: + if member.isdir(): + name, kind = member.name.rstrip("/"), "directory" + elif member.isfile(): + name, kind = member.name, "file" + if member.size < 0: + raise FileValidationError(archive.name, "archive member has negative size") + total += member.size + regular_files += 1 + else: + raise FileValidationError( + archive.name, "archive contains a link or special member: %s" % member.name) + try: + _relative_name(name) + except ValueError as error: + raise FileValidationError( + archive.name, "unsafe archive member path: %s" % member.name) from error + parts = name.split("/") + for index in range(1, len(parts)): + parent = "/".join(parts[:index]) + key = parent.casefold() + previous = nodes.get(key) + if previous is None: + nodes[key] = (parent, "directory", False) + elif previous[0] != parent or previous[1] != "directory": + raise FileValidationError(archive.name, "archive member paths collide: %s" % name) + key = name.casefold() + previous = nodes.get(key) + if kind == "directory" and previous is not None and previous[1] == "directory" and not previous[2]: + if previous[0] != name: + raise FileValidationError(archive.name, "archive member paths collide: %s" % name) + nodes[key] = (name, kind, True) + elif previous is not None: + raise FileValidationError(archive.name, "archive member paths collide: %s" % name) + else: + nodes[key] = (name, kind, True) + planned.append((member, name, kind)) + for name in extra_names: + parts = name.split("/") + for index in range(1, len(parts)): + parent = "/".join(parts[:index]) + key = parent.casefold() + previous = nodes.get(key) + if previous is None: + nodes[key] = (parent, "directory", False) + elif previous[0] != parent or previous[1] != "directory": + raise FileValidationError( + archive.name, "extra file path collides with archive content: %s" % name) + if name.casefold() in nodes: + raise FileValidationError( + archive.name, "extra file path collides with archive content: %s" % name) + nodes[name.casefold()] = (name, "file", True) + if not regular_files: + raise FileValidationError(archive.name, "archive contains no regular files") + if max_extracted_size is not None and total > max_extracted_size: + raise FileValidationError( + archive.name, "archive expands to more than %d bytes" % max_extracted_size) + return planned + + +def _extract_tar( + archive_path, destination, max_members=None, max_extracted_size=None, + extra_names=()): + try: + with tarfile.open(archive_path, "r:*") as archive: + planned = _archive_members( + archive, max_members, max_extracted_size, extra_names) + directories = sorted( + {name for _, name, kind in planned if kind == "directory"} | + {"/".join(name.split("/")[:index]) + for _, name, _ in planned for index in range(1, len(name.split("/")))}, + key=lambda name: (name.count("/"), name.casefold())) + for name in directories: + (destination / Path(*name.split("/"))).mkdir(exist_ok=True) + for member, name, kind in planned: + if kind == "directory": + continue + target = destination / Path(*name.split("/")) + source = archive.extractfile(member) + if source is None: + raise FileValidationError(archive_path, "could not read archive member %s" % name) + flags = os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_BINARY", 0) + flags |= getattr(os, "O_NOFOLLOW", 0) + descriptor = os.open(target, flags, 0o666) + count = 0 + try: + output = os.fdopen(descriptor, "wb") + except BaseException: + os.close(descriptor) + source.close() + raise + with source, output: + for chunk in iter(lambda: source.read(CHUNK_SIZE), b""): + output.write(chunk) + count += len(chunk) + if count != member.size: + raise FileValidationError( + archive_path, "archive member %s is truncated" % name) + except FileValidationError: + raise + except (tarfile.TarError, EOFError) as error: + raise FileValidationError(archive_path, "invalid tar archive: %s" % error) from error + + +def _validate_limit(name, value): + if value is not None and (isinstance(value, bool) or not isinstance(value, int) or value < 0): + raise ValueError("%s must be a non-negative integer" % name) + + +def _working_key(definition): + serializable = { + "sources": [{key: source[key] for key in ("identity", "sha256", "size")} + for source in definition["sources"]], + "expected_sha256": definition["expected_sha256"], + "expected_size": definition["expected_size"], + "extra_files": {name: hashlib.sha256(value).hexdigest() + for name, value in definition["extra_files"].items()}, + } + import json + return hashlib.sha256(json.dumps(serializable, sort_keys=True).encode("utf-8")).hexdigest() + + +def _part_expectations(definition, index): + source = definition["sources"][index] + digest, size = source["sha256"], source["size"] + if len(definition["sources"]) == 1: + digest = digest or definition["expected_sha256"] + size = size if size is not None else definition["expected_size"] + return digest, size + + +def _assemble_archive(working, definition, options): + from .download import fetch_file + + archive_path = working / "archive.tar" + aggregate = hashlib.sha256() + aggregate_size = 0 + recorded = [] + keep_parts = options.get("resume", False) + with archive_path.open("wb") as assembled: + for index, source in enumerate(definition["sources"]): + expected_sha256, expected_size = _part_expectations(definition, index) + temporary = None + if source["path"] is not None: + handle = _open_local_source(source["path"]) + else: + temporary = working / ("part-%06d" % index) + fetch_file( + source["url"], destination=temporary, raw=True, + expected_sha256=expected_sha256, expected_size=expected_size, + **options) + handle = _open_local_source(temporary) + digest, size = hashlib.sha256(), 0 + with handle: + for chunk in iter(lambda: handle.read(CHUNK_SIZE), b""): + assembled.write(chunk) + aggregate.update(chunk) + digest.update(chunk) + size += len(chunk) + aggregate_size += len(chunk) + actual = digest.hexdigest() + if expected_size is not None and size != expected_size: + raise FileValidationError(source["path"] or source["url"], "archive part size mismatch") + if expected_sha256 is not None and actual != expected_sha256: + raise FileValidationError(source["path"] or source["url"], "archive part SHA-256 mismatch") + recorded.append({ + "url": redact_url(source["identity"]), + "fingerprint": source["fingerprint"], + "sha256": actual, + "size": size, + }) + if temporary is not None and not keep_parts: + temporary.unlink() + actual_sha256 = aggregate.hexdigest() + if (definition["expected_size"] is not None and + aggregate_size != definition["expected_size"]): + raise FileValidationError(archive_path, "assembled archive size mismatch") + if (definition["expected_sha256"] is not None and + actual_sha256 != definition["expected_sha256"]): + raise FileValidationError(archive_path, "assembled archive SHA-256 mismatch") + return archive_path, {"sha256": actual_sha256, "size": aggregate_size}, recorded + + +def install_archive( + destination, sources, *, expected_sha256=None, expected_size=None, + extra_files=None, force=False, verified=True, download_options=None, + max_members=None, max_extracted_size=None): + """Safely install a complete tar archive tree as one immutable generation. + + ``sources`` is one URL/path or an ordered sequence. Source mappings accept + ``url``, ``path``, ``sha256``, and ``size``; supplying both URL and path uses + local bytes while retaining the URL as source identity. Split parts are + concatenated byte-for-byte in the supplied order. ``extra_files`` are + written after extraction and before the tree receipt, which lets consumers + publish their own source record in the same transaction. + """ + if not isinstance(force, bool) or not isinstance(verified, bool): + raise ValueError("force and verified must be booleans") + _validate_limit("max_members", max_members) + _validate_limit("max_extracted_size", max_extracted_size) + definition = _definition( + sources, expected_sha256, expected_size, extra_files, + require_verified=verified) + options = dict(download_options or {}) + allowed = {"timeout", "chunk_size", "progress_callback", "show_progress", + "max_retries", "retry_backoff", "retry_max_delay", "resume"} + if set(options) - allowed: + raise ValueError("unsupported archive download options: %s" % sorted(set(options) - allowed)) + if options.get("resume"): + for index, source in enumerate(definition["sources"]): + digest, size = _part_expectations(definition, index) + if source["path"] is None and (digest is None or size is None): + raise ValueError("resumable archive parts require sha256 and size") + + path = Path(destination) + inspection = inspect_archive( + path, sources, expected_sha256=expected_sha256, + expected_size=expected_size, extra_files=extra_files) + if not force and inspection.status == "available": + return Path(inspection.generation) + if inspection.status == "inaccessible": + raise inspection.error + if not force and inspection.status == "invalid": + raise FileValidationError(path, "invalid archive installation; use force=True to explicitly repair") from inspection.error + + path.parent.mkdir(parents=True, exist_ok=True) + key = hashlib.sha256(os.fsencode(path.name.casefold())).hexdigest()[:32] + with FileLock(str(path.parent / (".datacache-archive-lock-" + key))): + _initialize(path) + inspection = inspect_archive( + path, sources, expected_sha256=expected_sha256, + expected_size=expected_size, extra_files=extra_files) + if not force and inspection.status == "available": + return Path(inspection.generation) + if inspection.status == "invalid" and not force: + raise FileValidationError(path, "invalid archive installation; use force=True to explicitly repair") from inspection.error + if inspection.status in ("missing", "recovery-required", "invalid"): + recovered = _recover(path, definition) + if recovered is not None: + return Path(recovered.generation) + + resumable = options.get("resume", False) + working = path / (".staging-" + (_working_key(definition) if resumable else uuid4().hex)) + working.mkdir(mode=0o700, exist_ok=resumable) + if os.name == "posix": + info = working.lstat() + if (not stat.S_ISDIR(info.st_mode) or info.st_uid != os.getuid() or + stat.S_IMODE(info.st_mode) & 0o077): + raise FileValidationError(working, "archive staging must be a private directory") + tree = working / "tree" + if tree.exists(): + shutil.rmtree(tree) + tree.mkdir() + published = False + try: + archive_path, archive_record, source_records = _assemble_archive( + working, definition, options) + _extract_tar( + archive_path, tree, max_members, max_extracted_size, + definition["extra_files"]) + for name, value in definition["extra_files"].items(): + target = tree / Path(*name.split("/")) + target.parent.mkdir(parents=True, exist_ok=True) + flags = os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_BINARY", 0) + flags |= getattr(os, "O_NOFOLLOW", 0) + try: + descriptor = os.open(target, flags, 0o666) + except FileExistsError as error: + raise FileValidationError(target, "extra file collides with archive content") from error + with os.fdopen(descriptor, "wb") as output: + output.write(value) + files, directories, _ = _walk_tree(tree) + receipt = { + "format": FORMAT, + "kind": "archive-tree", + "fetched_at": datetime.now(timezone.utc).isoformat(), + "archive": archive_record, + "sources": source_records, + "files": files, + "directories": directories, + } + write_json(tree / MANIFEST, receipt, mode=_file_mode(tree)) + generation = uuid4().hex + final = path / "generations" / generation + os.replace(tree, final) + published = True + candidate = _inspect_generation(path, generation, definition) + write_json(path / CURRENT, {"generation": generation}, mode=_file_mode(path)) + return Path(candidate.generation) + finally: + if published or not resumable: + shutil.rmtree(working) + + +class VersionedArchiveRegistry: + """Versioned catalogue of complete archive-tree downloads. + + Each item has ``default_version`` and a ``versions`` mapping. A version is + one source, an ordered source sequence, or a mapping containing ``sources`` + plus archive expectations, consumer ``extra_files``, and extraction limits. + ``cache_root`` uses ``//`` stores; ``store_path`` can + instead preserve a consumer's established ``(name, version)`` layout. + Construction validates definitions without touching the filesystem. + """ + + def __init__( + self, archives, *, cache_root=None, store_path=None, verified=True): + if (cache_root is None) == (store_path is None): + raise ValueError("provide exactly one of cache_root or store_path") + if store_path is not None and not callable(store_path): + raise ValueError("store_path must be callable") + if not isinstance(verified, bool): + raise ValueError("verified must be a boolean") + if not isinstance(archives, dict) or not archives: + raise ValueError("archives must be a nonempty mapping") + self.verified = verified + self._path = (store_path if store_path is not None else + lambda name, version: Path(cache_root) / name / version) + self._archives = {} + version_options = { + "sources", "expected_sha256", "expected_size", "extra_files", + "max_members", "max_extracted_size", + } + for name, archive in archives.items(): + from .bundles import _component + _component(name) + if not isinstance(archive, dict): + raise ValueError("archive definitions must be mappings") + unknown = set(archive) - {"default_version", "versions", "description"} + if unknown: + raise ValueError("unknown archive definition options: %s" % sorted(unknown)) + versions = archive.get("versions") + if not isinstance(versions, dict) or not versions: + raise ValueError("each archive requires a nonempty versions mapping") + normalized_versions = {} + for version, value in versions.items(): + _component(version) + if isinstance(value, dict) and "sources" in value: + options = dict(value) + unknown = set(options) - version_options + if unknown: + raise ValueError("unknown archive version options: %s" % sorted(unknown)) + else: + options = {"sources": value} + _validate_limit("max_members", options.get("max_members")) + _validate_limit("max_extracted_size", options.get("max_extracted_size")) + definition = _definition( + options["sources"], options.get("expected_sha256"), + options.get("expected_size"), options.get("extra_files"), + require_verified=verified) + normalized_versions[version] = { + "sources": tuple({key: source[key] for key in + ("url", "path", "sha256", "size") + if source[key] is not None} + for source in definition["sources"]), + "expected_sha256": definition["expected_sha256"], + "expected_size": definition["expected_size"], + "extra_files": definition["extra_files"], + "max_members": options.get("max_members"), + "max_extracted_size": options.get("max_extracted_size"), + } + default = archive.get("default_version") + if default not in normalized_versions: + raise ValueError("default_version must name a concrete archive version") + self._archives[name] = { + "default_version": default, + "versions": normalized_versions, + "description": archive.get("description", ""), + } + + def resolve_version(self, name, version=None): + if name not in self._archives: + raise ValueError( + "unknown archive %r; known: %s" % ( + name, ", ".join(sorted(self._archives)))) + definition = self._archives[name] + version = definition["default_version"] if version is None else version + if version not in definition["versions"]: + raise ValueError("unknown version %r for %s" % (version, name)) + return version + + def store_path(self, name, version=None): + """Return the managed store path without creating or inspecting it.""" + version = self.resolve_version(name, version) + return Path(self._path(name, version)) + + def _version(self, name, version): + version = self.resolve_version(name, version) + return version, self._archives[name]["versions"][version] + + @staticmethod + def _local_sources(sources, source_paths): + if source_paths is None: + return sources + if isinstance(source_paths, (str, os.PathLike)): + source_paths = [source_paths] + elif not isinstance(source_paths, (list, tuple)): + raise ValueError("source_paths must be a path or ordered sequence") + if len(source_paths) != len(sources): + raise ValueError("source_paths must have one path for every archive part") + result = [] + for source, path in zip(sources, source_paths): + item = dict(source) + if "url" not in item: + item["url"] = _local_identity(item["path"]) + item["path"] = Path(path) + result.append(item) + return tuple(result) + + def inspect(self, name, version=None, *, verify_files=True): + """Inspect one version offline; no destination or lock is created.""" + version, definition = self._version(name, version) + return inspect_archive( + self.store_path(name, version), definition["sources"], + expected_sha256=definition["expected_sha256"], + expected_size=definition["expected_size"], + extra_files=definition["extra_files"], verify_files=verify_files) + + def download( + self, name, version=None, *, force=False, source_paths=None, + **download_options): + """Install/reuse one version and return its immutable extracted root.""" + version, definition = self._version(name, version) + sources = self._local_sources(definition["sources"], source_paths) + return install_archive( + self.store_path(name, version), sources, + expected_sha256=definition["expected_sha256"], + expected_size=definition["expected_size"], + extra_files=definition["extra_files"], force=force, + verified=self.verified, download_options=download_options, + max_members=definition["max_members"], + max_extracted_size=definition["max_extracted_size"]) + + def local_path(self, name, version=None, *, verify_files=False): + """Return an installed generation without downloading or repairing.""" + inspected = self.inspect(name, version, verify_files=verify_files) + if inspected.status == "missing": + raise FileNotFoundError(inspected.path) + if inspected.status == "inaccessible": + raise inspected.error + if inspected.status != "available": + raise FileValidationError( + inspected.path, "archive installation is %s" % inspected.status) + return Path(inspected.generation) + + def ensure(self, name, version=None, **download_options): + """Install if needed, then return the extracted generation path.""" + return self.download(name, version, **download_options) + + def is_cached(self, name, version=None, *, verify_files=False): + """Whether one version has an available published generation.""" + return self.inspect( + name, version, verify_files=verify_files).status == "available" + + def status(self, name=None): + """Return one offline status row per concrete archive version.""" + if name is not None: + self.resolve_version(name) + names = [name] + else: + names = sorted(self._archives) + rows = [] + for archive_name in names: + archive = self._archives[archive_name] + for version, definition in archive["versions"].items(): + inspection = self.inspect( + archive_name, version, verify_files=False) + identities = [source.get("url") or _local_identity(source["path"]) + for source in definition["sources"]] + rows.append({ + "name": archive_name, + "version": version, + "default": version == archive["default_version"], + "description": archive["description"], + "sources": [redact_url(identity) for identity in identities], + "downloaded_sources": list(inspection.source_urls), + "path": str(self.store_path(archive_name, version)), + "generation": inspection.generation, + "status": inspection.status, + "cached": inspection.status == "available", + "fetched_at": inspection.fetched_at, + "archive_size": inspection.archive_size, + "recorded_sha256": inspection.recorded_sha256, + "inspection": inspection, + }) + return rows diff --git a/datacache/file_registry.py b/datacache/file_registry.py index 1ef5438..225cd9d 100644 --- a/datacache/file_registry.py +++ b/datacache/file_registry.py @@ -161,7 +161,9 @@ def status(self) -> list[dict]: "available_versions": sorted(spec["urls"]), "cached": path.exists(), "cached_version": record.get("version") if record else None, + "url": record.get("url") if path.exists() else None, "bytes": record.get("bytes") if path.exists() else None, + "sha256": record.get("sha256") if path.exists() else None, "downloaded_at": record.get("downloaded_at") if path.exists() else None, "path": str(path), } diff --git a/datacache/version.py b/datacache/version.py index 6b0872c..638c121 100644 --- a/datacache/version.py +++ b/datacache/version.py @@ -1 +1 @@ -__version__ = "1.15.0" +__version__ = "1.16.0" diff --git a/docs/api.md b/docs/api.md index b336261..33dae3a 100644 --- a/docs/api.md +++ b/docs/api.md @@ -1,12 +1,12 @@ # Public API reference This reference covers every name exported in `datacache.__all__` and every -public `Cache` method in DataCache 1.15.0. Import these names from `datacache`. +public `Cache` method in DataCache 1.16.0. Import these names from `datacache`. Signatures below show all defaults; arguments after `*` are keyword-only. Method signatures omit `self` and are called on a `Cache` instance. Start with the [quickstart](../README.md#quickstart) for a first download. -The [download guide](downloads.md), [data guide](data.md), +The [download guide](downloads.md), [archive guide](archives.md), [data guide](data.md), [progress guide](progress.md), and [shared-cache guide](shared-caches.md) explain the longer workflows and compatibility guarantees. @@ -15,6 +15,8 @@ explain the longer workflows and compatibility guarantees. | Area | APIs | | --- | --- | | Downloading | [fetch_file](#fetch_file), [fetch_csv_dataframe](#fetch_csv_dataframe), [fetch_and_transform](#fetch_and_transform) | +| Complete archive trees | [install_archive](#install_archive), [inspect_archive](#inspect_archive), [ArchiveInspection](#archiveinspection), [VersionedArchiveRegistry](#versionedarchiveregistry) | +| Versioned file bundles | [install_bundle](#install_bundle), [inspect_bundle](#inspect_bundle), [BundleInspection](#bundleinspection), [VersionedDatasetRegistry](#versioneddatasetregistry), [VersionedFileRegistry](#versionedfileregistry) | | Paths and presence | [expected_path](#expected_path), [file_exists](#file_exists), [build_local_filename](#build_local_filename), [get_data_dir](#get_data_dir), [get_cache_root](#get_cache_root), [resolve_path](#resolve_path), [build_path](#build_path), [ensure_dir](#ensure_dir), [clear_cache](#clear_cache) | | Integrity and permissions | [validate_file](#validate_file), [inspect_file](#inspect_file), [inspect_files](#inspect_files), [make_file_readable](#make_file_readable) | | Results and exceptions | [FileInspection](#fileinspection), [CacheInspection](#cacheinspection), [FileValidationError](#filevalidationerror) | @@ -34,7 +36,9 @@ from contextlib import closing from pathlib import Path from tempfile import TemporaryDirectory import hashlib +import io import os +import tarfile import pandas as pd import datacache as dc @@ -48,6 +52,15 @@ source.write_bytes(contents) url = source.as_uri() sha256 = hashlib.sha256(contents).hexdigest() frame = pd.DataFrame({"id": [1, 2], "value": [10, 20]}) + +archive_source = root / "models.tar.bz2" +with tarfile.open(archive_source, "w:bz2") as archive: + member_contents = b'{"model": "example"}\n' + member = tarfile.TarInfo("models/model.json") + member.size = len(member_contents) + archive.addfile(member, io.BytesIO(member_contents)) +archive_contents = archive_source.read_bytes() +archive_sha256 = hashlib.sha256(archive_contents).hexdigest() ``` ## Downloading @@ -1064,13 +1077,6 @@ It is independent of the integer `version` used to identify a SQLite cache. print(dc.__version__) ``` -After finishing the examples, remove their temporary files: - -```python -temporary.cleanup() -``` - - ## discard_partial ```text @@ -1081,6 +1087,118 @@ Explicitly discard this user's private resumable bytes for an exact destination, under its writer lock. The installed file is unchanged; absent state is a no-op. See [resumable downloads](downloads.md#resumable-http-downloads). +## install_archive + +```text +install_archive( + destination, sources, *, expected_sha256=None, expected_size=None, + extra_files=None, force=False, verified=True, download_options=None, + max_members=None, max_extracted_size=None +) +``` + +Safely extract a complete tar archive into a private directory and publish the +tree as one immutable generation. `sources` is one URL, one `Path`, or an +explicitly ordered sequence. A mapping may contain `url`, `path`, `sha256`, and +`size`; supplying both `url` and `path` reads local bytes while retaining the URL +as the logical source identity. Multiple sources are concatenated byte-for-byte, +which supports historical split archives. + +By default the assembled archive needs trusted `expected_sha256` and +`expected_size`, or every part needs `sha256` and `size`. Set `verified=False` +explicitly for historical sources without published integrity metadata. Observed +hashes are still recorded and checked, but do not authenticate the source. + +`extra_files` maps safe relative names to text or bytes written after extraction +and before publication. It is intended for consumer receipts such as +`DOWNLOAD_INFO.csv`; collisions with archive content fail the installation. +`download_options` accepts timeout, chunk size, progress, retry, and resume +settings from `fetch_file`. Resumable URL parts each require a trusted hash and +size. Optional `max_members` and `max_extracted_size` limits are checked before +files are extracted. + +**Returns:** the immutable extracted-generation `Path`, not the managed store +path. Existing returned paths remain usable across forced refreshes. + +**Raises:** `ValueError` for invalid definitions or options; +`FileValidationError` for integrity failures, unsafe members, malformed tar +archives, or invalid installations; transport and filesystem errors otherwise. +An invalid current installation requires `force=True`. A legacy or otherwise +foreign destination is never claimed, even with force. + +```python +archive_store = root / "archive-store" +archive_generation = dc.install_archive( + archive_store, + archive_source, + expected_sha256=archive_sha256, + expected_size=len(archive_contents), + extra_files={"DOWNLOAD_INFO.csv": "url\n" + archive_source.as_uri() + "\n"}, +) +assert (archive_generation / "models/model.json").is_file() +``` + +See the [archive installation guide](archives.md) for split sources, safe-member +rules, publication semantics, and an MHCflurry adapter. + +## inspect_archive + +```text +inspect_archive( + destination, sources=None, *, expected_sha256=None, expected_size=None, + extra_files=None, verify_files=True +) +``` + +Inspect one installed archive generation without writes, locks, repair, or +network access. Omit `sources` for receipt-only consistency checking, which +returns `verified=False`. Supply the same source definition, expectations, and +consumer metadata used for installation to validate the requested identity. +Set `verify_files=False` for a fast published-generation/source-identity check +that does not hash the extracted tree. Fast results have an empty `files` +mapping and `verified=False`; use the default before asserting content integrity. + +## ArchiveInspection + +A frozen record with `path`, `status`, `verified`, `generation`, `files`, +`error`, `source_urls`, `fetched_at`, `archive_size`, and `recorded_sha256`. +Status is `available`, `missing`, `invalid`, `inaccessible`, or +`recovery-required`. `generation` is the extracted tree applications should use; +`files` maps every installed relative file name to a `FileInspection`, or is +empty for a fast `verify_files=False` inspection. Source URLs are redacted for +display. The recorded digest is observed receipt data, not trusted verification. + +## VersionedArchiveRegistry + +```text +VersionedArchiveRegistry( + archives, *, cache_root=None, store_path=None, verified=True +) +``` + +Reusable version catalogue for archive-tree consumers. Each named archive has a +`default_version`, optional description, and `versions` mapping. A concrete +version may be one source, an ordered source sequence, or a mapping containing +`sources`, `expected_sha256`, `expected_size`, `extra_files`, `max_members`, and +`max_extracted_size`. Definitions are validated without filesystem access. + +Select exactly one destination strategy. `cache_root` stores versions at +`//`. A two-argument `store_path(name, version)` callback +lets an adapter retain an established layout such as +`//` without putting layout policy in +DataCache. `verified=False` explicitly enables historical unpinned catalogues. + +| Method | Result | +| --- | --- | +| `resolve_version(name, version=None)` | Concrete version, applying the pinned default. | +| `store_path(name, version=None)` | Managed store `Path`, without creating or inspecting it. | +| `inspect(name, version=None, *, verify_files=True)` | Offline `ArchiveInspection`. | +| `download(name, version=None, *, force=False, source_paths=None, **download_options)` | Install/reuse and return the generation `Path`; `source_paths` overlays ordered local files while retaining catalogue identities. | +| `local_path(name, version=None, *, verify_files=False)` | Resolve a published generation without network, writes, or repair. | +| `ensure(name, version=None, **download_options)` | Install if needed and return the generation. | +| `is_cached(name, version=None, *, verify_files=False)` | Whether a published generation is available; full verification is optional. | +| `status(name=None)` | One read-only row per concrete version, including catalogue and downloaded source URLs, store/generation paths, default flag, status, fetch time, observed archive size/hash and inspection. | + ## install_bundle ```text @@ -1150,4 +1268,10 @@ trust boundaries, and the distinction from transactional generation bundles. | `is_cached(name, version=None)` | Presence only, not verified integrity. | | `download(name, version=None, *, force=False, **download_options)` | One fixed Path; fetch_file options control acquisition. Ordinary cache hits do not hash or write; explicit size/hash expectations are checked. | | `ensure(name, version=None, **download_options)` | Same download/reuse behavior and Path result. | -| `status()` | Legacy status dicts with name, description, default_version, available_versions, cached, cached_version, bytes, downloaded_at and path. | +| `status()` | Legacy status dicts with name, description, default_version, available_versions, cached, cached_version, URL, bytes, observed SHA-256, downloaded_at and path. | + +After finishing the examples, remove their temporary files: + +```python +temporary.cleanup() +``` diff --git a/docs/archives.md b/docs/archives.md new file mode 100644 index 0000000..aed93c8 --- /dev/null +++ b/docs/archives.md @@ -0,0 +1,197 @@ +# Archive directory installation + +`install_archive` safely publishes an entire tar archive tree. It is intended +for applications such as MHCflurry whose released weights and datasets already +have nested paths inside `.tar.bz2` files. DataCache owns transfer, extraction, +integrity receipts, concurrency and publication; the caller continues to own +catalogue versions and domain-specific validation. + +For a reusable catalogue, wrap it in `VersionedArchiveRegistry`; the lower-level +functions remain useful when a downstream library already owns version lookup. + +## One archive + +```python +from datacache import install_archive, inspect_archive + +tree = install_archive( + cache_root / "models_class1_presentation", + "https://example.org/models.tar.bz2", + expected_sha256=trusted_archive_sha256, + expected_size=trusted_archive_size, + download_options={"timeout": 60, "max_retries": 3, "show_progress": True}, +) +models = tree / "models" +state = inspect_archive( + cache_root / "models_class1_presentation", + "https://example.org/models.tar.bz2", + expected_sha256=trusted_archive_sha256, + expected_size=trusted_archive_size, +) +assert state.status == "available" and state.verified +``` + +The destination is a managed store, not the extracted directory. Installation +returns the immutable generation containing the archive's paths. Consumers must +resolve the generation once and use that returned path for an operation; they +must not construct `/` directly. + +## Ordered split archives and local inputs + +An ordered list is concatenated byte-for-byte before tar parsing: + +```python +urls = [ + "https://example.org/data.tar.bz2.part.aa", + "https://example.org/data.tar.bz2.part.ab", +] +tree = install_archive( + cache_root / "data_evaluation", + urls, + expected_sha256=trusted_concatenated_sha256, + expected_size=trusted_concatenated_size, +) +``` + +Each item can instead be `{url, sha256, size}` to verify parts independently. +If a command already has downloaded files, supply `path` as well as the logical +catalogue URL. The local bytes are not fetched again, while reuse still compares +the catalogue identities: + +```python +sources = [ + {"url": url, "path": already_downloaded / url.rsplit("/", 1)[-1]} + for url in urls +] +``` + +Strings are URLs; use a `Path` object for a local-only source. The full logical +URL is hashed for identity, while receipts omit credentials, query strings and +fragments. Avoid secrets in URL paths because paths are retained for display. + +## MHCflurry compatibility + +Historical MHCflurry catalogues lack published hashes and record the exact, +ordered URL list in `DOWNLOAD_INFO.csv`. An adapter can preserve that contract: + +```python +import csv +import io + +buffer = io.StringIO(newline="") +writer = csv.writer(buffer) +writer.writerow(["url"]) +writer.writerows([url] for url in urls) + +tree = install_archive( + release_directory / download_name, + sources, + verified=False, + extra_files={"DOWNLOAD_INFO.csv": buffer.getvalue()}, + download_options={"timeout": timeout, "max_retries": max_retries}, +) +``` + +MHCflurry should treat a fast +`inspect_archive(..., verify_files=False).status == "available"` using the same +sources and `extra_files` as installed, not mere destination existence. Its +`get_path` adapter should append member paths to `state.generation`. The fast +check validates the atomic publication receipt and requested source identity +without hashing large model files on every path lookup. Explicit diagnostics can +use the default `verify_files=True` for complete tree verification. This prevents +an initialized store or interrupted install from being mistaken for a complete +bundle. + +For a catalogue-wide adapter, `VersionedArchiveRegistry` supplies consistent +version resolution, download, path, inspection and status methods while allowing +MHCflurry's existing directory order: + +```python +from datacache import VersionedArchiveRegistry + +registry = VersionedArchiveRegistry( + archive_definitions, + store_path=lambda name, release: release_directory(release) / name, + verified=False, # Historical catalogue entries have no published hashes. +) +tree = registry.download( + "models_class1_presentation", "2.3.0", + show_progress=True, timeout=60, max_retries=3, +) +rows = registry.status("models_class1_presentation") +``` + +`source_paths=` on `download` accepts one local file per ordered catalogue URL, +covering MHCflurry's `--already-downloaded-dir` mode without changing source +identity or bypassing transactional extraction. + +Existing MHCflurry directories are deliberately not adopted or overwritten, +whether empty, complete, partial, or unreceipted. A migration can continue to +read a legacy directory and use archive stores only for new installations, or +explicitly validate/import legacy content in application code. DataCache never +claims it silently, including with `force=True`. + +## Extraction policy + +DataCache parses tar, gzip-compressed tar, bzip2-compressed tar and xz-compressed +tar archives using format detection. Before writing members it rejects: + +- absolute, parent-traversing, non-portable, and reserved DataCache paths; +- symbolic links, hard links, devices, FIFOs and every other special member; +- duplicate names, file/directory conflicts, and collisions ignoring case; +- archives without a regular file; and +- consumer `extra_files` that collide with archive content. + +Extraction writes regular files itself rather than calling `extractall`. Archive +ownership, timestamps and permission bits—including executable and set-ID bits— +are not applied. New files and directories use ordinary umask-derived modes. +`max_members` and `max_extracted_size` can impose application-specific resource +limits based on tar headers before extraction begins. + +After extraction, the generation receipt records the assembled archive's +observed hash and size, every part's observed hash, ordered source fingerprints, +and a sorted hash/size manifest of every installed file. Inspection rejects +missing, changed, additional, linked or special files and directory changes. +`ArchiveInspection` also exposes the redacted source URLs, fetch time, assembled +size and recorded digest for downstream `list` and `info` commands. + +## Trust and reuse + +Trusted verification requires either an assembled archive hash and size or a +hash and size for every part. Hash-pinned content can be reused across different +mirrors. `verified=False` permits historical unpinned sources; it verifies local +consistency against observed receipts and requires the same ordered full source +identities, but `ArchiveInspection.verified` remains false. + +A valid cache hit is entirely read-only. `inspect_archive` hashes the complete +tree by default; `verify_files=False` provides a receipt/source-identity-only +path lookup and never reports `verified=True`. Invalid installations require an +explicit `force=True`; inaccessible installations propagate their permission +error. Consumer metadata is part of the requested generation identity, so a +different `DOWNLOAD_INFO.csv` requires an explicit refresh. + +## Publication, recovery and platforms + +Writers serialize on a permanent sibling lock. Downloads, concatenation, +extraction, consumer metadata and the DataCache receipt complete in private +staging. The finished tree is renamed to `generations/` and `current.json` +is atomically replaced only after the generation validates. A failed refresh +leaves the previous pointer and generation active. On a first installation, an +interruption between the generation rename and pointer update is reported as +`recovery-required`; the next explicit installation recovers matching local +content without network access. If a refresh is interrupted at that point, the +old generation remains active and the completed unreferenced generation is +retained; an explicit retry with `force=True` performs the requested refresh. + +Generations are retained so paths already returned to readers remain valid. +DataCache does not garbage-collect them. Publication uses the cross-platform +`filelock` lock plus same-filesystem atomic renames and supports ordinary local +filesystems on Linux, macOS and Windows. Distributed coordination, arbitrary +network-filesystem semantics, and durability after sudden power loss are outside +the contract. + +With `download_options={"resume": True}`, remote part files persist privately +across attempts and each remote part must have a trusted hash and size. A single +part may use the assembled archive expectations. Retry, timeout and progress +options otherwise have the same meanings as `fetch_file`; progress covers +network transfer, not concatenation or extraction. diff --git a/docs/bundles.md b/docs/bundles.md index a0eda25..7357006 100644 --- a/docs/bundles.md +++ b/docs/bundles.md @@ -158,6 +158,11 @@ progress, disk-space and partial-discard details. ## Adopting from downstream libraries +- **MHCflurry:** use [`install_archive`](archives.md) for released `.tar.bz2` + trees and ordered historical parts. Resolve member paths from the returned + generation; preserve `DOWNLOAD_INFO.csv` with `extra_files`. Do not enumerate + model files as bundle assets or treat the managed store's existence as a + completed download. - **hitlist / tsarina:** use [VersionedFileRegistry](file_registry.md) to retain fixed paths, single-Path returns and legacy root manifests without moving old caches. For a deliberate migration to generation bundles, the existing diff --git a/docs/file_registry.md b/docs/file_registry.md index 6258a3a..c34666d 100644 --- a/docs/file_registry.md +++ b/docs/file_registry.md @@ -19,7 +19,7 @@ registry = VersionedFileRegistry({ }, cache_dir=lambda: Path("existing-cache")) expected = registry.local_path("reference") # works before installation; no writes -status = registry.status() # presence and legacy receipt, offline +status = registry.status() # presence, URL/hash receipt, offline path = registry.ensure("reference", timeout=60, record_provenance=True) ``` diff --git a/docs/integration.md b/docs/integration.md index 3ac8473..e708756 100644 --- a/docs/integration.md +++ b/docs/integration.md @@ -27,6 +27,34 @@ helper transform flags retain the public parsed-URL behavior. 7. Let original network and filesystem exceptions propagate, or retain them as causes when adding domain-specific context. `FileValidationError.path` and `.reason` provide structured validation details. +8. Use `install_archive` when a released tar file owns a directory tree. Treat + only `inspect_archive(...).status == "available"` as installed and resolve + member paths from the returned generation, never from the managed store. + +## Reusable downstream toolkit + +Downstream libraries should keep domain catalogues, CLI wording and generated +indexes, while delegating acquisition and source-store lifecycle to DataCache: + +| Downstream shape | DataCache API | Shared behavior | +| --- | --- | --- | +| One established file path, as in PyEnsembl annotation/FASTA sources | `fetch_file` or `VersionedFileRegistry` | Atomic download, decompression, retries, resume, progress, integrity and provenance | +| A released tar tree, as in MHCflurry weights/data | `VersionedArchiveRegistry` or `install_archive` | Ordered parts, safe extraction, consumer receipts, immutable publication, status and recovery | +| Several separately published files forming one source snapshot | `VersionedDatasetRegistry` or `install_bundle` | Pinned versions, all-files-before-pointer publication, inspection and recovery | + +All explicit download methods forward `show_progress`, callback, timeout and +bounded retry options to the shared transport. Registry construction, status, +inspection and local-path resolution are offline. Applications can therefore +offer consistent `list`/`info`/`download` behavior without maintaining private +network or extraction implementations. + +PyEnsembl can preserve its fixed paths and derived SQLite/FASTA indexes while +continuing to use `fetch_file`; direct callers can set `record_provenance=True` +and expose `inspect_file` results for download visibility. Those derived +artifacts remain application-owned. + +MHCflurry can preserve release selection, exact URL receipts and public model +paths through the archive registry's `store_path` callback and generation path. ## Shared OpenVax cache @@ -74,7 +102,10 @@ logging handlers or change the process umask. Download retry counts and waits are bounded, while `timeout` limits connection/read inactivity per attempt, not total elapsed time. -Use [versioned bundles](bundles.md) when files must be installed together. +Use [versioned bundles](bundles.md) when named files must be installed together, +or [archive installation](archives.md) when a tar archive defines the complete +tree. Archive stores use cross-platform locks; generation bundles currently +require POSIX `flock`. Individual atomic downloads do not provide multi-file transactions or distributed locking. Unhandled termination can leave staging files behind. SQLite operations depend on SQLite's filesystem locking and journaling. New database publication needs diff --git a/tests/test_archives.py b/tests/test_archives.py new file mode 100644 index 0000000..e792207 --- /dev/null +++ b/tests/test_archives.py @@ -0,0 +1,563 @@ +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Transactional whole-archive installation for directory-tree consumers.""" + +from concurrent.futures import ThreadPoolExecutor +from hashlib import sha256 +import io +import json +import os +from pathlib import Path +import tarfile +import threading + +import pytest + +from datacache import ( + FileValidationError, VersionedArchiveRegistry, inspect_archive, + install_archive, +) +from datacache import archives, download + + +def make_tar(path, members, *, mode="w:bz2"): + with tarfile.open(path, mode) as archive: + for name, value in members: + if isinstance(value, tarfile.TarInfo): + archive.addfile(value) + continue + info = tarfile.TarInfo(name) + if value is None: + info.type = tarfile.DIRTYPE + archive.addfile(info) + else: + info.size = len(value) + archive.addfile(info, io.BytesIO(value)) + return path + + +@pytest.fixture +def archive_file(tmp_path): + path = make_tar(tmp_path / "models.tar.bz2", [ + ("models", None), + ("models/a/model.json", b'{"model": "a"}\n'), + ("models/b/model.json", b'{"model": "b"}\n'), + ("README.txt", b"released weights\n"), + ]) + data = path.read_bytes() + return path, data, sha256(data).hexdigest() + + +def test_install_inspect_and_read_only_reuse(tmp_path, archive_file, monkeypatch): + source, data, digest = archive_file + destination = tmp_path / "download" + receipt = "url\nhttps://example.test/models.tar.bz2\n" + sources = [{"url": "https://example.test/models.tar.bz2", "path": source}] + + generation = install_archive( + destination, sources, expected_sha256=digest, expected_size=len(data), + extra_files={"DOWNLOAD_INFO.csv": receipt}) + + assert (generation / "models/a/model.json").read_text() == '{"model": "a"}\n' + assert (generation / "DOWNLOAD_INFO.csv").read_text() == receipt + inspected = inspect_archive( + destination, sources, expected_sha256=digest, expected_size=len(data), + extra_files={"DOWNLOAD_INFO.csv": receipt}) + assert inspected.status == "available" and inspected.verified + assert Path(inspected.generation) == generation + assert inspected.source_urls == ("https://example.test/models.tar.bz2",) + assert inspected.archive_size == len(data) + assert inspected.recorded_sha256 == digest + assert inspected.fetched_at.endswith("+00:00") + assert inspect_archive(destination).status == "available" + assert not inspect_archive(destination).verified + fast = inspect_archive( + destination, sources, expected_sha256=digest, expected_size=len(data), + extra_files={"DOWNLOAD_INFO.csv": receipt}, verify_files=False) + assert fast.status == "available" and not fast.verified + assert fast.files == {} and Path(fast.generation) == generation + + walk_tree = archives._walk_tree + monkeypatch.setattr( + archives, "_walk_tree", + lambda *args, **kwargs: (_ for _ in ()).throw( + AssertionError("fast inspection hashed the tree"))) + assert inspect_archive(destination, sources, verify_files=False).status == "available" + monkeypatch.setattr(archives, "_walk_tree", walk_tree) + + refreshed = install_archive( + destination, sources, expected_sha256=digest, expected_size=len(data), + extra_files={"DOWNLOAD_INFO.csv": receipt}, force=True) + assert refreshed != generation + assert (generation / "models/a/model.json").read_text() == '{"model": "a"}\n' + + def forbidden(*args, **kwargs): + raise AssertionError("a valid cache hit attempted a write, lock, or transfer") + + monkeypatch.setattr(archives, "FileLock", forbidden) + monkeypatch.setattr(download, "fetch_file", forbidden) + monkeypatch.setattr(archives, "write_json", forbidden) + assert install_archive( + destination, sources, expected_sha256=digest, expected_size=len(data), + extra_files={"DOWNLOAD_INFO.csv": receipt}) == refreshed + + +def test_ordered_split_archive_and_logical_source_urls(tmp_path, archive_file): + _, data, _ = archive_file + cuts = (len(data) // 3, 2 * len(data) // 3) + pieces = (data[:cuts[0]], data[cuts[0]:cuts[1]], data[cuts[1]:]) + sources, urls = [], [] + for index, piece in enumerate(pieces): + path = tmp_path / ("models.part.%d" % index) + path.write_bytes(piece) + url = "https://example.test/models.tar.bz2.part.%d?token=secret" % index + urls.append(url) + sources.append({"url": url, "path": path}) + csv_receipt = "url\n" + "".join(url.split("?", 1)[0] + "\n" for url in urls) + destination = tmp_path / "split" + + generation = install_archive( + destination, sources, verified=False, + extra_files={"DOWNLOAD_INFO.csv": csv_receipt}) + + assert (generation / "models/b/model.json").is_file() + assert (generation / "DOWNLOAD_INFO.csv").read_text() == csv_receipt + state = inspect_archive(destination, sources, extra_files={"DOWNLOAD_INFO.csv": csv_receipt}) + assert state.status == "available" and not state.verified + manifest = json.loads((generation / archives.MANIFEST).read_text()) + assert [item["url"] for item in manifest["sources"]] == [url.split("?", 1)[0] for url in urls] + assert [item["fingerprint"] for item in manifest["sources"]] != [ + sha256(url.split("?", 1)[0].encode()).hexdigest() for url in urls] + + assert inspect_archive( + destination, list(reversed(sources)), extra_files={"DOWNLOAD_INFO.csv": csv_receipt}).status == "invalid" + with pytest.raises(FileValidationError, match="force=True"): + install_archive( + destination, list(reversed(sources)), verified=False, + extra_files={"DOWNLOAD_INFO.csv": csv_receipt}) + + +def test_trusted_part_hashes_allow_mirror_reuse(tmp_path, archive_file): + _, data, _ = archive_file + middle = len(data) // 2 + pieces = [data[:middle], data[middle:]] + sources = [] + for index, piece in enumerate(pieces): + path = tmp_path / ("part-%d" % index) + path.write_bytes(piece) + sources.append({ + "url": "https://one.example/part-%d" % index, + "path": path, + "sha256": sha256(piece).hexdigest(), + "size": len(piece), + }) + destination = tmp_path / "download" + generation = install_archive(destination, sources) + mirrored = [dict(source, url=source["url"].replace("one", "two")) + for source in sources] + + state = inspect_archive(destination, mirrored) + assert state.status == "available" and state.verified + assert install_archive(destination, mirrored) == generation + + +def test_url_parts_use_raw_fetch_and_forward_download_options(tmp_path, archive_file, monkeypatch): + source, data, digest = archive_file + calls = [] + original = download.fetch_file + + def recording(url, **options): + calls.append((url, dict(options))) + return original(url, **options) + + monkeypatch.setattr(download, "fetch_file", recording) + generation = install_archive( + tmp_path / "download", source.as_uri(), + expected_sha256=digest, expected_size=len(data), + download_options={"timeout": 12, "max_retries": 4}) + + assert generation.is_dir() + assert len(calls) == 1 + assert calls[0][1]["raw"] is True + assert calls[0][1]["timeout"] == 12 + assert calls[0][1]["max_retries"] == 4 + assert calls[0][1]["expected_sha256"] == digest + + +@pytest.mark.parametrize("bad_member", [ + tarfile.TarInfo("../escape"), + tarfile.TarInfo("/absolute"), + tarfile.TarInfo("models/link"), + tarfile.TarInfo("models/fifo"), +]) +def test_unsafe_members_are_rejected_without_publication(tmp_path, bad_member): + if bad_member.name.endswith("link"): + bad_member.type = tarfile.SYMTYPE + bad_member.linkname = "../outside" + elif bad_member.name.endswith("fifo"): + bad_member.type = tarfile.FIFOTYPE + else: + bad_member.size = 1 + path = make_tar(tmp_path / "bad.tar.bz2", [(bad_member.name, bad_member)]) + outside = tmp_path / "escape" + destination = tmp_path / "download" + + with pytest.raises(FileValidationError): + install_archive(destination, path, verified=False) + + assert not outside.exists() + assert inspect_archive(destination).status == "missing" + assert not list((destination / "generations").iterdir()) + + +@pytest.mark.parametrize("members", [ + [("A.txt", b"a"), ("a.txt", b"b")], + [("models/file", b"a"), ("models/file", b"b")], + [("models/file", b"a"), ("models/file/child", b"b")], + [(".datacache-archive-manifest.json", b"forged")], +]) +def test_colliding_and_reserved_member_names_are_rejected(tmp_path, members): + path = make_tar(tmp_path / "bad.tar.bz2", members) + with pytest.raises(FileValidationError): + install_archive(tmp_path / "download", path, verified=False) + + +def test_extra_file_collision_is_not_published(tmp_path, archive_file): + source, _, _ = archive_file + destination = tmp_path / "download" + with pytest.raises(FileValidationError, match="collides"): + install_archive( + destination, source, verified=False, + extra_files={"README.txt": "consumer receipt"}) + assert inspect_archive(destination).status == "missing" + + +def test_extra_file_casefolded_parent_collision_is_rejected(tmp_path): + archive = make_tar(tmp_path / "archive.tar.bz2", [("Models/model.json", b"model")]) + with pytest.raises(FileValidationError, match="collides"): + install_archive( + tmp_path / "download", archive, verified=False, + extra_files={"models/DOWNLOAD_INFO.csv": "url\n"}) + + +def test_failed_refresh_preserves_old_generation(tmp_path, archive_file, monkeypatch): + old_source, old_data, old_digest = archive_file + destination = tmp_path / "download" + old = install_archive( + destination, old_source, expected_sha256=old_digest, + expected_size=len(old_data)) + pointer = (destination / archives.CURRENT).read_bytes() + replacement = make_tar(tmp_path / "replacement.tar.bz2", [("new.txt", b"new")]) + replacement_data = replacement.read_bytes() + + def fail(*args, **kwargs): + raise OSError("injected extraction failure") + + monkeypatch.setattr(archives, "_extract_tar", fail) + with pytest.raises(OSError, match="injected"): + install_archive( + destination, replacement, expected_sha256=sha256(replacement_data).hexdigest(), + expected_size=len(replacement_data), force=True) + + assert (destination / archives.CURRENT).read_bytes() == pointer + assert inspect_archive( + destination, old_source, expected_sha256=old_digest, + expected_size=len(old_data)).status == "available" + assert (old / "README.txt").is_file() + + +def test_interrupted_first_extraction_is_hidden_and_retryable( + tmp_path, archive_file, monkeypatch): + source, data, digest = archive_file + destination = tmp_path / "download" + original = archives._extract_tar + + def interrupt(archive_path, tree, *args, **kwargs): + (tree / "partial.txt").write_text("not complete") + raise OSError("interrupted extraction") + + monkeypatch.setattr(archives, "_extract_tar", interrupt) + with pytest.raises(OSError, match="interrupted extraction"): + install_archive( + destination, source, expected_sha256=digest, expected_size=len(data)) + assert inspect_archive(destination).status == "missing" + assert not list((destination / "generations").iterdir()) + + monkeypatch.setattr(archives, "_extract_tar", original) + generation = install_archive( + destination, source, expected_sha256=digest, expected_size=len(data)) + assert (generation / "models/a/model.json").is_file() + + +def test_interrupted_pointer_publication_recovers_without_source(tmp_path, archive_file, monkeypatch): + source, data, digest = archive_file + destination = tmp_path / "download" + original = archives.write_json + + def interrupt(path, value, **kwargs): + if Path(path).name == archives.CURRENT: + raise KeyboardInterrupt("after generation rename") + return original(path, value, **kwargs) + + monkeypatch.setattr(archives, "write_json", interrupt) + with pytest.raises(KeyboardInterrupt): + install_archive( + destination, source, expected_sha256=digest, expected_size=len(data)) + assert inspect_archive( + destination, source, expected_sha256=digest, + expected_size=len(data)).status == "recovery-required" + + monkeypatch.setattr(archives, "write_json", original) + source.unlink() + + def forbidden(*args, **kwargs): + raise AssertionError("recovery attempted to reacquire the archive") + + monkeypatch.setattr(archives, "_assemble_archive", forbidden) + generation = install_archive( + destination, source, expected_sha256=digest, expected_size=len(data)) + assert generation.is_dir() + assert inspect_archive( + destination, source, expected_sha256=digest, + expected_size=len(data)).verified + + +def test_failed_manifest_write_leaves_no_apparently_installed_tree( + tmp_path, archive_file, monkeypatch): + source, data, digest = archive_file + destination = tmp_path / "download" + original = archives.write_json + + def fail(path, value, **kwargs): + if Path(path).name == archives.MANIFEST: + raise OSError("receipt publication failed") + return original(path, value, **kwargs) + + monkeypatch.setattr(archives, "write_json", fail) + with pytest.raises(OSError, match="receipt publication"): + install_archive( + destination, source, expected_sha256=digest, expected_size=len(data)) + assert inspect_archive(destination).status == "missing" + assert not list((destination / "generations").iterdir()) + + monkeypatch.setattr(archives, "write_json", original) + generation = install_archive( + destination, source, expected_sha256=digest, expected_size=len(data)) + assert generation.is_dir() + + +def test_foreign_directories_are_never_claimed_even_with_force(tmp_path, archive_file): + source, data, digest = archive_file + for name, contents in (("empty", None), ("legacy", "keep me")): + destination = tmp_path / name + destination.mkdir() + if contents: + (destination / "legacy.txt").write_text(contents) + with pytest.raises((FileNotFoundError, FileValidationError)): + install_archive( + destination, source, expected_sha256=digest, + expected_size=len(data), force=True) + assert not contents or (destination / "legacy.txt").read_text() == contents + + +def test_tree_changes_and_extra_files_are_detected(tmp_path, archive_file): + source, data, digest = archive_file + destination = tmp_path / "download" + generation = install_archive( + destination, source, expected_sha256=digest, expected_size=len(data)) + (generation / "README.txt").write_text("changed") + assert inspect_archive(destination).status == "invalid" + (generation / "unexpected.txt").write_text("extra") + assert inspect_archive(destination).status == "invalid" + with pytest.raises(FileValidationError, match="force=True"): + install_archive( + destination, source, expected_sha256=digest, expected_size=len(data)) + repaired = install_archive( + destination, source, expected_sha256=digest, + expected_size=len(data), force=True) + assert repaired != generation + assert inspect_archive(destination).verified is False + assert inspect_archive( + destination, source, expected_sha256=digest, + expected_size=len(data)).verified + + +def test_links_added_after_installation_are_invalid_and_never_followed( + tmp_path, archive_file): + source, data, digest = archive_file + destination = tmp_path / "download" + generation = install_archive( + destination, source, expected_sha256=digest, expected_size=len(data)) + outside = tmp_path / "outside.txt" + outside.write_text("outside") + target = generation / "README.txt" + target.unlink() + try: + target.symlink_to(outside) + except (OSError, NotImplementedError): + pytest.skip("symlinks are unavailable") + assert inspect_archive(destination).status == "invalid" + assert outside.read_text() == "outside" + + +@pytest.mark.skipif(os.name != "posix", reason="POSIX permission modes") +@pytest.mark.parametrize("creation_mask", [0o022, 0o002, 0o077]) +def test_published_tree_uses_normal_creation_permissions( + tmp_path, archive_file, creation_mask): + source, data, digest = archive_file + previous = os.umask(creation_mask) + try: + destination = tmp_path / "download" + generation = install_archive( + destination, source, expected_sha256=digest, expected_size=len(data)) + finally: + os.umask(previous) + assert destination.stat().st_mode & 0o777 == 0o777 & ~creation_mask + assert generation.stat().st_mode & 0o777 == 0o777 & ~creation_mask + assert (generation / "README.txt").stat().st_mode & 0o777 == 0o666 & ~creation_mask + assert (generation / archives.MANIFEST).stat().st_mode & 0o777 == 0o666 & ~creation_mask + assert (destination / archives.CURRENT).stat().st_mode & 0o777 == 0o666 & ~creation_mask + + +def test_concurrent_installers_converge_on_one_generation(tmp_path, archive_file): + source, data, digest = archive_file + destination = tmp_path / "download" + barrier = threading.Barrier(4) + + def install(_): + barrier.wait(timeout=5) + return install_archive( + destination, source, expected_sha256=digest, expected_size=len(data)) + + with ThreadPoolExecutor(max_workers=4) as pool: + results = list(pool.map(install, range(4))) + + assert len(set(results)) == 1 + assert len(list((destination / "generations").iterdir())) == 1 + + with ThreadPoolExecutor(max_workers=2) as pool: + refreshed = list(pool.map( + lambda _: install_archive( + destination, source, expected_sha256=digest, + expected_size=len(data), force=True), + range(2))) + assert len(set(refreshed)) == 2 + assert all(path.is_dir() for path in results + refreshed) + assert inspect_archive( + destination, source, expected_sha256=digest, + expected_size=len(data)).verified + + +def test_limits_and_verification_requirements_fail_before_publication(tmp_path, archive_file): + source, data, digest = archive_file + with pytest.raises(ValueError, match="verified archive"): + install_archive(tmp_path / "unverified", source) + assert not (tmp_path / "unverified").exists() + with pytest.raises(FileValidationError, match="members"): + install_archive( + tmp_path / "members", source, expected_sha256=digest, + expected_size=len(data), max_members=1) + with pytest.raises(FileValidationError, match="expands"): + install_archive( + tmp_path / "size", source, expected_sha256=digest, + expected_size=len(data), max_extracted_size=1) + + +def test_resumable_split_downloads_require_part_integrity(tmp_path, archive_file): + source, data, digest = archive_file + half = len(data) // 2 + first, second = tmp_path / "one", tmp_path / "two" + first.write_bytes(data[:half]) + second.write_bytes(data[half:]) + with pytest.raises(ValueError, match="resumable archive parts"): + install_archive( + tmp_path / "download", [first.as_uri(), second.as_uri()], + expected_sha256=digest, expected_size=len(data), + download_options={"resume": True}) + + +def test_versioned_archive_registry_supports_consumer_layout_and_status( + tmp_path, archive_file): + source, data, digest = archive_file + url = "https://example.test/models.tar.bz2" + receipt = "url\n%s\n" % url + paths = [] + + def consumer_path(name, version): + paths.append((name, version)) + return tmp_path / "releases" / version / name + + registry = VersionedArchiveRegistry({ + "models": { + "default_version": "2.3.0", + "description": "Presentation models", + "versions": { + "2.3.0": { + "sources": url, + "expected_sha256": digest, + "expected_size": len(data), + "extra_files": {"DOWNLOAD_INFO.csv": receipt}, + }, + "2.2.0": { + "sources": "https://example.test/older.tar.bz2", + "expected_sha256": digest, + "expected_size": len(data), + }, + }, + }, + }, store_path=consumer_path) + + assert registry.resolve_version("models") == "2.3.0" + assert registry.store_path("models") == tmp_path / "releases/2.3.0/models" + assert not (tmp_path / "releases").exists() + rows = registry.status() + assert [(row["version"], row["status"]) for row in rows] == [ + ("2.3.0", "missing"), ("2.2.0", "missing")] + assert rows[0]["default"] and rows[0]["sources"] == [url] + assert rows[0]["downloaded_sources"] == [] + + generation = registry.download("models", source_paths=source) + assert (generation / "models/a/model.json").is_file() + assert registry.is_cached("models") + assert registry.local_path("models") == generation + assert registry.inspect("models").verified + assert not registry.is_cached("models", "2.2.0") + installed = registry.status("models")[0] + assert installed["generation"] == str(generation) + assert installed["archive_size"] == len(data) + assert installed["downloaded_sources"] == [url] + assert ("models", "2.3.0") in paths + + source.unlink() + assert registry.local_path("models") == generation + with pytest.raises(FileNotFoundError): + registry.local_path("models", "2.2.0") + with pytest.raises(ValueError, match="unknown archive"): + registry.resolve_version("unknown") + + +def test_versioned_archive_registry_default_layout_and_unverified_history( + tmp_path, archive_file): + source, _, _ = archive_file + registry = VersionedArchiveRegistry({ + "data": { + "default_version": "historical", + "versions": {"historical": source}, + }, + }, cache_root=tmp_path / "cache", verified=False) + + assert registry.store_path("data") == tmp_path / "cache/data/historical" + generation = registry.ensure("data") + assert generation.is_dir() + assert registry.inspect("data").status == "available" + assert not registry.inspect("data").verified diff --git a/tests/test_file_registry.py b/tests/test_file_registry.py index 4e135f6..98b05b5 100644 --- a/tests/test_file_registry.py +++ b/tests/test_file_registry.py @@ -33,7 +33,8 @@ def test_missing_lookup_and_status_are_read_only(registry, tmp_path): assert reg.status() == [{ 'name': 'thing', 'description': 'reference', 'default_version': 'v2', 'available_versions': ['v1', 'v2'], 'cached': False, 'cached_version': None, - 'bytes': None, 'downloaded_at': None, 'path': str(reg.local_path('thing')), + 'url': None, 'bytes': None, 'sha256': None, 'downloaded_at': None, + 'path': str(reg.local_path('thing')), }] assert not (tmp_path / 'cache').exists() @@ -55,7 +56,10 @@ def test_download_receipt_and_multiple_versions(registry, tmp_path): second = reg.ensure('thing') assert first.exists() and second.exists() assert first != second - assert reg.status()[0]['cached_version'] == 'v2' + status = reg.status()[0] + assert status['cached_version'] == 'v2' + assert status['url'] == source.as_uri() + assert status['sha256'] == hashlib.sha256(b'original').hexdigest() assert not list(manifest_path.parent.glob('.datacache-json-*')) From 18b82b337642b9d62264dd5a5ab6ec9db1fb4e01 Mon Sep 17 00:00:00 2001 From: Alex Rubinsteyn Date: Wed, 30 Sep 2026 15:50:42 -0400 Subject: [PATCH 2/3] Keep unsafe archive fixtures portable --- tests/test_archives.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/tests/test_archives.py b/tests/test_archives.py index e792207..ab48ea3 100644 --- a/tests/test_archives.py +++ b/tests/test_archives.py @@ -206,8 +206,6 @@ def test_unsafe_members_are_rejected_without_publication(tmp_path, bad_member): bad_member.linkname = "../outside" elif bad_member.name.endswith("fifo"): bad_member.type = tarfile.FIFOTYPE - else: - bad_member.size = 1 path = make_tar(tmp_path / "bad.tar.bz2", [(bad_member.name, bad_member)]) outside = tmp_path / "escape" destination = tmp_path / "download" From 0b06637c9ba4f95d919b60a05a45d54bf63d4f1b Mon Sep 17 00:00:00 2001 From: Alex Rubinsteyn Date: Wed, 30 Sep 2026 22:56:50 -0400 Subject: [PATCH 3/3] Fix archive review findings before release --- CHANGELOG.md | 5 ++ datacache/_filesystem.py | 5 +- datacache/archives.py | 66 ++++++++++------- docs/api.md | 6 +- docs/archives.md | 19 +++-- tests/test_archives.py | 153 ++++++++++++++++++++++++++++++++++++++- 6 files changed, 220 insertions(+), 34 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index e5d4dff..818dcd3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,8 @@ concatenated byte-for-byte; URL and local-path sources share the same API. Extraction rejects traversal, links, special files, path collisions and reserved metadata names, and ignores archive ownership and permission bits. + Standard `.` / `./` tar paths are normalized safely; resource limits stop + header scanning at the first violation, and local FIFOs cannot block a writer. - Add `VersionedArchiveRegistry`, a catalogue facade with pinned defaults, explicit versions, reusable status rows, fast local path resolution, consumer layout callbacks, local predownloaded-part overrides, and the same progress, @@ -17,6 +19,9 @@ same transaction with `extra_files`. Archive receipts record the ordered source identities, assembled and per-part observed hashes, and every extracted file. Trusted expectations remain distinct from observed consistency checks. + Reuse checks the complete consumer metadata inventory, including removals. + Large manifests are read without the small control-record size cap and + validated in staging before the generation is published. - Serialize archive writers with a cross-platform file lock, retain immutable generations across refreshes, and recover a complete generation after an interrupted pointer publication. Legacy or otherwise foreign directories are diff --git a/datacache/_filesystem.py b/datacache/_filesystem.py index 38675ef..4676763 100644 --- a/datacache/_filesystem.py +++ b/datacache/_filesystem.py @@ -58,9 +58,10 @@ def file_lock(path, *, private=False): def read_json(path, limit=1024 * 1024): + """Read a regular JSON file; ``limit=None`` permits growing tree manifests.""" with os.fdopen(open_regular(path), 'rb') as handle: - data = handle.read(limit + 1) - if len(data) > limit: + data = handle.read() if limit is None else handle.read(limit + 1) + if limit is not None and len(data) > limit: raise ValueError('JSON record is too large') return json.loads(data) diff --git a/datacache/archives.py b/datacache/archives.py index 8a18813..ab7087c 100644 --- a/datacache/archives.py +++ b/datacache/archives.py @@ -197,11 +197,13 @@ def _manifest(receipt): archive = receipt.get("archive") sources = receipt.get("sources") files = receipt.get("files") + extra_files = receipt.get("extra_files") directories = receipt.get("directories") fetched_at = receipt.get("fetched_at") if (not isinstance(archive, dict) or not isinstance(sources, list) or not sources or not isinstance(files, dict) or not files or - not isinstance(directories, list) or not isinstance(fetched_at, str)): + not isinstance(extra_files, dict) or not isinstance(directories, list) or + not isinstance(fetched_at, str)): raise ValueError("invalid archive manifest") _validate_expectations(archive.get("sha256"), archive.get("size")) if archive.get("sha256") is None or archive.get("size") is None: @@ -230,6 +232,9 @@ def _manifest(receipt): if spec["sha256"] is None or spec["size"] is None: raise ValueError("extracted file receipt does not identify observed bytes") normalized_files[name] = spec + for name, spec in extra_files.items(): + if name not in normalized_files or spec != normalized_files[name]: + raise ValueError("consumer extra file receipt disagrees with tree manifest") normalized_directories = [] folded = set() for name in directories: @@ -242,7 +247,7 @@ def _manifest(receipt): file_keys = {name.casefold() for name in normalized_files} if len(file_keys) != len(normalized_files) or file_keys & folded: raise ValueError("colliding paths in archive manifest") - return archive, normalized_sources, normalized_files, normalized_directories, fetched_at + return archive, normalized_sources, normalized_files, normalized_directories, fetched_at, extra_files def _hash_handle(handle): @@ -254,12 +259,8 @@ def _hash_handle(handle): def _open_local_source(path): - flags = os.O_RDONLY | getattr(os, "O_BINARY", 0) | getattr(os, "O_NOFOLLOW", 0) - descriptor = os.open(path, flags) + descriptor = open_regular(path, os.O_RDONLY | getattr(os, "O_BINARY", 0)) try: - info = os.fstat(descriptor) - if not stat.S_ISREG(info.st_mode): - raise FileValidationError(path, "expected a regular archive part") return os.fdopen(descriptor, "rb") except BaseException: os.close(descriptor) @@ -307,7 +308,7 @@ def _walk_tree(root, receipt=None): return observed, directories, inspections -def _matches_definition(receipt_archive, receipt_sources, receipt_files, definition): +def _matches_definition(receipt_archive, receipt_sources, receipt_extras, definition): expected_sha256 = definition["expected_sha256"] expected_size = definition["expected_size"] if expected_sha256 is not None and receipt_archive["sha256"] != expected_sha256: @@ -326,18 +327,20 @@ def _matches_definition(receipt_archive, receipt_sources, receipt_files, definit expected = [source["fingerprint"] for source in definition["sources"]] if fingerprints != expected: raise FileValidationError("archive", "ordered archive source identities disagree") + if set(receipt_extras) != set(definition["extra_files"]): + raise FileValidationError("archive", "consumer extra file names disagree with requested definition") for name, value in definition["extra_files"].items(): expected = {"sha256": hashlib.sha256(value).hexdigest(), "size": len(value)} - if receipt_files.get(name) != expected: + if receipt_extras[name] != expected: raise FileValidationError(name, "extra file disagrees with requested content") -def _inspect_generation(store, generation, definition, verify_files=True): - directory = _generation(store, generation) - receipt = read_json(directory / MANIFEST) - archive, sources, files, directories, fetched_at = _manifest(receipt) +def _inspect_tree(store, directory, definition, verify_files=True): + # Tree inventories grow with the archive; the writer has no fixed size cap. + receipt = read_json(directory / MANIFEST, limit=None) + archive, sources, files, directories, fetched_at, extras = _manifest(receipt) if definition is not None: - _matches_definition(archive, sources, files, definition) + _matches_definition(archive, sources, extras, definition) metadata = { "source_urls": tuple(source["url"] for source in sources), "fetched_at": fetched_at, @@ -355,6 +358,11 @@ def _inspect_generation(store, generation, definition, verify_files=True): str(directory), inspections, **metadata) +def _inspect_generation(store, generation, definition, verify_files=True): + return _inspect_tree( + store, _generation(store, generation), definition, verify_files=verify_files) + + def inspect_archive( destination, sources=None, *, expected_sha256=None, expected_size=None, extra_files=None, verify_files=True): @@ -423,16 +431,15 @@ def _recover(path, definition): def _archive_members( archive, max_members=None, max_extracted_size=None, extra_names=()): - members = archive.getmembers() - if not members: - raise FileValidationError(archive.name, "archive is empty") - if max_members is not None and len(members) > max_members: - raise FileValidationError(archive.name, "archive has more than %d members" % max_members) nodes = {} planned = [] total = 0 regular_files = 0 - for member in members: + member_count = 0 + for member in archive: + member_count += 1 + if max_members is not None and member_count > max_members: + raise FileValidationError(archive.name, "archive has more than %d members" % max_members) if member.isdir(): name, kind = member.name.rstrip("/"), "directory" elif member.isfile(): @@ -440,10 +447,18 @@ def _archive_members( if member.size < 0: raise FileValidationError(archive.name, "archive member has negative size") total += member.size + if max_extracted_size is not None and total > max_extracted_size: + raise FileValidationError( + archive.name, "archive expands to more than %d bytes" % max_extracted_size) regular_files += 1 else: raise FileValidationError( archive.name, "archive contains a link or special member: %s" % member.name) + # Standard tar tools use '.' for the root and './' for relative paths. + while name.startswith("./"): + name = name[2:] + if kind == "directory" and name == ".": + continue try: _relative_name(name) except ValueError as error: @@ -469,6 +484,8 @@ def _archive_members( else: nodes[key] = (name, kind, True) planned.append((member, name, kind)) + if not member_count: + raise FileValidationError(archive.name, "archive is empty") for name in extra_names: parts = name.split("/") for index in range(1, len(parts)): @@ -486,9 +503,6 @@ def _archive_members( nodes[name.casefold()] = (name, "file", True) if not regular_files: raise FileValidationError(archive.name, "archive contains no regular files") - if max_extracted_size is not None and total > max_extracted_size: - raise FileValidationError( - archive.name, "archive expands to more than %d bytes" % max_extracted_size) return planned @@ -711,16 +725,18 @@ def install_archive( "archive": archive_record, "sources": source_records, "files": files, + "extra_files": {name: files[name] for name in definition["extra_files"]}, "directories": directories, } write_json(tree / MANIFEST, receipt, mode=_file_mode(tree)) + # Read back and validate the receipt and complete tree while private. + _inspect_tree(path, tree, definition) generation = uuid4().hex final = path / "generations" / generation os.replace(tree, final) published = True - candidate = _inspect_generation(path, generation, definition) write_json(path / CURRENT, {"generation": generation}, mode=_file_mode(path)) - return Path(candidate.generation) + return final finally: if published or not resumable: shutil.rmtree(working) diff --git a/docs/api.md b/docs/api.md index 33dae3a..c02ed42 100644 --- a/docs/api.md +++ b/docs/api.md @@ -1112,10 +1112,14 @@ hashes are still recorded and checked, but do not authenticate the source. `extra_files` maps safe relative names to text or bytes written after extraction and before publication. It is intended for consumer receipts such as `DOWNLOAD_INFO.csv`; collisions with archive content fail the installation. +The complete metadata inventory participates in generation identity: additions, +removals, and content changes require an explicit refresh. `download_options` accepts timeout, chunk size, progress, retry, and resume settings from `fetch_file`. Resumable URL parts each require a trusted hash and size. Optional `max_members` and `max_extracted_size` limits are checked before -files are extracted. +files are extracted, as headers arrive; scanning stops when a limit is exceeded. +Benign leading `./` paths are normalized and root `.` directory entries ignored. +Local special files such as FIFOs are rejected without blocking. **Returns:** the immutable extracted-generation `Path`, not the managed store path. Existing returned paths remain usable across forced refreshes. diff --git a/docs/archives.md b/docs/archives.md index aed93c8..ac23e2a 100644 --- a/docs/archives.md +++ b/docs/archives.md @@ -134,7 +134,9 @@ claims it silently, including with `force=True`. ## Extraction policy DataCache parses tar, gzip-compressed tar, bzip2-compressed tar and xz-compressed -tar archives using format detection. Before writing members it rejects: +tar archives using format detection. Leading `./` components are normalized and +the `.` root directory entry is ignored, so archives created with +`tar -cf archive.tar -C tree .` work. Before writing members it rejects: - absolute, parent-traversing, non-portable, and reserved DataCache paths; - symbolic links, hard links, devices, FIFOs and every other special member; @@ -146,11 +148,17 @@ Extraction writes regular files itself rather than calling `extractall`. Archive ownership, timestamps and permission bits—including executable and set-ID bits— are not applied. New files and directories use ordinary umask-derived modes. `max_members` and `max_extracted_size` can impose application-specific resource -limits based on tar headers before extraction begins. +limits. They are enforced as each header arrives, before traversing the payload +of a member that exceeds the limit; scanning stops at the first violation. +The member count includes root directory entries even though they are not +extracted. Local sources must be regular files; links and special inputs such +as FIFOs are rejected without waiting for a writer. After extraction, the generation receipt records the assembled archive's observed hash and size, every part's observed hash, ordered source fingerprints, -and a sorted hash/size manifest of every installed file. Inspection rejects +and a sorted hash/size manifest of every installed file. The complete consumer +`extra_files` inventory is recorded separately. Manifest reading supports the +same sizes as writing, including inventories larger than 1 MiB. Inspection rejects missing, changed, additional, linked or special files and directory changes. `ArchiveInspection` also exposes the redacted source URLs, fetch time, assembled size and recorded digest for downstream `list` and `info` commands. @@ -168,14 +176,15 @@ tree by default; `verify_files=False` provides a receipt/source-identity-only path lookup and never reports `verified=True`. Invalid installations require an explicit `force=True`; inaccessible installations propagate their permission error. Consumer metadata is part of the requested generation identity, so a -different `DOWNLOAD_INFO.csv` requires an explicit refresh. +different, added, or removed metadata entry requires an explicit refresh. ## Publication, recovery and platforms Writers serialize on a permanent sibling lock. Downloads, concatenation, extraction, consumer metadata and the DataCache receipt complete in private staging. The finished tree is renamed to `generations/` and `current.json` -is atomically replaced only after the generation validates. A failed refresh +is atomically replaced only after the complete tree and its written receipt +have been read back and validated in staging. A failed refresh leaves the previous pointer and generation active. On a first installation, an interruption between the generation rename and pointer update is reported as `recovery-required`; the next explicit installation recovers matching local diff --git a/tests/test_archives.py b/tests/test_archives.py index ab48ea3..acc218e 100644 --- a/tests/test_archives.py +++ b/tests/test_archives.py @@ -14,10 +14,13 @@ from concurrent.futures import ThreadPoolExecutor from hashlib import sha256 +import gzip import io import json import os from pathlib import Path +import subprocess +import sys import tarfile import threading @@ -92,7 +95,9 @@ def test_install_inspect_and_read_only_reuse(tmp_path, archive_file, monkeypatch archives, "_walk_tree", lambda *args, **kwargs: (_ for _ in ()).throw( AssertionError("fast inspection hashed the tree"))) - assert inspect_archive(destination, sources, verify_files=False).status == "available" + assert inspect_archive( + destination, sources, extra_files={"DOWNLOAD_INFO.csv": receipt}, + verify_files=False).status == "available" monkeypatch.setattr(archives, "_walk_tree", walk_tree) refreshed = install_archive( @@ -112,6 +117,113 @@ def forbidden(*args, **kwargs): extra_files={"DOWNLOAD_INFO.csv": receipt}) == refreshed +def test_standard_dot_prefixed_tar_tree_installs_and_inspects(tmp_path): + source = make_tar(tmp_path / "standard.tar.bz2", [ + (".", None), ("././", None), ("./models/", None), + ("./models/model.json", b"model"), ("././README.txt", b"readme"), + ]) + data = source.read_bytes() + destination = tmp_path / "download" + generation = install_archive( + destination, source, expected_sha256=sha256(data).hexdigest(), + expected_size=len(data)) + + assert (generation / "models/model.json").read_bytes() == b"model" + inspected = inspect_archive(destination) + assert inspected.status == "available" + assert set(inspected.files) == {"models/model.json", "README.txt"} + + +@pytest.mark.parametrize("members", [ + [("./file", b"one"), ("file", b"two")], + [("./Models/file", b"one"), ("models/other", b"two")], + [("./../escape", b"unsafe")], + [("./.datacache-archive-manifest.json", b"forged")], +]) +def test_dot_prefix_normalization_preserves_path_safety(tmp_path, members): + source = make_tar(tmp_path / "unsafe.tar.bz2", members) + destination = tmp_path / "download" + with pytest.raises(FileValidationError): + install_archive(destination, source, verified=False) + assert inspect_archive(destination).status == "missing" + assert not (tmp_path / "escape").exists() + + +def test_large_tree_manifest_installs_and_recovers_offline(tmp_path): + source = make_tar(tmp_path / "large.tar.bz2", [ + ("file-%05d" % index, b"") for index in range(12000) + ]) + data = source.read_bytes() + options = {"expected_sha256": sha256(data).hexdigest(), "expected_size": len(data)} + destination = tmp_path / "download" + generation = install_archive(destination, source, **options) + assert (generation / archives.MANIFEST).stat().st_size > 1024 * 1024 + inspected = inspect_archive(destination, source, **options) + assert inspected.status == "available" and inspected.verified + assert len(inspected.files) == 12000 + + (destination / archives.CURRENT).unlink() + source.unlink() + assert inspect_archive(destination).status == "recovery-required" + assert install_archive(destination, source, **options) == generation + assert inspect_archive(destination, source, verify_files=False, **options).status == "available" + + +@pytest.mark.skipif(os.name != "posix", reason="POSIX FIFO sources") +def test_fifo_source_fails_promptly_and_releases_writer_lock(tmp_path, archive_file): + fifo = tmp_path / "source.fifo" + os.mkfifo(fifo) + destination = tmp_path / "download" + # A subprocess timeout makes a blocking open a bounded test failure. + script = """ +import sys +from pathlib import Path +from datacache import FileValidationError, install_archive +try: + install_archive(Path(sys.argv[1]), Path(sys.argv[2]), verified=False) +except FileValidationError: + pass +else: + raise AssertionError("FIFO archive source was accepted") +""" + subprocess.run( + [sys.executable, "-c", script, str(destination), str(fifo)], + timeout=5, check=True, capture_output=True, text=True) + assert inspect_archive(destination).status == "missing" + source, data, digest = archive_file + assert install_archive( + destination, source, expected_sha256=digest, expected_size=len(data)).is_dir() + + +@pytest.mark.parametrize("new_extras", [ + {}, {"DOWNLOAD_INFO.csv": "url\noriginal\n"}, + {"DOWNLOAD_INFO.csv": "url\nchanged\n", "INFO.txt": "old"}, + {"NEW_INFO.txt": "new"}, +]) +def test_changed_or_removed_consumer_metadata_requires_refresh( + tmp_path, archive_file, new_extras): + source, data, digest = archive_file + options = {"expected_sha256": digest, "expected_size": len(data)} + destination = tmp_path / "download" + old_extras = {"DOWNLOAD_INFO.csv": "url\noriginal\n", "INFO.txt": "old"} + old = install_archive(destination, source, extra_files=old_extras, **options) + manifest = json.loads((old / archives.MANIFEST).read_text()) + assert set(manifest["extra_files"]) == set(old_extras) + for verify_files in (True, False): + assert inspect_archive( + destination, source, extra_files=new_extras, + verify_files=verify_files, **options).status == "invalid" + with pytest.raises(FileValidationError, match="force=True"): + install_archive(destination, source, extra_files=new_extras, **options) + + refreshed = install_archive( + destination, source, extra_files=new_extras, force=True, **options) + assert refreshed != old + assert inspect_archive(destination, source, extra_files=new_extras, **options).verified + assert {name for name in old_extras if (refreshed / name).exists()} == set(old_extras) & set(new_extras) + assert all((old / name).read_text() == value for name, value in old_extras.items()) + + def test_ordered_split_archive_and_logical_source_urls(tmp_path, archive_file): _, data, _ = archive_file cuts = (len(data) // 3, 2 * len(data) // 3) @@ -354,6 +466,24 @@ def fail(path, value, **kwargs): assert generation.is_dir() +def test_unreadable_generated_manifest_is_rejected_before_generation_rename( + tmp_path, archive_file, monkeypatch): + source, data, digest = archive_file + destination = tmp_path / "download" + original = archives.write_json + + def corrupt(path, value, **kwargs): + original(path, value, **kwargs) + if Path(path).name == archives.MANIFEST: + Path(path).write_text("invalid JSON") + + monkeypatch.setattr(archives, "write_json", corrupt) + with pytest.raises(ValueError): + install_archive(destination, source, expected_sha256=digest, expected_size=len(data)) + assert inspect_archive(destination).status == "missing" + assert not list((destination / "generations").iterdir()) + + def test_foreign_directories_are_never_claimed_even_with_force(tmp_path, archive_file): source, data, digest = archive_file for name, contents in (("empty", None), ("legacy", "keep me")): @@ -471,6 +601,27 @@ def test_limits_and_verification_requirements_fail_before_publication(tmp_path, expected_size=len(data), max_extracted_size=1) +def test_member_limit_stops_before_parsing_later_headers(tmp_path): + # Parsing the deliberately invalid third header must never be necessary. + source = tmp_path / "many.tar" + source.write_bytes( + tarfile.TarInfo("first").tobuf() + tarfile.TarInfo("second").tobuf() + + b"invalid header".ljust(512, b"\0")) + with pytest.raises(FileValidationError, match="more than 1 members"): + install_archive(tmp_path / "download", source, verified=False, max_members=1) + + +def test_expansion_limit_stops_before_traversing_compressed_member_data(tmp_path): + member = tarfile.TarInfo("huge") + member.size = 2 ** 30 + # The advertised payload is absent: scanning to the next header would fail. + source = tmp_path / "huge.tar.gz" + source.write_bytes(gzip.compress(member.tobuf())) + with pytest.raises(FileValidationError, match="expands to more than 1024 bytes"): + install_archive( + tmp_path / "download", source, verified=False, max_extracted_size=1024) + + def test_resumable_split_downloads_require_part_integrity(tmp_path, archive_file): source, data, digest = archive_file half = len(data) // 2