diff --git a/docs/backends.rst b/docs/backends.rst index b9b8c0a..deb47f8 100644 --- a/docs/backends.rst +++ b/docs/backends.rst @@ -186,6 +186,7 @@ Use a storage backend running inside a BorgStore REST server process: - Authentication: Optional Basic Auth is supported. - hash: runs the hexdigest computation server-side. Using algorithm "blake3" requires the optional ``blake3`` package to be installed **on the server**. +- gather: runs server-side, all ranges are read with one roundtrip. - defrag: runs the defragmentation helper server-side. - atime / mtime: supported (if backend used by server supports it). mtime is stamped by the server side (the store operation executes there). diff --git a/docs/changes.rst b/docs/changes.rst index c1a12d8..821413c 100644 --- a/docs/changes.rst +++ b/docs/changes.rst @@ -1,6 +1,26 @@ Changelog ========= +Version 0.7.0 (not released yet) +-------------------------------- + +New features: + +- gather: read multiple byte ranges (from one or multiple items) with one call, #211. + The REST backend does it with one roundtrip (server-side gather), other backends + do one partial load per range. + +Fixes: + +- caching: a partial load with a negative offset from a cached namespace returned + wrong data (e.g. nothing for the last N bytes of an item) if the item was loaded + from the primary backend (mirror mode, or a cache miss in writethrough mode). + +Other changes: + +- defrag: implemented on top of gather. + + Version 0.6.4 (2026-09-24) -------------------------- diff --git a/docs/servers.rst b/docs/servers.rst index c034466..7f89185 100644 --- a/docs/servers.rst +++ b/docs/servers.rst @@ -12,6 +12,7 @@ cloud storage servers: hashsum (from http header X-Content-hash-sha256) - server-side hash computation (e.g. sha256, or blake3 if the optional ``blake3`` package is installed on the server) for item content +- server-side gather (reads multiple byte ranges from multiple items with one request) - server-side defragmentation helper (copies blocks to new items) Running the server on host:port diff --git a/docs/store.rst b/docs/store.rst index 223008f..40f32d3 100644 --- a/docs/store.rst +++ b/docs/store.rst @@ -14,6 +14,10 @@ API can be much simpler: - store: write a new item into the store (providing its key/value pair). - load: read a value from the store (given its key); partial loads specifying an offset and/or size are supported. +- gather: read multiple byte ranges (from one or multiple items in the same + namespace) with one call, returning their contents concatenated in the order + given. The caller knows the sizes it requested, so it can split the result + (e.g. into memoryview slices). A short read raises ``ReadRangeError``. - info: get information about an item via its key (exists, size, ...). - hash: computes the hexdigest for the content of an item (given its key). Supported algorithms are all algorithms supported by ``hashlib`` (e.g. @@ -193,4 +197,7 @@ Scalability chunks before storing them in the store. - Partial loads improve performance by avoiding a full load if only part of the value is needed (e.g., a header with metadata). +- gather improves performance if multiple parts of values are needed: a remote + backend that supports it (e.g. REST) reads all the ranges with one roundtrip, + while a partial load per range would cost one roundtrip each. diff --git a/src/borgstore/backends/_base.py b/src/borgstore/backends/_base.py index da9f69f..7b479a9 100644 --- a/src/borgstore/backends/_base.py +++ b/src/borgstore/backends/_base.py @@ -10,6 +10,7 @@ from ..constants import MAX_NAME_LENGTH, TMP_SUFFIX, HID_SUFFIX from ..utils import hashing +from .errors import ReadRangeError # atime is the last read access UNIX timestamp [s] or 0 if not implemented. # mtime is the last modification UNIX timestamp [s] or 0 if not implemented - it must be @@ -47,6 +48,32 @@ def to_bytes(value: StoreValue) -> bytes: return value if isinstance(value, bytes) else bytes(value) +def validate_sources(sources) -> list: + """Validate the sources given to gather / defrag, return them as a list of (name, offset, size) tuples. + + Each source is a (name, offset, size) tuple (or list, e.g. when it comes from JSON): + name is an item name [str], offset is an int (negative: counted from the end of the item), + size is a non-negative int (the exact amount of bytes wanted). + """ + # always build a new list and let the caller use it instead of : if is a + # generator (or another iterator), the validation consumes it, so iterating over + # again would silently yield nothing. + result = [] + for source in sources: + try: + name, offset, size = source + except (TypeError, ValueError): + raise ValueError(f"source must be a (name, offset, size) tuple, got {source!r}") from None + if not isinstance(name, str): + raise ValueError(f"source name must be a str, got {name!r}") + if not isinstance(offset, int) or isinstance(offset, bool): + raise ValueError(f"source offset must be an int, got {offset!r}") + if not isinstance(size, int) or isinstance(size, bool) or size < 0: + raise ValueError(f"source size must be a non-negative int, got {size!r}") + result.append((name, offset, size)) + return result + + def validate_name(name): """Validate a backend key/name.""" # this is used before an object is accepted for storage and @@ -156,6 +183,33 @@ def delete(self, name: str) -> None: def move(self, curr_name: str, new_name: str) -> None: """rename curr_name to new_name (overwrite target)""" + def gather(self, sources) -> bytes: + """ + Read multiple byte ranges (from one or multiple items) and return their contents + concatenated, in the order given. + + is a list of (name, offset, size) tuples, see validate_sources. The item names + are backend names (with namespace and nesting, as for load). A short read raises + ReadRangeError. For an empty list, b"" is returned. + + The caller knows the sizes of the requested ranges, so it can split the returned data + (e.g. into memoryview slices). + """ + # default implementation: one (partial) load per range, works for all backends. + # might be overridden for performance (e.g. to save one roundtrip per range). + sources = validate_sources(sources) + data_parts = [] + for name, offset, size in sources: + if size == 0: + continue # nothing to read (and some backends reject an empty range request) + chunk = self.load(name, offset=offset, size=size) + if len(chunk) != size: + raise ReadRangeError( + f"Read range error from {name} (requested {size} bytes at offset {offset}, got {len(chunk)})" + ) + data_parts.append(chunk) + return b"".join(data_parts) + def defrag(self, sources, *, target=None, algorithm=None, namespace=None, levels=0) -> str: """ Similar to the higher-level Store.defrag method, with these differences: @@ -168,20 +222,11 @@ def defrag(self, sources, *, target=None, algorithm=None, namespace=None, levels Returns the target item name. """ - # default implementation: slow, but works for all backends. - # might be overridden for performance. + # default implementation: gather the ranges, then store them as a new item. + # works for all backends, might be overridden (e.g. to run it remotely). from ..utils.nesting import nest - from .errors import ReadRangeError - data_parts = [] - for source, offset, size in sources: - chunk = self.load(source, offset=offset, size=size) - if len(chunk) != size: - raise ReadRangeError( - f"Read range error from {source} (requested {size} bytes at offset {offset}, got {len(chunk)})" - ) - data_parts.append(chunk) - data = b"".join(data_parts) + data = self.gather(sources) if target is None: if algorithm is None: raise ValueError("Either target or algorithm must be given for defrag") diff --git a/src/borgstore/backends/posixfs.py b/src/borgstore/backends/posixfs.py index b2b1ee1..8e4406c 100644 --- a/src/borgstore/backends/posixfs.py +++ b/src/borgstore/backends/posixfs.py @@ -21,7 +21,7 @@ except ImportError: fcntl = None # not available on Windows -from ._base import BackendBase, ItemInfo, validate_name, validate_value +from ._base import BackendBase, ItemInfo, validate_name, validate_value, validate_sources from .errors import BackendError, BackendAlreadyExists, BackendDoesNotExist, BackendMustNotBeOpen, BackendMustBeOpen from .errors import ObjectNotFound, PermissionDenied, QuotaExceeded from ..constants import TMP_SUFFIX, QUOTA_STORE_NAME, QUOTA_PERSIST_DELTA, QUOTA_PERSIST_INTERVAL @@ -323,17 +323,23 @@ def _rename_to_new_name(): except FileNotFoundError: raise ObjectNotFound(curr_name) from None + def gather(self, sources) -> bytes: + if not self.opened: + raise BackendMustBeOpen() + sources = validate_sources(sources) + # check the permissions of all sources before reading anything + for name in dict.fromkeys(name for name, _, _ in sources): + self._check_permission(name, "r") + return super().gather(sources) + def defrag(self, sources, *, target=None, algorithm=None, namespace=None, levels=0) -> str: if not self.opened: raise BackendMustBeOpen() - # check all permissions before doing anything + # check the target permission before doing anything (the sources are checked by gather). prefix = namespace.rstrip("/") + "/" if namespace else "" # if target is not given, an item named like content-hash is created in same namespace. check_target = target if target else prefix + "01234567" self._check_permission(check_target, "W") - names = [prefix + source[0] for source in sources] - for name in names: - self._check_permission(name, "r") return super().defrag(sources, target=target, algorithm=algorithm, namespace=namespace, levels=levels) def hash(self, name: str, algorithm: str = "sha256") -> str: diff --git a/src/borgstore/backends/rest.py b/src/borgstore/backends/rest.py index 6d4a2fa..868b660 100644 --- a/src/borgstore/backends/rest.py +++ b/src/borgstore/backends/rest.py @@ -29,7 +29,7 @@ except ImportError: pass -from ._base import BackendBase, ItemInfo, validate_name, validate_value, StoreValue +from ._base import BackendBase, ItemInfo, validate_name, validate_value, validate_sources, StoreValue from ._utils import make_range_header, ignore_sigint from .errors import ( ObjectNotFound, @@ -583,6 +583,17 @@ def move(self, curr_name: str, new_name: str) -> None: response = self._request("post", self._url(""), params={"cmd": "move", "current": curr_name, "new": new_name}) self._handle_response(response, f"{curr_name} -> {new_name}") + @with_reconnect + def gather(self, sources) -> bytes: + self._assert_open() + sources = validate_sources(sources) + for name, _, _ in sources: + validate_name(name) + data = json.dumps(sources).encode("utf-8") + response = self._request("post", self._url(""), params={"cmd": "gather"}, data=data) + self._handle_response(response, "gather") + return response.content + @with_reconnect def defrag(self, sources, *, target=None, algorithm=None, namespace=None, levels=0) -> str: self._assert_open() diff --git a/src/borgstore/server/rest.py b/src/borgstore/server/rest.py index de5327e..5f5113c 100644 --- a/src/borgstore/server/rest.py +++ b/src/borgstore/server/rest.py @@ -225,6 +225,20 @@ def do_POST(self): self._handle_exception(e, "quota") return + if cmd == "gather": + try: + content_length = int(self.headers.get("Content-Length", 0)) + body = self.rfile.read(content_length) + sources = json.loads(body) + with self.server.backend_lock, self.server.backend: + data = self.server.backend.gather(sources) + self.respond(HTTP.OK, data=data, content_type="application/octet-stream") + except ValueError as e: + self.send_error(HTTP.BAD_REQUEST, str(e)) + except Exception as e: + self._handle_exception(e, "gather") + return + if cmd == "defrag": target = self.query.get("target", [None])[0] algorithm = self.query.get("algorithm", [None])[0] diff --git a/src/borgstore/store.py b/src/borgstore/store.py index 63c14bd..8665785 100644 --- a/src/borgstore/store.py +++ b/src/borgstore/store.py @@ -22,7 +22,7 @@ from typing import Generator, Iterator, NamedTuple, Optional from .utils.nesting import nest, unnest -from .backends._base import ItemInfo, BackendBase, StoreValue, validate_value +from .backends._base import ItemInfo, BackendBase, StoreValue, validate_value, validate_sources from .backends.errors import ObjectNotFound, NoBackendGiven, BackendURLInvalid, ReadRangeError # noqa from .backends.posixfs import get_file_backend from .backends.rclone import get_rclone_backend @@ -469,19 +469,21 @@ def stats(self): - Write buffering or cached reads might give a wrong impression. """ st = dict(self._stats) # copy Counter -> generic dict - for key in "info", "load", "store", "delete", "move", "list": + for key in "info", "load", "store", "delete", "move", "list", "gather": # make sure key is present, even if method was not called st[f"{key}_calls"] = st.get(f"{key}_calls", 0) # convert integer ns timings to float s st[f"{key}_time"] = st.get(f"{key}_time", 0) / 1e9 - for key in "load", "store": + for key in "load", "store", "gather": v = st.get(f"{key}_volume", 0) t = st.get(f"{key}_time", 0) st[f"{key}_throughput"] = v / t if t else 0 st["backend_load_calls"] = st.get("backend_load_calls", 0) + st["backend_gather_calls"] = st.get("backend_gather_calls", 0) st["backend_store_calls"] = st.get("backend_store_calls", 0) st["backend_delete_calls"] = st.get("backend_delete_calls", 0) st["backend_load_volume"] = st.get("backend_load_volume", 0) + st["backend_gather_volume"] = st.get("backend_gather_volume", 0) st["backend_store_volume"] = st.get("backend_store_volume", 0) st["cache_disabled"] = self._cache_disabled st["cache_hits"] = st.get("cache_hits", 0) @@ -575,38 +577,94 @@ def load(self, name: str, *, size=None, offset=0, deleted=False) -> bytes: with self._stats_updater("load", f"load({name!r}, offset={offset}, size={size}, deleted={deleted})"): cache_policy = self._cache_policy_for(name) nested_name = self.find(name, deleted=deleted) - if cache_policy.mode == CacheMode.C_WRITETHROUGH: - # try a partial read from the cache first, matching the requested range. - cached_value = self._cache_load(nested_name, size=size, offset=offset) - if cached_value is not None: - self._stats_update_volume("load", len(cached_value)) - return cached_value - # cache miss: do a full load from the primary backend and populate the cache. - full_value = self._backend_call( - lambda: self.backend.load(nested_name, size=None, offset=0), - key="load", - volume=lambda value: len(value), - ) - self._cache_store(nested_name, full_value) - elif cache_policy.mode == CacheMode.C_MIRROR: - full_value = self._backend_call( - lambda: self.backend.load(nested_name, size=None, offset=0), - key="load", - volume=lambda value: len(value), - ) - self._cache_store(nested_name, full_value) + if cache_policy.mode in {CacheMode.C_WRITETHROUGH, CacheMode.C_MIRROR}: + result = self._cached_load(nested_name, cache_policy.mode, size=size, offset=offset) else: result = self._backend_call( lambda: self.backend.load(nested_name, size=size, offset=offset), key="load", volume=lambda value: len(value), ) - self._stats_update_volume("load", len(result)) - return result - result = full_value[offset : (None if size is None else offset + size)] self._stats_update_volume("load", len(result)) return result + def _cached_load(self, nested_name: str, mode: CacheMode, *, size=None, offset=0) -> bytes: + """load (a range of) the value of an item in a cached namespace, see load.""" + if mode == CacheMode.C_WRITETHROUGH: + # try a partial read from the cache first, matching the requested range. + cached_value = self._cache_load(nested_name, size=size, offset=offset) + if cached_value is not None: + return cached_value + # cache miss (or mirror mode): do a full load from the primary backend and populate the cache. + full_value = self._backend_call( + lambda: self.backend.load(nested_name, size=None, offset=0), key="load", volume=lambda value: len(value) + ) + self._cache_store(nested_name, full_value) + if offset < 0: + # a negative offset counts from the end of the item. make it absolute, otherwise the end of + # the slice (offset + size) is not right: a range ending at the end of the item would be empty. + offset = max(len(full_value) + offset, 0) + return full_value[offset : (None if size is None else offset + size)] + + @_locked + def gather(self, sources, *, namespace=None, deleted=False) -> bytes: + """ + read multiple byte ranges (from one or multiple items in the same namespace) and return + their contents concatenated, in the order given. item names are always without namespace. + + sources is a list of (name, offset, size) tuples, as for defrag. size must be given + (an int), a short read raises ReadRangeError. the caller knows the sizes it requested, + so it can split the result (e.g. into memoryview slices). + + a backend that supports it (e.g. rest) reads all the ranges with one roundtrip, while + a partial load per range would cost one roundtrip each. + """ + sources = validate_sources(sources) + with self._stats_updater( + "gather", f"gather({len(sources)} ranges, namespace={namespace!r}, deleted={deleted})" + ): + prefix = (namespace + "/") if namespace else "" + nested_names = {} + for name, _, _ in sources: + if name not in nested_names: + nested_names[name] = self.find(prefix + name, deleted=deleted) + # ranges of items in a cached namespace are read like load does it (from the cache, or by + # loading the whole item and caching it), all other ranges are gathered from the backend + # with one call. + parts: list = [None] * len(sources) + backend_sources = [] + for i, (name, offset, size) in enumerate(sources): + mode = self._cache_policy_for(prefix + name).mode + if mode in {CacheMode.C_WRITETHROUGH, CacheMode.C_MIRROR}: + part = self._cached_load(nested_names[name], mode, size=size, offset=offset) + if len(part) != size: + raise ReadRangeError( + f"Read range error from {name} (requested {size} bytes at offset {offset}, got {len(part)})" + ) + parts[i] = part + else: + backend_sources.append((nested_names[name], offset, size)) + gathered = b"" + if backend_sources: + gathered = self._backend_call( + lambda: self.backend.gather(backend_sources), key="gather", volume=lambda value: len(value) + ) + expected_size = sum(size for _, _, size in backend_sources) + if len(gathered) != expected_size: + raise ReadRangeError( + f"Read range error: gather returned {len(gathered)} bytes, expected {expected_size}" + ) + if len(backend_sources) == len(sources): + result = gathered # the usual case: no copy needed + else: + view, pos = memoryview(gathered), 0 + for i, (_, _, size) in enumerate(sources): + if parts[i] is None: + parts[i], pos = view[pos : pos + size], pos + size + result = b"".join(parts) + self._stats_update_volume("gather", len(result)) + return result + def _cache_store(self, nested_name: str, value: StoreValue) -> None: if self.cache_backend is None or self._cache_disabled: return diff --git a/tests/test_backends.py b/tests/test_backends.py index 2c945f7..0bf6bda 100644 --- a/tests/test_backends.py +++ b/tests/test_backends.py @@ -35,6 +35,7 @@ BackendMustBeOpen, BackendMustNotBeOpen, ObjectNotFound, + ReadRangeError, ) from borgstore.backends.posixfs import PosixFS, get_file_backend from borgstore.backends.sftp import Sftp, get_sftp_backend @@ -960,6 +961,34 @@ def test_posixfs_failed_write_leaves_no_tmpfile(posixfs_backend_created, size): assert list_names(backend, ROOTNS) == ["key"] +def test_gather(tested_backends, request): + backend = get_backend_from_fixture(tested_backends, request) + with backend: + backend.store("test/item1", b"0123456789") + backend.store("test/item2", b"abcdefghij") + sources = [("test/item2", 5, 2), ("test/item1", 2, 3), ("test/item1", 0, 1), ("test/item2", -3, 3)] + assert backend.gather(sources) == b"fg234" + b"0" + b"hij" + assert backend.gather(source for source in sources) == b"fg234" + b"0" + b"hij" + assert backend.gather([]) == b"" + assert backend.gather([("test/item1", 0, 0)]) == b"" + + # short read + with pytest.raises(ReadRangeError): + backend.gather([("test/item1", 2, 3), ("test/item2", 5, 20)]) + + # nonexistent object + with pytest.raises(ObjectNotFound): + backend.gather([("test/item1", 0, 1), ("test/nonexistent", 0, 1)]) + + # invalid sources + with pytest.raises(ValueError): + backend.gather([("test/item1", 0, None)]) + + # Test must be open + with pytest.raises(BackendMustBeOpen): + backend.gather([("test/item1", 0, 1)]) + + def test_hash(tested_backends, request): backend = get_backend_from_fixture(tested_backends, request) with backend: diff --git a/tests/test_cache.py b/tests/test_cache.py index 65dd4c6..68e8c4f 100644 --- a/tests/test_cache.py +++ b/tests/test_cache.py @@ -1266,3 +1266,90 @@ def client(seed): assert cache_usage(cache_root) <= size finally: store.destroy() + + +def test_cache_gather(tmp_path): + """gather serves ranges of cached namespaces like load does, the rest is gathered from the backend.""" + config = make_config({"data/": {"cache": "writethrough"}, "meta/": {"cache": "mirror"}}) + store, _ = make_store(tmp_path, config=config) + store.create() + + def stats_delta(before, keys): + return {key: store.stats.get(key, 0) - before.get(key, 0) for key in keys} + + try: + with store: + store.store("data/00000000", b"0123456789") + store.store("meta/00000000", b"abcdefghij") + store.store("config/00000000", b"ABCDEFGHIJ") + store.cache_invalidate("data/00000000") + store.cache_invalidate("meta/00000000") + keys = ["cache_hits", "cache_misses", "cache_store_calls", "backend_load_calls", "backend_gather_calls"] + + # writethrough: the first range misses the cache and loads the whole item into it, + # the second range of the same item is already served from the cache. + before = store.stats + assert store.gather([("00000000", 2, 3), ("00000000", 7, 2)], namespace="data") == b"23478" + assert stats_delta(before, keys) == dict( + cache_hits=1, cache_misses=1, cache_store_calls=1, backend_load_calls=1, backend_gather_calls=0 + ) + # writethrough: cache hits now + before = store.stats + assert store.gather([("00000000", 2, 3), ("00000000", 7, 2)], namespace="data") == b"23478" + assert stats_delta(before, keys) == dict( + cache_hits=2, cache_misses=0, cache_store_calls=0, backend_load_calls=0, backend_gather_calls=0 + ) + # short read from the cached item + with pytest.raises(store_module.ReadRangeError): + store.gather([("00000000", 7, 20)], namespace="data") + + # mirror: always loaded from the primary (and cached) + before = store.stats + assert store.gather([("00000000", 5, 2)], namespace="meta") == b"fg" + assert stats_delta(before, keys) == dict( + cache_hits=0, cache_misses=0, cache_store_calls=1, backend_load_calls=1, backend_gather_calls=0 + ) + + # not cached: gathered from the backend with one call + before = store.stats + assert store.gather([("00000000", 1, 2), ("00000000", 8, 2)], namespace="config") == b"BCIJ" + assert stats_delta(before, keys) == dict( + cache_hits=0, cache_misses=0, cache_store_calls=0, backend_load_calls=0, backend_gather_calls=1 + ) + assert stats_delta(before, ["backend_gather_volume", "gather_volume"]) == dict( + backend_gather_volume=4, gather_volume=4 + ) + + # mixed (no namespace: item names include the namespace): the order of the ranges is kept + before = store.stats + sources = [("config/00000000", 0, 2), ("data/00000000", 0, 2), ("config/00000000", 8, 2)] + assert store.gather(sources) == b"AB01IJ" + assert stats_delta(before, keys) == dict( + cache_hits=1, cache_misses=0, cache_store_calls=0, backend_load_calls=0, backend_gather_calls=1 + ) + assert stats_delta(before, ["backend_gather_volume", "gather_volume"]) == dict( + backend_gather_volume=4, gather_volume=6 + ) + finally: + store.destroy() + + +@pytest.mark.parametrize("mode", ["writethrough", "mirror"]) +def test_cache_load_negative_offset(tmp_path, mode): + """A negative offset counts from the end of the item, also if the item is loaded from the primary backend.""" + store, _ = make_store(tmp_path, config=make_config({"data/": {"cache": mode}})) + store.create() + try: + with store: + name = "data/00000000" + store.store(name, b"0123456789") + for offset, size, expected in [(-3, 3, b"789"), (-3, 2, b"78"), (-3, None, b"789")]: + store.cache_invalidate(name) + # writethrough: the first load is a cache miss (loads from the primary), the second a cache hit. + # mirror: both load from the primary. + assert store.load(name, offset=offset, size=size) == expected + assert store.load(name, offset=offset, size=size) == expected + store.cache_invalidate(name) + assert store.gather([("00000000", -3, 3), ("00000000", 0, 1)], namespace="data") == b"7890" + finally: + store.destroy() diff --git a/tests/test_posixfs_permissions.py b/tests/test_posixfs_permissions.py index eb69031..f399b21 100644 --- a/tests/test_posixfs_permissions.py +++ b/tests/test_posixfs_permissions.py @@ -137,3 +137,29 @@ def test_recursive_permission_shadowing(tmp_path): with pytest.raises(PermissionDenied): fs.store("restricted/item2", DATA2) fs.close() + + +def test_gather_permissions(tmp_path): + fs = PosixFS(path=tmp_path, permissions={"": "w"}) # permissions needed for setup + fs.create() + fs.open() + fs.mkdir("dir1") + fs.mkdir("dir2") + fs.store("dir1/file", DATA1) + fs.store("dir2/file", DATA2) + # r granted for both dirs + fs.permissions = {"dir1": "r", "dir2": "r"} + assert fs.gather([("dir1/file", 0, 2), ("dir2/file", 2, 3)]) == b"da" + b"ta2" + # r denied for dir2: all sources are checked before anything is read + fs.permissions = {"dir1": "r", "dir2": ""} + with pytest.raises(PermissionDenied): + fs.gather([("dir1/file", 0, 2), ("dir2/file", 2, 3)]) + assert fs.gather([("dir1/file", 0, 2)]) == b"da" + # defrag needs r for the sources (checked by gather) and W for the target + fs.permissions = {"dir1": "rW", "dir2": ""} + assert fs.defrag([("dir1/file", 0, 2)], target="dir1/target") == "dir1/target" + with pytest.raises(PermissionDenied): + fs.defrag([("dir1/file", 0, 2), ("dir2/file", 2, 3)], target="dir1/target") + with pytest.raises(PermissionDenied): + fs.defrag([("dir1/file", 0, 2)], target="dir2/target") + fs.close() diff --git a/tests/test_server_rest.py b/tests/test_server_rest.py index 0ecfcb1..a0688e2 100644 --- a/tests/test_server_rest.py +++ b/tests/test_server_rest.py @@ -329,6 +329,53 @@ def test_rest_server_hash_blake3(rest_server_with_auth): be.close() +def test_rest_server_gather(tmp_path): + backend_url = tmp_path.as_uri() + address, port = "127.0.0.1", 0 + username, password = "testuser", "testpassword" + + server, thread = start_server(backend_url, address, port, username, password) + host, assigned_port = server.server_address + url = f"http://{host}:{assigned_port}/" + headers = {"Accept": "application/vnd.x.borgstore.rest.v1"} + auth = (username, password) + + try: + requests.post(url + "?cmd=create", auth=auth, headers=headers).raise_for_status() + requests.post(url + "file1", data=b"0123456789", auth=auth, headers=headers).raise_for_status() + requests.post(url + "file2", data=b"abcdefghij", auth=auth, headers=headers).raise_for_status() + + # gather "234" from file1 (offset 2, size 3) and "fg" from file2 (offset 5, size 2) + sources = [("file1", 2, 3), ("file2", 5, 2)] + response = requests.post(url + "?cmd=gather", data=json.dumps(sources), auth=auth, headers=headers) + response.raise_for_status() + assert response.status_code == 200 + assert response.headers["Content-Type"] == "application/octet-stream" + assert response.content == b"234fg" + + # empty list + response = requests.post(url + "?cmd=gather", data=json.dumps([]), auth=auth, headers=headers) + response.raise_for_status() + assert response.content == b"" + + # short read + response = requests.post(url + "?cmd=gather", data=json.dumps([("file1", 2, 20)]), auth=auth, headers=headers) + assert response.status_code == 416 + assert "requested 20 bytes" in response.text + + # nonexistent source + response = requests.post(url + "?cmd=gather", data=json.dumps([("file3", 0, 1)]), auth=auth, headers=headers) + assert response.status_code == 404 + + # invalid json / invalid sources + for body in [b"this is not json", json.dumps([("file1", 0)]), json.dumps([("file1", 0, "3")])]: + response = requests.post(url + "?cmd=gather", data=body, auth=auth, headers=headers) + assert response.status_code == 400 + finally: + server.shutdown() + server.server_close() + + def test_rest_server_defrag(tmp_path): backend_url = tmp_path.as_uri() address, port = "127.0.0.1", 0 @@ -432,6 +479,31 @@ def test_rest_server_defrag(tmp_path): server.server_close() +def test_rest_backend_gather(rest_server_with_auth): + be = rest_server_with_auth + be.create() + be.open() + try: + be.store("file1", b"0123456789") + be.store("file2", b"abcdefghij") + + sources = [("file2", 5, 2), ("file1", 2, 3), ("file1", 0, 1), ("file2", -3, 3)] + assert be.gather(sources) == b"fg234" + b"0" + b"hij" + assert be.gather([]) == b"" + + with pytest.raises(ReadRangeError) as exc_info: + be.gather([("file1", 2, 20)]) + assert "requested 20 bytes" in str(exc_info.value) + + with pytest.raises(ObjectNotFound): + be.gather([("file1", 0, 1), ("file3", 0, 1)]) + + with pytest.raises(ValueError): + be.gather([("file1", 0, None)]) + finally: + be.close() + + def test_rest_backend_defrag(rest_server_with_auth): be = rest_server_with_auth be.create() diff --git a/tests/test_store.py b/tests/test_store.py index 085a902..4002625 100644 --- a/tests/test_store.py +++ b/tests/test_store.py @@ -17,7 +17,7 @@ from .test_backends import blake3, blake3_is_available from borgstore.constants import ROOTNS -from borgstore.store import Store, ItemInfo, ReadRangeError +from borgstore.store import Store, ItemInfo, ObjectNotFound, ReadRangeError CONFIG = {"zero/": {"levels": [0]}, "one/": {"levels": [1]}, "two/": {"levels": [2]}} # Layout used for most tests @@ -161,6 +161,64 @@ def test_defrag_nested(posixfs_store_created): assert "requested 20 bytes" in str(exc_info.value) +def test_gather_nested(posixfs_store_created): + ns = "two" # nested! CONFIG has {"two/": {"levels": [2]}} + v1 = b"0123456789" + v2 = b"abcdefghij" + with posixfs_store_created as store: + # the stats have the gather keys, even if gather was not called yet + assert store.stats["gather_calls"] == 0 + assert store.stats["backend_gather_calls"] == 0 + assert store.stats["backend_gather_volume"] == 0 + + store.store(ns + "/file1", v1) + store.store(ns + "/file2", v2) + + # ranges from multiple items, multiple ranges from the same item, order is kept + sources = [("file2", 5, 2), ("file1", 2, 3), ("file1", 0, 1), ("file2", -3, 3)] + assert store.gather(sources, namespace=ns) == b"fg" + b"234" + b"0" + b"hij" + assert store.stats["gather_calls"] == 1 + assert store.stats["gather_volume"] == 9 + assert store.stats["backend_gather_calls"] == 1 + assert store.stats["backend_gather_volume"] == 9 + + # the caller can split the result, knowing the sizes + data = memoryview(store.gather(sources, namespace=ns)) + pieces, pos = [], 0 + for _, _, size in sources: + pieces.append(bytes(data[pos : pos + size])) + pos += size + assert pieces == [b"fg", b"234", b"0", b"hij"] + + # a generator works, too (validating the sources must not consume them) + assert store.gather((source for source in sources), namespace=ns) == b"fg234" + b"0" + b"hij" + + # empty ranges and an empty list + assert store.gather([("file1", 3, 0)], namespace=ns) == b"" + assert store.gather([], namespace=ns) == b"" + + # short read + with pytest.raises(ReadRangeError) as exc_info: + store.gather([("file1", 2, 3), ("file2", 5, 20)], namespace=ns) + assert "Read range error from" in str(exc_info.value) + assert "requested 20 bytes" in str(exc_info.value) + + # unknown item + with pytest.raises(ObjectNotFound): + store.gather([("file1", 0, 1), ("nonexistent", 0, 1)], namespace=ns) + + # invalid sources + for sources in [[("file1", 0)], [("file1", "0", 1)], [("file1", 0, -1)], [("file1", 0, None)]]: + with pytest.raises(ValueError): + store.gather(sources, namespace=ns) + + # soft-deleted item + store.move(ns + "/file1", delete=True) + with pytest.raises(ObjectNotFound): + store.gather([("file1", 2, 3)], namespace=ns) + assert store.gather([("file1", 2, 3)], namespace=ns, deleted=True) == b"234" + + @pytest.mark.skipif(not blake3_is_available, reason="blake3 package is not installed") def test_defrag_nested_blake3(posixfs_store_created): ns = "two" # nested! CONFIG has {"two/": {"levels": [2]}}