From 77839961467f04dcb0859b48a5afa8e39670e12d Mon Sep 17 00:00:00 2001 From: Prince Kumar Date: Fri, 7 Aug 2026 18:21:24 +0000 Subject: [PATCH 1/2] feat(Metadata cache): Unified metadata cache for both list and info --- fsspec/dircache.py | 242 +++++++++++++++++++++++++--------- fsspec/spec.py | 29 +++- fsspec/tests/test_dircache.py | 116 ++++++++++++++++ 3 files changed, 323 insertions(+), 64 deletions(-) create mode 100644 fsspec/tests/test_dircache.py diff --git a/fsspec/dircache.py b/fsspec/dircache.py index eca19566b..6ffb0ad68 100644 --- a/fsspec/dircache.py +++ b/fsspec/dircache.py @@ -1,27 +1,26 @@ import time +from collections import OrderedDict from collections.abc import MutableMapping -from functools import lru_cache class DirCache(MutableMapping): """ - Caching of directory listings, in a structure like:: - - {"path0": [ - {"name": "path0/file0", - "size": 123, - "type": "file", - ... - }, - {"name": "path0/file1", - }, - ... - ], - "path1": [...] - } - - Parameters to this class control listing expiry or indeed turn - caching off + Unified Entry-Index Caching of directory listings and file metadata. + + Decouples single object metadata storage (`_entries`) from directory tree + indexing (`_children`) and listing completeness (`_fully_cached_dirs`). + + Parameters + ---------- + use_listings_cache: bool + If False, this cache never returns items, but always reports KeyError, + and setting items has no effect. + listings_expiry_time: int or float (optional) + Time in seconds that a listing is considered valid. If None, + listings do not expire. + max_paths: int (optional) + The maximum number of path entries retained in cache; 'recent' + refers to when the entry was set or accessed. """ def __init__( @@ -31,65 +30,186 @@ def __init__( max_paths=None, **kwargs, ): - """ - - Parameters - ---------- - use_listings_cache: bool - If False, this cache never returns items, but always reports KeyError, - and setting items has no effect - listings_expiry_time: int or float (optional) - Time in seconds that a listing is considered valid. If None, - listings do not expire. - max_paths: int (optional) - The number of most recent listings that are considered valid; 'recent' - refers to when the entry was set. - """ - self._cache = {} - self._times = {} - if max_paths: - self._q = lru_cache(max_paths + 1)(lambda key: self._cache.pop(key, None)) self.use_listings_cache = use_listings_cache self.listings_expiry_time = listings_expiry_time self.max_paths = max_paths + # Primary metadata store: path -> (file_info_dict, expiry_timestamp) + self._entries = OrderedDict() + # Parent-child index: parent_path -> set(child_paths) + self._children = {} + # Directory listing completeness: dir_path -> expiry_timestamp + self._fully_cached_dirs = {} + + @staticmethod + def _parent(path: str) -> str: + clean = path.rstrip("/") + if "/" not in clean: + return "" + return clean.rsplit("/", 1)[0] + + def get_info(self, path: str): + """ + O(1) lookup for single item metadata without requiring parent directory listing. + """ + if not self.use_listings_cache: + return None + + path = path.rstrip("/") + if path not in self._entries: + return None + + info, expiry = self._entries[path] + if self.listings_expiry_time is not None and time.time() > expiry: + self._evict_entry(path) + return None + + self._entries.move_to_end(path) + return info + + def save_info(self, path: str, info: dict): + """ + Cache single item info without claiming full directory listing completeness. + """ + if not self.use_listings_cache: + return + + path = path.rstrip("/") + parent = self._parent(path) + expiry = ( + time.time() + self.listings_expiry_time + if self.listings_expiry_time is not None + else float("inf") + ) + + self._entries[path] = (info, expiry) + self._entries.move_to_end(path) + + if parent not in self._children: + self._children[parent] = set() + self._children[parent].add(path) + + self._enforce_capacity() + def __getitem__(self, item): - if self.listings_expiry_time is not None: - if self._times.get(item, 0) - time.time() < -self.listings_expiry_time: - del self._cache[item] - if self.max_paths: - self._q(item) - return self._cache[item] # maybe raises KeyError + if not self.use_listings_cache: + raise KeyError(item) - def clear(self): - self._cache.clear() + path = item.rstrip("/") - def __len__(self): - return len(self._cache) + # Check full directory listing + if path in self._fully_cached_dirs: + expiry = self._fully_cached_dirs[path] + if self.listings_expiry_time is not None and time.time() > expiry: + self._invalidate_dir(path) + raise KeyError(item) - def __contains__(self, item): - try: - self[item] - return True - except KeyError: - return False + children = self._children.get(path, set()) + res = [] + for child in list(children): + info = self.get_info(child) + if info is None: + # Child was evicted; directory listing is no longer complete + self._fully_cached_dirs.pop(path, None) + raise KeyError(item) + res.append(info) + return sorted(res, key=lambda x: x.get("name", "")) + + # Fallback to single item info lookup + info = self.get_info(path) + if info is not None: + return [info] + + raise KeyError(item) def __setitem__(self, key, value): if not self.use_listings_cache: return - if self.max_paths: - self._q(key) - self._cache[key] = value - if self.listings_expiry_time is not None: - self._times[key] = time.time() + + dir_path = key.rstrip("/") + expiry = ( + time.time() + self.listings_expiry_time + if self.listings_expiry_time is not None + else float("inf") + ) + + if isinstance(value, list): + child_paths = set() + for item in value: + child_path = item.get("name", "").rstrip("/") + if child_path: + self._entries[child_path] = (item, expiry) + self._entries.move_to_end(child_path) + child_paths.add(child_path) + + parent = self._parent(child_path) + if parent not in self._children: + self._children[parent] = set() + self._children[parent].add(child_path) + + self._children[dir_path] = child_paths + self._fully_cached_dirs[dir_path] = expiry + elif isinstance(value, dict): + self.save_info(dir_path, value) + + self._enforce_capacity() def __delitem__(self, key): - del self._cache[key] + path = key.rstrip("/") + found = False - def __iter__(self): - entries = list(self._cache) + if path in self._fully_cached_dirs: + self._invalidate_dir(path) + found = True + + if path in self._entries: + self._evict_entry(path) + found = True - return (k for k in entries if k in self) + if not found: + raise KeyError(key) + + def _invalidate_dir(self, dir_path: str): + self._fully_cached_dirs.pop(dir_path, None) + children = self._children.pop(dir_path, set()) + for child in children: + if child in self._fully_cached_dirs: + self._invalidate_dir(child) + if child in self._entries: + del self._entries[child] + + def _evict_entry(self, path: str): + self._entries.pop(path, None) + parent = self._parent(path) + if parent in self._children: + self._children[parent].discard(path) + self._fully_cached_dirs.pop(parent, None) + + def _enforce_capacity(self): + if self.max_paths is not None: + while len(self._entries) > self.max_paths: + oldest_path, _ = self._entries.popitem(last=False) + self._evict_entry(oldest_path) + + def clear(self): + self._entries.clear() + self._children.clear() + self._fully_cached_dirs.clear() + + def __len__(self): + return len(self._entries) + + def __contains__(self, item): + path = item.rstrip("/") + if path in self._fully_cached_dirs: + return True + return self.get_info(path) is not None + + def __iter__(self): + keys = list(self._fully_cached_dirs) + [ + k for k in self._entries if k not in self._fully_cached_dirs + ] + return iter(keys) def __reduce__(self): return ( diff --git a/fsspec/spec.py b/fsspec/spec.py index 94d286b04..e92fa864e 100644 --- a/fsspec/spec.py +++ b/fsspec/spec.py @@ -744,19 +744,42 @@ def info(self, path, **kwargs): directory, or something else) and other FS-specific keys. """ path = self._strip_protocol(path) + if not kwargs.get("refresh", False): + try: + cached_info = self.dircache.get_info(path) + if cached_info is not None: + return cached_info + except AttributeError: + pass + out = self.ls(self._parent(path), detail=True, **kwargs) out = [o for o in out if o["name"].rstrip("/") == path] if out: - return out[0] + res = out[0] + try: + self.dircache.save_info(path, res) + except AttributeError: + pass + return res out = self.ls(path, detail=True, **kwargs) path = path.rstrip("/") out1 = [o for o in out if o["name"].rstrip("/") == path] if len(out1) == 1: if "size" not in out1[0]: out1[0]["size"] = None - return out1[0] + res = out1[0] + try: + self.dircache.save_info(path, res) + except AttributeError: + pass + return res elif len(out1) > 1 or out: - return {"name": path, "size": 0, "type": "directory"} + res = {"name": path, "size": 0, "type": "directory"} + try: + self.dircache.save_info(path, res) + except AttributeError: + pass + return res else: raise FileNotFoundError(path) diff --git a/fsspec/tests/test_dircache.py b/fsspec/tests/test_dircache.py new file mode 100644 index 000000000..8d758b018 --- /dev/null +++ b/fsspec/tests/test_dircache.py @@ -0,0 +1,116 @@ +import time +import unittest +from fsspec.dircache import DirCache + + +class TestDirCache(unittest.TestCase): + def test_basic_ls_set_and_get(self): + cache = DirCache() + files = [ + {"name": "dir/file1", "type": "file", "size": 100}, + {"name": "dir/file2", "type": "file", "size": 200}, + ] + cache["dir"] = files + + # Lookup directory + self.assertIn("dir", cache) + retrieved = cache["dir"] + self.assertEqual(len(retrieved), 2) + self.assertEqual(retrieved[0]["name"], "dir/file1") + self.assertEqual(retrieved[1]["name"], "dir/file2") + + def test_single_info_lookup(self): + cache = DirCache() + files = [ + {"name": "dir/file1", "type": "file", "size": 100}, + {"name": "dir/file2", "type": "file", "size": 200}, + ] + cache["dir"] = files + + # O(1) single item lookup + info1 = cache.get_info("dir/file1") + self.assertIsNotNone(info1) + self.assertEqual(info1["size"], 100) + + # Non-existent item + self.assertIsNone(cache.get_info("dir/file3")) + + def test_save_info_standalone(self): + cache = DirCache() + # Save info for a standalone file without dir listing + cache.save_info("dir/file_standalone", {"name": "dir/file_standalone", "type": "file", "size": 500}) + + # Single item lookup hits + info = cache.get_info("dir/file_standalone") + self.assertIsNotNone(info) + self.assertEqual(info["size"], 500) + + # Full dir listing for parent should raise KeyError (not fully cached) + self.assertNotIn("dir", cache._fully_cached_dirs) + with self.assertRaises(KeyError): + _ = cache["dir"] + + def test_dict_assignment_shortcut(self): + cache = DirCache() + info = {"name": "path/file.txt", "type": "file", "size": 1234} + cache["path/file.txt"] = info + + self.assertEqual(cache.get_info("path/file.txt"), info) + self.assertIn("path/file.txt", cache) + + def test_child_deletion_invalidates_parent_listing(self): + cache = DirCache() + files = [ + {"name": "dir/file1", "type": "file", "size": 100}, + {"name": "dir/file2", "type": "file", "size": 200}, + ] + cache["dir"] = files + + # Delete single file + del cache["dir/file1"] + self.assertIsNone(cache.get_info("dir/file1")) + + # Parent directory full listing should now be invalidated + self.assertNotIn("dir", cache._fully_cached_dirs) + with self.assertRaises(KeyError): + _ = cache["dir"] + + def test_eviction_invalidates_parent_listing(self): + cache = DirCache(max_paths=2) + # Store listing of 2 items under 'dir' + cache["dir"] = [ + {"name": "dir/file1", "type": "file", "size": 100}, + {"name": "dir/file2", "type": "file", "size": 200}, + ] + self.assertIn("dir", cache._fully_cached_dirs) + + # Adding a 3rd item forces LRU eviction of file1 + cache.save_info("dir2/file3", {"name": "dir2/file3", "type": "file", "size": 300}) + + # file1 was evicted, so 'dir' completeness MUST be invalidated + self.assertNotIn("dir", cache._fully_cached_dirs) + + def test_ttl_expiry(self): + cache = DirCache(listings_expiry_time=0.1) + cache.save_info("file1", {"name": "file1", "size": 10}) + self.assertIsNotNone(cache.get_info("file1")) + + time.sleep(0.15) + self.assertIsNone(cache.get_info("file1")) + + def test_clear_and_len(self): + cache = DirCache() + cache["dir"] = [ + {"name": "dir/file1", "type": "file", "size": 100}, + {"name": "dir/file2", "type": "file", "size": 200}, + ] + self.assertEqual(len(cache), 2) + + cache.clear() + self.assertEqual(len(cache), 0) + self.assertNotIn("dir", cache) + self.assertIsNone(cache.get_info("dir/file1")) + + +if __name__ == "__main__": + unittest.main() From 94aece83948e8549db6aca919ae396be787e128b Mon Sep 17 00:00:00 2001 From: Prince Kumar Date: Fri, 7 Aug 2026 18:33:43 +0000 Subject: [PATCH 2/2] thread safe --- fsspec/dircache.py | 88 ++++++++++++++++++----------------- fsspec/tests/test_dircache.py | 83 ++++++++++++++++++++++++++++++++- 2 files changed, 126 insertions(+), 45 deletions(-) diff --git a/fsspec/dircache.py b/fsspec/dircache.py index 6ffb0ad68..7b8864ce4 100644 --- a/fsspec/dircache.py +++ b/fsspec/dircache.py @@ -1,11 +1,22 @@ +import functools +import threading import time -from collections import OrderedDict +from collections import OrderedDict, defaultdict from collections.abc import MutableMapping +def _locked(func): + @functools.wraps(func) + def wrapper(self, *args, **kwargs): + with self._lock: + return func(self, *args, **kwargs) + + return wrapper + + class DirCache(MutableMapping): """ - Unified Entry-Index Caching of directory listings and file metadata. + Thread-safe Unified Entry-Index Caching of directory listings and file metadata. Decouples single object metadata storage (`_entries`) from directory tree indexing (`_children`) and listing completeness (`_fully_cached_dirs`). @@ -34,11 +45,9 @@ def __init__( self.listings_expiry_time = listings_expiry_time self.max_paths = max_paths - # Primary metadata store: path -> (file_info_dict, expiry_timestamp) + self._lock = threading.RLock() self._entries = OrderedDict() - # Parent-child index: parent_path -> set(child_paths) - self._children = {} - # Directory listing completeness: dir_path -> expiry_timestamp + self._children = defaultdict(set) self._fully_cached_dirs = {} @staticmethod @@ -48,10 +57,16 @@ def _parent(path: str) -> str: return "" return clean.rsplit("/", 1)[0] + def _calc_expiry(self) -> float: + return ( + time.time() + self.listings_expiry_time + if self.listings_expiry_time is not None + else float("inf") + ) + + @_locked def get_info(self, path: str): - """ - O(1) lookup for single item metadata without requiring parent directory listing. - """ + """O(1) thread-safe lookup for single item metadata.""" if not self.use_listings_cache: return None @@ -60,37 +75,29 @@ def get_info(self, path: str): return None info, expiry = self._entries[path] - if self.listings_expiry_time is not None and time.time() > expiry: + if time.time() > expiry: self._evict_entry(path) return None self._entries.move_to_end(path) return info - def save_info(self, path: str, info: dict): - """ - Cache single item info without claiming full directory listing completeness. - """ + @_locked + def save_info(self, path: str, info: dict, expiry: float | None = None): + """Thread-safe cache for single item info.""" if not self.use_listings_cache: return path = path.rstrip("/") parent = self._parent(path) - expiry = ( - time.time() + self.listings_expiry_time - if self.listings_expiry_time is not None - else float("inf") - ) + expiry = expiry if expiry is not None else self._calc_expiry() self._entries[path] = (info, expiry) self._entries.move_to_end(path) - - if parent not in self._children: - self._children[parent] = set() self._children[parent].add(path) - self._enforce_capacity() + @_locked def __getitem__(self, item): if not self.use_listings_cache: raise KeyError(item) @@ -99,17 +106,14 @@ def __getitem__(self, item): # Check full directory listing if path in self._fully_cached_dirs: - expiry = self._fully_cached_dirs[path] - if self.listings_expiry_time is not None and time.time() > expiry: + if time.time() > self._fully_cached_dirs[path]: self._invalidate_dir(path) raise KeyError(item) - children = self._children.get(path, set()) res = [] - for child in list(children): + for child in list(self._children.get(path, set())): info = self.get_info(child) if info is None: - # Child was evicted; directory listing is no longer complete self._fully_cached_dirs.pop(path, None) raise KeyError(item) res.append(info) @@ -122,38 +126,27 @@ def __getitem__(self, item): raise KeyError(item) + @_locked def __setitem__(self, key, value): if not self.use_listings_cache: return dir_path = key.rstrip("/") - expiry = ( - time.time() + self.listings_expiry_time - if self.listings_expiry_time is not None - else float("inf") - ) - if isinstance(value, list): + expiry = self._calc_expiry() child_paths = set() for item in value: child_path = item.get("name", "").rstrip("/") if child_path: - self._entries[child_path] = (item, expiry) - self._entries.move_to_end(child_path) + self.save_info(child_path, item, expiry=expiry) child_paths.add(child_path) - parent = self._parent(child_path) - if parent not in self._children: - self._children[parent] = set() - self._children[parent].add(child_path) - self._children[dir_path] = child_paths self._fully_cached_dirs[dir_path] = expiry elif isinstance(value, dict): self.save_info(dir_path, value) - self._enforce_capacity() - + @_locked def __delitem__(self, key): path = key.rstrip("/") found = False @@ -169,6 +162,7 @@ def __delitem__(self, key): if not found: raise KeyError(key) + @_locked def _invalidate_dir(self, dir_path: str): self._fully_cached_dirs.pop(dir_path, None) children = self._children.pop(dir_path, set()) @@ -178,33 +172,41 @@ def _invalidate_dir(self, dir_path: str): if child in self._entries: del self._entries[child] + @_locked def _evict_entry(self, path: str): self._entries.pop(path, None) parent = self._parent(path) if parent in self._children: self._children[parent].discard(path) + if not self._children[parent]: + del self._children[parent] self._fully_cached_dirs.pop(parent, None) + @_locked def _enforce_capacity(self): if self.max_paths is not None: while len(self._entries) > self.max_paths: oldest_path, _ = self._entries.popitem(last=False) self._evict_entry(oldest_path) + @_locked def clear(self): self._entries.clear() self._children.clear() self._fully_cached_dirs.clear() + @_locked def __len__(self): return len(self._entries) + @_locked def __contains__(self, item): path = item.rstrip("/") if path in self._fully_cached_dirs: return True return self.get_info(path) is not None + @_locked def __iter__(self): keys = list(self._fully_cached_dirs) + [ k for k in self._entries if k not in self._fully_cached_dirs diff --git a/fsspec/tests/test_dircache.py b/fsspec/tests/test_dircache.py index 8d758b018..05c350b10 100644 --- a/fsspec/tests/test_dircache.py +++ b/fsspec/tests/test_dircache.py @@ -1,5 +1,8 @@ +import random import time import unittest +from concurrent.futures import ThreadPoolExecutor + from fsspec.dircache import DirCache @@ -38,7 +41,10 @@ def test_single_info_lookup(self): def test_save_info_standalone(self): cache = DirCache() # Save info for a standalone file without dir listing - cache.save_info("dir/file_standalone", {"name": "dir/file_standalone", "type": "file", "size": 500}) + cache.save_info( + "dir/file_standalone", + {"name": "dir/file_standalone", "type": "file", "size": 500}, + ) # Single item lookup hits info = cache.get_info("dir/file_standalone") @@ -85,7 +91,9 @@ def test_eviction_invalidates_parent_listing(self): self.assertIn("dir", cache._fully_cached_dirs) # Adding a 3rd item forces LRU eviction of file1 - cache.save_info("dir2/file3", {"name": "dir2/file3", "type": "file", "size": 300}) + cache.save_info( + "dir2/file3", {"name": "dir2/file3", "type": "file", "size": 300} + ) # file1 was evicted, so 'dir' completeness MUST be invalidated self.assertNotIn("dir", cache._fully_cached_dirs) @@ -111,6 +119,77 @@ def test_clear_and_len(self): self.assertNotIn("dir", cache) self.assertIsNone(cache.get_info("dir/file1")) + def test_high_stress_concurrency(self): + """ + Stress test thread safety under heavy contention: + 50 threads performing 2000 random operations (reads, writes, deletes, iterations) + with a small max_paths limit to continuously force LRU evictions during active access. + """ + max_capacity = 25 + cache = DirCache(max_paths=max_capacity, listings_expiry_time=1.0) + errors = [] + + def worker(thread_id): + try: + for i in range(40): + folder_id = random.randint(0, 5) + item_id = random.randint(0, 50) + dir_name = f"folder_{folder_id}" + file_name = f"{dir_name}/file_{item_id}.dat" + + op = random.choice( + [ + "save_info", + "get_info", + "set_dir", + "get_dir", + "del_item", + "iter_cache", + "contains", + ] + ) + + if op == "save_info": + cache.save_info(file_name, {"name": file_name, "size": item_id}) + elif op == "get_info": + _ = cache.get_info(file_name) + elif op == "set_dir": + cache[dir_name] = [ + {"name": f"{dir_name}/item_{k}", "type": "file", "size": k} + for k in range(3) + ] + elif op == "get_dir": + try: + _ = cache[dir_name] + except KeyError: + pass + elif op == "del_item": + try: + del cache[file_name] + except KeyError: + pass + elif op == "iter_cache": + _ = list(cache) + elif op == "contains": + _ = file_name in cache + _ = dir_name in cache + + # Enforce capacity check + self.assertLessEqual(len(cache), max_capacity) + + except Exception as exc: + errors.append(exc) + + with ThreadPoolExecutor(max_workers=50) as executor: + futures = [executor.submit(worker, tid) for tid in range(50)] + for f in futures: + f.result() + + self.assertEqual( + len(errors), 0, f"Concurrent execution produced errors: {errors}" + ) + self.assertLessEqual(len(cache), max_capacity) + if __name__ == "__main__": unittest.main()