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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions docs/backends.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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).
20 changes: 20 additions & 0 deletions docs/changes.rst
Original file line number Diff line number Diff line change
@@ -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)
--------------------------

Expand Down
1 change: 1 addition & 0 deletions docs/servers.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 7 additions & 0 deletions docs/store.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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.

69 changes: 57 additions & 12 deletions src/borgstore/backends/_base.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 <sources>: if <sources> is a
# generator (or another iterator), the validation consumes it, so iterating over <sources>
# 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
Expand Down Expand Up @@ -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.

<sources> 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:
Expand All @@ -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")
Expand Down
16 changes: 11 additions & 5 deletions src/borgstore/backends/posixfs.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down
13 changes: 12 additions & 1 deletion src/borgstore/backends/rest.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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()
Expand Down
14 changes: 14 additions & 0 deletions src/borgstore/server/rest.py
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
Loading
Loading