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
78 changes: 60 additions & 18 deletions docs/store_caching.rst
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,8 @@ Each namespace configuration dictionary can have:
The default is ``None`` (no age limit).
- ``size``: optional maximum size in bytes. It sets a per-namespace cache
size budget enforced by evicting least-recently-used items until the
namespace total size is within the configured budget.
namespace total size is within the configured budget. Items bigger than
``size`` are not cached.

Example::

Expand Down Expand Up @@ -58,16 +59,54 @@ Behavior
- Cache keys are identical to primary backend keys (same nesting).
- Soft-deleted items are cached under the same ``.del`` name as primary.
- Soft delete/undelete renames cache entries as well.
- On ``Store.open()`` and ``Store.close()``, cache-enabled namespaces are scanned
to clean up the cache. Cleanup order per namespace is:
- Cache failures are non-fatal and logged as warnings.

1. remove expired cache objects when ``max_age`` is configured,
2. if ``size`` is configured, evict the least-recently-used remaining items
until the namespace total size is ``<= size``.
Eviction
--------

Expired entries are always removed first, even if total size is already below
the ``size`` limit.
- Cache failures are non-fatal and logged as warnings.
For each namespace that has a ``max_age`` or ``size`` limit, the ``Store`` keeps
an in-memory index of the cached items (name, size, time of last use) while it
is opened. The index is ordered by the last use *by this store*: a cache hit,
putting an item into the cache or moving it counts as using it.

- On ``Store.open()`` and ``Store.close()``, the namespace is scanned (listed) and
cleaned up. When scanning, items the store did not know yet are ordered by
their ``ItemInfo.atime``.
- Before an item is put into the cache (by ``store()`` or by a ``load()`` that
was a cache miss), room is made for it, so the namespace total size stays
``<= size`` also while the store is in use.

Cleanup order per namespace is:

1. remove expired cache objects when ``max_age`` is configured,
2. if ``size`` is configured, evict the least-recently-used remaining items
until the namespace total size (plus the size of the item to put into the
cache) is ``<= size``.

Expired entries are always removed first, even if total size is already below
the ``size`` limit. A cache hit never expires anything: an expired item that is
still in the cache is served (and that counts as using it).

Shared caches
-------------

Multiple clients (``Store`` instances, also in different processes) may use the
same cache at the same time, e.g. a cache directory shared by several
processes working with the same content-hash addressed data:

- The posixfs backend stores items atomically, so a client never sees or
evicts an incompletely written cache item.
- If another client has evicted an item, that is a cache miss and the item
gets cached again.
- A client only knows what it has put into the cache itself and what it has
seen when it last scanned the namespace. Thus, a namespace with a ``size``
limit is scanned again after the client has put more than ``size / 4`` bytes
into it. With N clients, the namespace total size can temporarily reach about
``size * (1 + N / 4)``.
- Clients do not see each other's cache hits (see the ``atime`` limitation
below), so a client might evict an item another client frequently uses.
- If the clients use different limits for the same namespace, the smallest
limits win.

Manual Cache Invalidation
-------------------------
Expand All @@ -94,22 +133,25 @@ clients, or if cache corruption is suspected), you can use the
Limitations
-----------

- Eviction by ``max_age`` or ``size`` is open-time and close-time only
(``Store.open()`` / ``Store.close()``), not continuous during
``store()``/``load()`` operations.
- No proactive cache validation/revalidation.
- If an object is deleted in the primary backend by another client, the local
cache will still have a stale object.
- ``max_age`` and LRU-by-``size`` depend on backend ``ItemInfo.atime`` support,
currently that is supported by ``posixfs`` and ``REST`` backends.
- For items a ``Store`` has not used itself since it was opened (items cached
in a previous session or by another client), ``max_age`` and LRU-by-``size``
depend on backend ``ItemInfo.atime`` support, currently that is supported by
``posixfs`` and ``REST`` backends. Filesystems often do not update the atime
for each read (e.g. ``relatime`` or ``noatime`` mounts), so it can be older
than the real last use.
If ``atime`` is 0 (not implemented):

- using ``max_age`` would empty the cache on ``Store.open()`` or ``Store.close()``
- using ``size`` would not work in LRU order, because order can't be
determined
- using ``max_age`` would remove these items from the cache when it is
scanned
- using ``size`` would not evict these items in LRU order, because their
order can't be determined
- If a partial range ``load`` call for an object in a cached namespace causes
a cache miss, the full object will be read from the primary backend and the
cache will be populated with the full object.
cache will be populated with the full object (if it is not bigger than
``size``).

Statistics
----------
Expand Down
186 changes: 157 additions & 29 deletions src/borgstore/store.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@
"""

from binascii import hexlify
from collections import Counter
from collections import Counter, OrderedDict
from contextlib import contextmanager
import enum
from functools import wraps
Expand All @@ -33,6 +33,10 @@

logger = logging.getLogger(__name__)

# a cache namespace with a size limit is rescanned after this store has inserted more than
# size / CACHE_RESCAN_DIVISOR bytes into it, see Store._cache_scan.
CACHE_RESCAN_DIVISOR = 4


class CacheMode(enum.Enum):
C_OFF = "off"
Expand All @@ -57,6 +61,44 @@ class CachePolicy(NamedTuple):
size: Optional[int]


class CacheIndex:
"""
In-memory view of one cache namespace that has a max_age or size limit.

It tells the Store what to evict without listing the cache backend for each operation.
It is this store's view only: what other clients sharing the same cache add, use or
evict is only seen when the namespace is scanned, see Store._cache_scan.
"""

def __init__(self, namespace: str, policy: CachePolicy):
self.namespace = namespace
self.policy = policy
# nested name -> (size, last access timestamp), least recently used entry first.
self.entries: OrderedDict = OrderedDict()
self.total = 0 # sum of the entries' sizes
self.inserted = 0 # bytes this store has put into the cache since the last scan

def add(self, name: str, size: int, last_access: float) -> None:
"""add (or replace) an entry as the most recently used one."""
self.remove(name)
self.entries[name] = (size, last_access)
self.total += size

def remove(self, name: str) -> None:
entry = self.entries.pop(name, None)
if entry is not None:
self.total -= entry[0]

def replace_all(self, entries) -> None:
"""replace the contents by entries, an iterable of (name, size, last_access)."""
self.entries = OrderedDict(
(name, (size, last_access))
for name, size, last_access in sorted(entries, key=lambda entry: (entry[2], entry[0]))
)
self.total = sum(size for size, _ in self.entries.values())
self.inserted = 0


def get_backend(url, permissions=None, quota=None):
"""Parse backend URL and return a backend instance (or None)."""
backend = get_file_backend(url, permissions=permissions, quota=quota)
Expand Down Expand Up @@ -153,6 +195,8 @@ def __init__(
if self.cache_backend is None:
raise BackendURLInvalid(f"Invalid or unsupported Cache Backend URL: {cache_url}")
self._cache_disabled = False
# namespace -> CacheIndex, only for namespaces with a max_age or size limit, only while opened.
self._cache_indexes: dict = {}
self.cache_namespaces = [
entry
for entry in sorted(
Expand Down Expand Up @@ -207,6 +251,12 @@ def _cache_policy_for(self, name: str) -> CachePolicy:
return policy
return CachePolicy(mode=CacheMode.C_OFF, max_age=None, size=None)

def _cache_index_for(self, name: str) -> Optional[CacheIndex]:
for namespace, policy in self.cache_namespaces:
if name.startswith(namespace):
return self._cache_indexes.get(namespace)
return None

@_locked
def set_levels(self, levels: dict, create: bool = False) -> None:
if not levels or not isinstance(levels, dict):
Expand Down Expand Up @@ -281,14 +331,20 @@ def open(self) -> None:
logger.warning(f"borgstore: cache open failed, disabling cache: {err!r}")
self._cache_disabled = True
else:
self._cache_cleanup_expired()
self._cache_indexes = {
namespace: CacheIndex(namespace, policy)
for namespace, policy in self.cache_namespaces
if policy.max_age is not None or policy.size is not None
}
self._cache_cleanup()

@_locked
def close(self) -> None:
self.backend.close()
if self.cache_backend is not None:
if not self._cache_disabled:
self._cache_cleanup_expired()
self._cache_cleanup()
self._cache_indexes = {}
try:
self.cache_backend.close()
except Exception as err:
Expand Down Expand Up @@ -425,17 +481,31 @@ def _cache_load(self, nested_name: str, *, size=None, offset=0) -> Optional[byte
if self.cache_backend is None or self._cache_disabled:
return None
self._stats["cache_load_calls"] += 1
index = self._cache_index_for(nested_name)
try:
value = self.cache_backend.load(nested_name, size=size, offset=offset)
except ObjectNotFound:
self._stats["cache_misses"] += 1
if index is not None:
index.remove(nested_name) # another client has evicted it
return None
except Exception as err:
logger.warning(f"borgstore: cache load failed for {nested_name!r}: {err!r}")
self._stats["cache_errors"] += 1
return None
self._stats["cache_hits"] += 1
self._stats["cache_load_volume"] += len(value)
if index is not None:
entry = index.entries.get(nested_name)
if entry is not None:
item_size = entry[0]
else:
# another client has cached it. value might be only a part of the item.
try:
item_size = self.cache_backend.info(nested_name).size
except Exception:
item_size = len(value)
index.add(nested_name, item_size, time.time())
return value

@_locked
Expand Down Expand Up @@ -478,13 +548,30 @@ def load(self, name: str, *, size=None, offset=0, deleted=False) -> bytes:
def _cache_store(self, nested_name: str, value: StoreValue) -> None:
if self.cache_backend is None or self._cache_disabled:
return
index = self._cache_index_for(nested_name)
if index is not None:
size_limit = index.policy.size
if size_limit is not None and len(value) > size_limit:
# it can never fit. also make sure the cache does not keep a previous value.
self._cache_delete(nested_name)
return
if size_limit is not None and index.inserted > size_limit / CACHE_RESCAN_DIVISOR:
self._cache_scan(index)
# make room before storing, so the cache does not exceed its size limit.
# the value replaces a previous one (if any), so that does not count.
index.remove(nested_name)
self._cache_evict(index, needed=len(value))
self._stats["cache_store_calls"] += 1
try:
self.cache_backend.store(nested_name, value)
self._stats["cache_store_volume"] += len(value)
except Exception as err:
logger.warning(f"borgstore: cache store failed for {nested_name!r}: {err!r}")
self._stats["cache_errors"] += 1
else:
if index is not None:
index.add(nested_name, len(value), time.time())
index.inserted += len(value)

@_locked
def store(self, name: str, value: StoreValue) -> None:
Expand All @@ -511,6 +598,11 @@ def _cache_delete(self, nested_name: str) -> None:
if self.cache_backend is None or self._cache_disabled:
return
self._stats["cache_delete_calls"] += 1
index = self._cache_index_for(nested_name)
if index is not None:
# also if deleting fails: the eviction must not try the same item again and again,
# the next scan brings back an item that is still there.
index.remove(nested_name)
try:
self.cache_backend.delete(nested_name)
except ObjectNotFound:
Expand Down Expand Up @@ -571,13 +663,25 @@ def cache_invalidate(self, name: str, *, deleted: bool = False) -> None:
def _cache_move(self, old_nested: str, new_nested: str) -> None:
if self.cache_backend is None or self._cache_disabled:
return
old_index = self._cache_index_for(old_nested)
entry = old_index.entries.get(old_nested) if old_index is not None else None
try:
self.cache_backend.move(old_nested, new_nested)
except ObjectNotFound:
pass
if old_index is not None:
old_index.remove(old_nested) # another client has evicted it
except Exception as err:
logger.warning(f"borgstore: cache move failed for {old_nested!r}->{new_nested!r}: {err!r}")
self._stats["cache_errors"] += 1
else:
if old_index is not None:
old_index.remove(old_nested)
new_index = self._cache_index_for(new_nested)
if new_index is not None:
if entry is not None:
new_index.add(new_nested, entry[0], time.time()) # moving counts as using it
else:
new_index.remove(new_nested) # unknown item, the next scan or cache hit adds it

@_locked
def move(
Expand Down Expand Up @@ -745,28 +849,52 @@ def defrag(self, sources, *, target=None, algorithm=None, namespace=None, delete
)
return unnest(backend_target, namespace=prefix).removeprefix(prefix)

def _cache_cleanup_expired(self) -> None:
now = time.time()
for namespace, policy in self.cache_namespaces:
if policy.max_age is None and policy.size is None:
continue
try:
items = [info for info in self._cache_list(namespace.rstrip("/")) if not info.directory]
if policy.max_age is not None:
remaining_items = []
for info in items:
if not info.atime or (now - info.atime) > policy.max_age:
self._cache_delete(info.name)
else:
remaining_items.append(info)
items = remaining_items
if policy.size is not None:
total_size = sum(info.size for info in items)
for info in sorted(items, key=lambda entry: (entry.atime, entry.name)):
if total_size <= policy.size:
break
self._cache_delete(info.name)
total_size -= info.size
except Exception as err:
logger.warning(f"borgstore: cache cleanup failed for namespace {namespace!r}: {err!r}")
self._stats["cache_errors"] += 1
def _cache_scan(self, index: CacheIndex) -> None:
"""
Bring index in line with what the cache backend really has in that namespace.

Other clients sharing the cache add and evict items, too. Items this store did not
know yet are sorted in by their atime, for known items the more recent one of our last
access and the atime counts (a backend's atime can lag behind or be not implemented).
"""
namespace = index.namespace
index.inserted = 0 # even if scanning fails, do not retry it for each store operation
entries = []
try:
for info in self._cache_list(namespace.rstrip("/")):
known = index.entries.get(info.name)
last_access = max(known[1], info.atime) if known is not None else info.atime
entries.append((info.name, info.size, last_access))
index.replace_all(entries)
except ObjectNotFound:
# nothing was cached in this namespace yet (or a directory vanished while listing).
index.replace_all(entries)
except Exception as err:
logger.warning(f"borgstore: cache scan failed for namespace {namespace!r}: {err!r}")
self._stats["cache_errors"] += 1

def _cache_evict(self, index: CacheIndex, *, needed: int = 0) -> None:
"""
Evict items that are older than max_age, then evict least recently used items until
<needed> more bytes fit into the size limit.
"""
policy = index.policy
if policy.max_age is not None:
now = time.time()
while index.entries:
name, (size, last_access) = next(iter(index.entries.items()))
# last_access is 0 if the item is only known from a backend that has no atime.
if last_access and (now - last_access) <= policy.max_age:
break
index.remove(name)
self._cache_delete(name)
if policy.size is not None:
while index.entries and index.total + needed > policy.size:
name = next(iter(index.entries))
index.remove(name)
self._cache_delete(name)

def _cache_cleanup(self) -> None:
for index in self._cache_indexes.values():
self._cache_scan(index)
self._cache_evict(index)
Loading
Loading