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
2 changes: 2 additions & 0 deletions docs/changes.rst
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@ Fixes:
Other changes:

- defrag: implemented on top of gather.
- Store.defrag / Store.gather: all items must be in the given namespace, a "/" in
an item name (or the defrag target name) raises ValueError.


Version 0.6.4 (2026-09-24)
Expand Down
8 changes: 5 additions & 3 deletions docs/store.rst
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ API can be much simpler:
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``.
The namespace is given separately, the item names must not contain "/".
- 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 All @@ -26,9 +27,10 @@ API can be much simpler:
- delete: immediately remove an item from the store (given its key).
- move: implements renaming, soft delete/undelete, and moving to the current
nesting level.
- defrag: general purpose defragmentation helper (copies blocks to new items).
If the target name is computed from the content, the same algorithms as for
hash are supported.
- defrag: general purpose defragmentation helper (copies blocks to new items
in the same namespace). The namespace is given separately, the item names
must not contain "/". If the target name is computed from the content, the
same algorithms as for hash are supported.
- quota: return quota limit and usage (-1 if quotas not enabled or not supported)
- stats: API call counters, time spent in API methods, data volume/throughput.
- latency/bandwidth emulator: see :ref:`store-latency-bandwidth-emulator`.
Expand Down
89 changes: 49 additions & 40 deletions src/borgstore/store.py
Original file line number Diff line number Diff line change
Expand Up @@ -610,61 +610,69 @@ def _cached_load(self, nested_name: str, mode: CacheMode, *, size=None, offset=0
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.
their contents concatenated, in the order given.

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,
sources is a list of (name, offset, size) tuples, as for defrag: all items must be in the
given namespace, item names are without namespace and must not contain "/". 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)
mapped_sources = self._find_sources(sources, namespace=namespace, deleted=deleted)
with self._stats_updater(
"gather", f"gather({len(sources)} ranges, namespace={namespace!r}, deleted={deleted})"
"gather", f"gather({len(mapped_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)
mode = self._cache_policy_for(prefix).mode
if mode in {CacheMode.C_WRITETHROUGH, CacheMode.C_MIRROR}:
# cached namespace: read the ranges like load does it (from the cache, or by
# loading the whole item and caching it).
parts = []
for nested_name, offset, size in mapped_sources:
part = self._cached_load(nested_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)})"
f"Read range error from {nested_name} "
f"(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)
parts.append(part)
result = b"".join(parts)
elif mapped_sources:
# gather all ranges from the backend with one call.
result = self._backend_call(
lambda: self.backend.gather(mapped_sources), key="gather", volume=lambda value: len(value)
)
expected_size = sum(size for _, _, size in backend_sources)
if len(gathered) != expected_size:
expected_size = sum(size for _, _, size in mapped_sources)
if len(result) != expected_size:
raise ReadRangeError(
f"Read range error: gather returned {len(gathered)} bytes, expected {expected_size}"
f"Read range error: gather returned {len(result)} 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)
result = b""
self._stats_update_volume("gather", len(result))
return result

def _find_sources(self, sources, *, namespace, deleted) -> list:
"""
validate the sources of gather / defrag and return them with the nested (backend) item names.

all items must be in the given namespace: a "/" in an item name is rejected, because such an
item could be in a deeper namespace (with other nesting levels and cache policy).
"""
sources = validate_sources(sources)
prefix = (namespace + "/") if namespace else ""
nested_names: dict[str, str] = {}
mapped_sources = []
for name, offset, size in sources:
if name not in nested_names:
if "/" in name:
raise ValueError(f"item name must not contain '/' (namespace is given separately): {name!r}")
nested_names[name] = self.find(prefix + name, deleted=deleted)
mapped_sources.append((nested_names[name], offset, size))
return mapped_sources

def _cache_store(self, nested_name: str, value: StoreValue) -> None:
if self.cache_backend is None or self._cache_disabled:
return
Expand Down Expand Up @@ -945,7 +953,8 @@ def hash(self, name: str, algorithm: str = "sha256", *, deleted: bool = False) -
def defrag(self, sources, *, target=None, algorithm=None, namespace=None, deleted=False) -> str:
"""
efficiently create a new item (target) by combining blocks from existing items (sources)
in the same namespace. item and target names are always without namespace.
in the same namespace. all items must be in the given namespace, item and target names are
without namespace and must not contain "/".

sources is a list of (name, block_offset, block_length) tuples. blocks will be processed
in order of appearance in the list and their contents will be appended to the target item.
Expand All @@ -957,10 +966,10 @@ def defrag(self, sources, *, target=None, algorithm=None, namespace=None, delete
returns the target name.
"""
prefix = (namespace + "/") if namespace else ""
mapped_sources = [
(self.find(prefix + source, deleted=deleted), offset, size) for source, offset, size in sources
]
mapped_sources = self._find_sources(sources, namespace=namespace, deleted=deleted)
if target is not None:
if "/" in target:
raise ValueError(f"item name must not contain '/' (namespace is given separately): {target!r}")
target = self.find(prefix + target, deleted=deleted)

# Note: defrag does not interact with the cache. It creates a new item from
Expand Down
11 changes: 0 additions & 11 deletions tests/test_cache.py
Original file line number Diff line number Diff line change
Expand Up @@ -1319,17 +1319,6 @@ def stats_delta(before, keys):
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()

Expand Down
19 changes: 19 additions & 0 deletions tests/test_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -219,6 +219,25 @@ def test_gather_nested(posixfs_store_created):
assert store.gather([("file1", 2, 3)], namespace=ns, deleted=True) == b"234"


def test_gather_defrag_one_namespace(posixfs_store_created):
# all items must be in the given namespace, so item names must not contain "/".
with posixfs_store_created as store:
store.store("one/file1", b"0123456789")
store.store("two/file2", b"abcdefghij")
# an item of the namespace works:
assert store.gather([("file1", 2, 3)], namespace="one") == b"234"
assert store.defrag([("file1", 2, 3)], target="target", namespace="one") == "target"
# an item of another namespace (or an item name including the namespace) is rejected:
for namespace, name in [("one", "../two/file2"), (None, "two/file2"), (None, "one/file1")]:
with pytest.raises(ValueError, match="must not contain '/'"):
store.gather([(name, 2, 3)], namespace=namespace)
with pytest.raises(ValueError, match="must not contain '/'"):
store.defrag([(name, 2, 3)], target="target", namespace=namespace)
# also for the defrag target:
with pytest.raises(ValueError, match="must not contain '/'"):
store.defrag([("file1", 2, 3)], target="two/target", namespace="one")


@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]}}
Expand Down
Loading