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
33 changes: 21 additions & 12 deletions docs/store_caching.rst
Original file line number Diff line number Diff line change
Expand Up @@ -101,8 +101,15 @@ processes working with the same content-hash addressed data:
- 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)``.
into it. So, while N clients are putting items into the cache, the namespace
total size usually is above ``size``, it can reach about
``size * (1 + N / 4)``. It is within ``size`` again when the last of these
clients has closed the store.
- Scanning a namespace with a lot of items takes a while. To not block the
store for that long, the scan is done in steps while the store is in use:
each item that is put into the cache continues the scan for about 5 ms, the
store uses the cache as usual between these steps. ``Store.open()`` and
``Store.close()`` scan the namespace in one go.
- 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
Expand Down Expand Up @@ -138,16 +145,18 @@ Limitations
cache will still have a stale object.
- 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 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
depend on the timestamps the cache backend gives in ``ItemInfo``:

- ``atime`` (``posixfs``, ``sftp`` 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), ``mtime`` is used instead (``s3``
backend). That is the time when the item was put into the cache, so
``max_age`` refers to that time and ``size`` evicts the items first that
were cached first.
- If ``mtime`` is 0 also (``rclone`` backend), using ``max_age`` would remove
these items from the cache when it is scanned and 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 (if it is not bigger than
Expand Down
139 changes: 113 additions & 26 deletions src/borgstore/store.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@
import os
import threading
import time
from typing import Iterator, NamedTuple, Optional
from typing import Generator, Iterator, NamedTuple, Optional

from .utils.nesting import nest, unnest
from .backends._base import ItemInfo, BackendBase, StoreValue, validate_value
Expand All @@ -36,6 +36,9 @@
# 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
# such a rescan is done in steps, so it does not block the store for long if the cache has a lot of
# items: each item that is put into the cache continues the scan for about that long [s].
CACHE_SCAN_STEP_TIME = 0.005


class CacheMode(enum.Enum):
Expand All @@ -61,6 +64,17 @@ class CachePolicy(NamedTuple):
size: Optional[int]


class CacheScan:
"""State of a scan of a cache namespace that is in progress, see Store._cache_scan."""

def __init__(self, infos: Generator[ItemInfo, None, None]):
self.infos = infos # lists the namespace (items and directories)
self.seen: dict = {} # nested name -> (size, last access timestamp), as listed
self.changed: set = set() # names the store has added, used or removed since the scan started
self.steps = 0
self.time = 0.0 # time spent scanning [s]


class CacheIndex:
"""
In-memory view of one cache namespace that has a max_age or size limit.
Expand All @@ -76,7 +90,8 @@ def __init__(self, namespace: str, policy: CachePolicy):
# 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
self.inserted = 0 # bytes this store has put into the cache since it started the last scan
self.scan: Optional[CacheScan] = None # the scan that is in progress

def add(self, name: str, size: int, last_access: float) -> None:
"""add (or replace) an entry as the most recently used one."""
Expand All @@ -88,6 +103,8 @@ def remove(self, name: str) -> None:
entry = self.entries.pop(name, None)
if entry is not None:
self.total -= entry[0]
if self.scan is not None:
self.scan.changed.add(name)

def replace_all(self, entries) -> None:
"""replace the contents by entries, an iterable of (name, size, last_access)."""
Expand All @@ -96,8 +113,53 @@ def replace_all(self, entries) -> None:
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())

def start_scan(self, infos: Generator[ItemInfo, None, None]) -> None:
self.scan = CacheScan(infos)
self.inserted = 0

def abort_scan(self) -> None:
if self.scan is not None:
self.scan.infos.close()
self.scan = None

def finish_scan(self) -> tuple[int, int]:
"""
Bring the index in line with what the scan has seen, return (added, removed) entry counts.

The scan has listed the namespace while the store went on using the cache. It might have
seen an item before the store removed it or have missed an item the store added after that
directory was listed. Thus, for names the store has changed since the scan started, the
index is right and the scan is ignored. For all other names, the scan is right:

- known items the scan has not seen were evicted by another client.
- unknown items the scan has seen were cached by another client.
- for known items, the more recent one of our last access and the listed one counts
(another client might have used the item).
"""
scan, self.scan = self.scan, None
gone = [] # names
new, updated = [], [] # (name, (size, last_access))
for name, (size, last_access) in self.entries.items():
if name not in scan.changed:
seen = scan.seen.get(name)
if seen is None:
gone.append(name)
elif seen[0] != size or seen[1] > last_access:
updated.append((name, (seen[0], max(last_access, seen[1]))))
for name, seen in scan.seen.items():
if name not in self.entries and name not in scan.changed:
new.append((name, seen))
for name in gone:
self.remove(name)
if new or updated:
# these must be sorted in by their last access. if there are none (as usual if the
# cache is not shared), the entries are in the right order already.
entries = dict(self.entries)
entries.update(new + updated)
self.replace_all((name, size, last_access) for name, (size, last_access) in entries.items())
return len(new), len(gone)


def get_backend(url, permissions=None, quota=None):
"""Parse backend URL and return a backend instance (or None)."""
Expand Down Expand Up @@ -555,8 +617,10 @@ def _cache_store(self, nested_name: str, value: StoreValue) -> None:
# 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)
if index.scan is not None or (
size_limit is not None and index.inserted > size_limit / CACHE_RESCAN_DIVISOR
):
self._cache_scan(index, max_time=CACHE_SCAN_STEP_TIME)
# 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)
Expand Down Expand Up @@ -722,13 +786,16 @@ def move(
if self._cache_policy_for(name).mode in {CacheMode.C_WRITETHROUGH, CacheMode.C_MIRROR}:
self._cache_move(nested_name, nested_new_name)

def _cache_list(self, name: str) -> Iterator[ItemInfo]:
def _cache_list(self, name: str, *, dirs: bool = False) -> Generator[ItemInfo, None, None]:
"""list all cached items below <name> (recursively), with dirs=True also the directories."""
if self.cache_backend is None:
return
for info in self.cache_backend.list(name):
if info.directory:
subdir_name = (name + "/" + info.name) if name else info.name
yield from self._cache_list(subdir_name)
if dirs:
yield info._replace(name=subdir_name)
yield from self._cache_list(subdir_name, dirs=dirs)
else:
full_name = (name + "/" + info.name) if name else info.name
yield info._replace(name=full_name)
Expand Down Expand Up @@ -849,29 +916,51 @@ def defrag(self, sources, *, target=None, algorithm=None, namespace=None, delete
)
return unnest(backend_target, namespace=prefix).removeprefix(prefix)

def _cache_scan(self, index: CacheIndex) -> None:
def _cache_scan(self, index: CacheIndex, *, max_time: Optional[float] = None) -> 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).
Bring index in line with what the cache backend really has in that namespace:
other clients sharing the cache add, use and evict items, too.

Scanning means listing the whole namespace, which takes long if it has a lot of items.
If max_time [s] is given, the scan is done in steps: a call starts a scan or continues
the scan that is in progress for about that time, the call that gets to the end of the
listing finishes the scan and updates the index, see CacheIndex.finish_scan.
Between the steps, the store uses the cache and the index as usual.
"""
started = time.perf_counter()
namespace = index.namespace
index.inserted = 0 # even if scanning fails, do not retry it for each store operation
entries = []
if index.scan is None:
# this also resets index.inserted, so a failing scan is not retried for each store operation.
index.start_scan(self._cache_list(namespace.rstrip("/"), dirs=True))
scan = index.scan
scan.steps += 1
finished = False
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)
# a step lists at least one item or directory, so the scan makes progress.
for info in scan.infos:
if not info.directory:
# if the backend has no atime, the mtime (when the item was cached) is the best guess.
scan.seen[info.name] = (info.size, info.atime or info.mtime)
if max_time is not None and time.perf_counter() - started >= max_time:
break
else:
finished = True
except ObjectNotFound:
# nothing was cached in this namespace yet (or a directory vanished while listing).
index.replace_all(entries)
finished = True
except Exception as err:
logger.warning(f"borgstore: cache scan failed for namespace {namespace!r}: {err!r}")
self._stats["cache_errors"] += 1
index.abort_scan()
return
if finished:
added, removed = index.finish_scan()
scan.time += time.perf_counter() - started
if finished:
logger.debug(
f"borgstore: cache scan of namespace {namespace!r} -> {len(index.entries)} items "
f"({added} new, {removed} gone) in {scan.time * 1e3:0.1f}ms ({scan.steps} steps)"
)

def _cache_evict(self, index: CacheIndex, *, needed: int = 0) -> None:
"""
Expand All @@ -883,18 +972,16 @@ def _cache_evict(self, index: CacheIndex, *, needed: int = 0) -> 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.
# last_access is 0 if the item is only known from a backend that has neither atime nor mtime.
if last_access and (now - last_access) <= policy.max_age:
break
index.remove(name)
self._cache_delete(name)
self._cache_delete(name) # also removes it from the index
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)
self._cache_delete(next(iter(index.entries))) # also removes it from the index

def _cache_cleanup(self) -> None:
for index in self._cache_indexes.values():
index.abort_scan() # a scan that was done in steps does not know the latest changes by other clients
self._cache_scan(index)
self._cache_evict(index)
Loading
Loading