diff --git a/docs/store_caching.rst b/docs/store_caching.rst index 025fd12..c129ef8 100644 --- a/docs/store_caching.rst +++ b/docs/store_caching.rst @@ -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 @@ -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 diff --git a/src/borgstore/store.py b/src/borgstore/store.py index 6870517..63c14bd 100644 --- a/src/borgstore/store.py +++ b/src/borgstore/store.py @@ -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 @@ -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): @@ -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. @@ -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.""" @@ -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).""" @@ -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).""" @@ -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) @@ -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 (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) @@ -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: """ @@ -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) diff --git a/tests/test_cache.py b/tests/test_cache.py index 85c5aa4..65dd4c6 100644 --- a/tests/test_cache.py +++ b/tests/test_cache.py @@ -1,5 +1,10 @@ """Tests for Store optional cache behavior.""" +import logging +import random +import sys +import threading + import pytest import borgstore.store as store_module @@ -903,8 +908,9 @@ def test_cache_shared_hit_on_item_of_other_client(tmp_path): store.destroy() -def test_cache_shared_rescan_sees_other_clients_items(tmp_path): +def test_cache_shared_rescan_sees_other_clients_items(tmp_path, monkeypatch): """After inserting more than size / CACHE_RESCAN_DIVISOR bytes, the store scans the shared cache.""" + monkeypatch.setattr(store_module, "CACHE_SCAN_STEP_TIME", 60) # scan it in one step, also if the machine is slow store, cache_root = make_limited_store(tmp_path, size=1000) store.create() try: @@ -953,7 +959,10 @@ def test_cache_eviction_errors_do_not_fail_main_operations(tmp_path): try: with store: - def failing_delete(_backend_name): + failed = [] + + def failing_delete(backend_name): + failed.append(backend_name) raise RuntimeError("boom") original_delete = store.cache_backend.delete @@ -963,7 +972,10 @@ def failing_delete(_backend_name): store.store(data_name(i), bytes([i]) * 100) for i in range(5): assert store.load(data_name(i)) == bytes([i]) * 100 - assert store.stats["cache_errors"] >= 1 + # 3 of the 5 items only fit after evicting another item. an eviction that failed is not + # tried again, except if a scan has found the item again. + assert len(failed) >= 3 + assert store.stats["cache_errors"] == len(failed) finally: store.cache_backend.delete = original_delete # closing the store scans the cache and evicts what could not be evicted before: @@ -996,3 +1008,261 @@ def test_cache_first_open_does_not_count_errors(tmp_path): assert store.stats["cache_errors"] == 0 finally: store.destroy() + + +def test_cache_scan_in_steps(tmp_path, monkeypatch, caplog): + """While the store is in use, a scan is done in steps: each item put into the cache continues it.""" + monkeypatch.setattr(store_module, "CACHE_SCAN_STEP_TIME", 0) # a step lists one item or directory + store, cache_root = make_limited_store(tmp_path, size=1000) + store.create() + try: + with store: + index = store._cache_indexes["data/"] + fill_shared_cache(tmp_path, [(data_name(100 + i), b"o" * 100) for i in range(10)]) + scanning = [] + with caplog.at_level(logging.DEBUG, logger="borgstore.store"): + for i in range(100): + store.store(data_name(i), bytes([i]) * 100) + assert store.load(data_name(i)) == bytes([i]) * 100 + scanning.append(index.scan is not None) + if True in scanning and not scanning[-1]: + break + assert scanning.count(True) > 1 # the scan took multiple steps ... + assert not scanning[-1] # ... and was finished + # the store knows the other client's items now and has made room for its own item: + assert index.total == cache_usage(cache_root) <= 1000 + messages = [record.getMessage() for record in caplog.records] + messages = [message for message in messages if "cache scan of namespace 'data/'" in message] + assert len(messages) == 1 + assert "10 new, 0 gone" in messages[0] + assert f"{scanning.count(True) + 1} steps" in messages[0] + finally: + store.destroy() + + +def test_cache_scan_in_steps_keeps_changes_by_this_store(tmp_path, monkeypatch): + """What the store changes while a scan is in progress is not overridden by what the scan has (not) seen.""" + monkeypatch.setattr(store_module, "CACHE_SCAN_STEP_TIME", 0) # a step lists one item or directory + store, cache_root = make_limited_store(tmp_path, size=10000) + store.create() + try: + with store: + index = store._cache_indexes["data/"] + for i in range(3): + store.store(data_name(i), bytes([i]) * 100) + fill_shared_cache(tmp_path, [(data_name(9), b"o" * 100)]) # another client caches an item + nested = {i: store.find(data_name(i)) for i in (0, 1, 2, 3, 9)} + # all items are in the same directory. the scan lists 2 directories, then the items 0, 1, 2 and 9. + for _ in range(3): + store._cache_scan(index, max_time=0) + assert set(index.scan.seen) == {nested[0]} + store.delete(data_name(0)) # the scan has seen this item + (cache_root / nested[2]).unlink() # another client evicts an item the scan has not seen yet + # the scan has listed the directory already, so it will not see this item. storing it continues the scan. + store.store(data_name(3), b"n" * 100) + assert set(index.scan.seen) == {nested[0], nested[1]} + while index.scan is not None: + store._cache_scan(index, max_time=0) + assert set(index.entries) == {nested[1], nested[3], nested[9]} + assert index.total == cache_usage(cache_root) == 300 + finally: + store.destroy() + + +def test_close_with_a_scan_in_progress(tmp_path, monkeypatch): + monkeypatch.setattr(store_module, "CACHE_SCAN_STEP_TIME", 0) # a step lists one item or directory + store, cache_root = make_limited_store(tmp_path, size=1000) + store.create() + try: + with store: + index = store._cache_indexes["data/"] + for i in range(7): + store.store(data_name(i), bytes([i]) * 100) + assert index.scan is not None and index.scan.seen # the scan has listed the directory of the items + # another client fills the cache. the scan in progress does not see these items any more. + fill_shared_cache(tmp_path, [(data_name(100 + i), b"o" * 100) for i in range(10)]) + assert cache_usage(cache_root) == 1700 + # closing the store scanned the cache from the start and cleaned it up: + assert cache_usage(cache_root) == 1000 + assert store.stats["cache_errors"] == 0 + finally: + store.destroy() + + +def test_cache_scan_errors_do_not_fail_main_operations(tmp_path): + store, cache_root = make_limited_store(tmp_path, size=1000) + store.create() + try: + with store: + index = store._cache_indexes["data/"] + original_list = store.cache_backend.list + + def failing_list(backend_name): + if backend_name.count("/") == 2: # the scan fails after it has listed 2 directories + raise RuntimeError("boom") + yield from original_list(backend_name) + + store.cache_backend.list = failing_list + try: + for i in range(6): + store.store(data_name(i), bytes([i]) * 100) # the 4th item starts a scan + assert store.load(data_name(i)) == bytes([i]) * 100 + assert store.stats["cache_errors"] == 1 # the failed scan is not tried again for each item + assert index.scan is None + assert index.total == cache_usage(cache_root) == 600 + finally: + store.cache_backend.list = original_list + finally: + store.destroy() + + +def test_cache_scan_uses_mtime_if_there_is_no_atime(tmp_path, monkeypatch): + """If the cache backend has no atime, max_age and size eviction go by the time an item was cached.""" + store, cache_root = make_limited_store(tmp_path, max_age=50, size=250) + store.create() + names_values = [(data_name(i), bytes([i]) * 100) for i in range(4)] + fill_shared_cache(tmp_path, names_values) + nested_names = [store.find(name) for name, _value in names_values] + mtimes = {nested_names[0]: 900.0, nested_names[1]: 990.0, nested_names[2]: 980.0, nested_names[3]: 995.0} + monkeypatch.setattr("borgstore.store.time.time", lambda: 1000.0) + original_list = store.cache_backend.list + + def wrapped_list(backend_name): + for info in original_list(backend_name): + full_name = (backend_name + "/" + info.name) if backend_name else info.name + if full_name in mtimes: + yield info._replace(atime=0, mtime=mtimes[full_name]) + else: + yield info + + store.cache_backend.list = wrapped_list + try: + with store: + # item 0 was older than max_age. of the other 3 items, only 2 fit into size: item 2 was the oldest one. + cached = {info.name for info in store._cache_list("data")} + assert cached == {nested_names[1], nested_names[3]} + finally: + store.cache_backend.list = original_list + store.destroy() + + +def test_cache_scan_keeps_more_recent_last_access(tmp_path, monkeypatch): + """For a known item, a scan only changes the last access if another client has used it more recently.""" + now = 1000.0 + monkeypatch.setattr("borgstore.store.time.time", lambda: now) + store, cache_root = make_limited_store(tmp_path, size=1000) + store.create() + try: + with store: + index = store._cache_indexes["data/"] + for i in range(3): + now = 1000.0 + i + store.store(data_name(i), bytes([i]) * 100) + nested_names = [store.find(data_name(i)) for i in range(3)] + # another client has used item 0 after us. the atimes of the items 1 and 2 lag behind our last access. + atimes = {nested_names[0]: 1010.0, nested_names[1]: 500.0, nested_names[2]: 400.0} + sizes = {nested_names[0]: 100, nested_names[1]: 150, nested_names[2]: 100} # item 1 has another size now + original_list = store.cache_backend.list + + def wrapped_list(backend_name): + for info in original_list(backend_name): + full_name = (backend_name + "/" + info.name) if backend_name else info.name + if full_name in atimes: + yield info._replace(atime=atimes[full_name], size=sizes[full_name]) + else: + yield info + + store.cache_backend.list = wrapped_list + try: + store._cache_scan(index) + finally: + store.cache_backend.list = original_list + assert list(index.entries.items()) == [ + (nested_names[1], (150, 1001.0)), + (nested_names[2], (100, 1002.0)), + (nested_names[0], (100, 1010.0)), + ] + assert index.total == 350 + finally: + store.destroy() + + +def test_cache_miss_of_item_deleted_by_other_client(tmp_path): + """If another client has deleted an item, the index does not keep an entry for it.""" + store, cache_root = make_limited_store(tmp_path, size=1000) + store.create() + try: + with store: + index = store._cache_indexes["data/"] + for i in range(2): + store.store(data_name(i), bytes([i]) * 100) + other, _ = make_store(tmp_path, config=make_config({"data/": {"cache": CacheMode.C_WRITETHROUGH}})) + with other: + other.delete(data_name(0)) + with pytest.raises(ObjectNotFound): + store.load(data_name(0)) + assert set(index.entries) == {store.find(data_name(1))} + assert index.total == cache_usage(cache_root) == 100 + finally: + store.destroy() + + +def test_cache_size_limit_in_mirror_mode(tmp_path): + store, cache_root = make_store(tmp_path, config=make_config({"data/": {"cache": CacheMode.C_MIRROR, "size": 250}})) + store.create() + try: + with store: + for i in range(5): + store.store(data_name(i), bytes([i]) * 100) + assert cache_usage(cache_root) <= 250 + for i in range(5): + assert store.load(data_name(i)) == bytes([i]) * 100 # from the primary, also updates the mirror + assert cache_usage(cache_root) <= 250 + assert store.stats["cache_hits"] == 0 + assert store._cache_indexes["data/"].total == cache_usage(cache_root) == 200 + finally: + store.destroy() + + +def test_cache_shared_by_concurrent_clients(tmp_path): + """Multiple clients use the same size limited cache at the same time, each one with its own view of it.""" + n_clients, n_items, item_size, size = 4, 40, 1000, 10000 + names_values = [(data_name(i), bytes([i]) * item_size) for i in range(n_items)] + store, cache_root = make_limited_store(tmp_path, size=size) + store.create() + with make_store(tmp_path, with_cache_backend=False)[0] as uncached: + for name, value in names_values: + uncached.store(name, value) + all_stats, errors = [], [] + + def client(seed): + try: + rnd = random.Random(seed) + client_store, _ = make_limited_store(tmp_path, size=size) + with client_store: + for _ in range(300): + name, value = names_values[rnd.randrange(n_items)] + if rnd.random() < 0.5: + assert client_store.load(name) == value + else: + offset = rnd.randrange(item_size - 10) + assert client_store.load(name, offset=offset, size=10) == value[offset : offset + 10] + all_stats.append(client_store.stats) + except BaseException as err: + errors.append(err) + + try: + threads = [threading.Thread(target=client, args=(seed,)) for seed in range(n_clients)] + for thread in threads: + thread.start() + for thread in threads: + thread.join() + assert not errors + assert len(all_stats) == n_clients + assert all(stats["cache_hits"] > 0 and stats["cache_misses"] > 0 for stats in all_stats) + if sys.platform != "win32": # windows can not delete or replace a file another client has opened + assert all(stats["cache_errors"] == 0 for stats in all_stats) + # while the clients were working, the cache could be bigger than size. the client that closed last cleaned up: + assert cache_usage(cache_root) <= size + finally: + store.destroy()