From 9d1ec029a672d6d5fc68413d42b308549e0d393d Mon Sep 17 00:00:00 2001 From: Jeremy Rapin Date: Thu, 17 Sep 2026 13:20:32 +0200 Subject: [PATCH 1/3] Improve deletion in cache dict (now needs context manager) --- exca/cachedict/core.py | 28 ++++++++++++++----- exca/cachedict/test_cachedict.py | 40 +++++++++++++++++++--------- exca/cachedict/test_legacy_format.py | 6 +++-- exca/map.py | 8 ++++-- exca/steps/backends.py | 4 +-- 5 files changed, 60 insertions(+), 26 deletions(-) diff --git a/exca/cachedict/core.py b/exca/cachedict/core.py index 6c8274ff..b0c6b6f7 100644 --- a/exca/cachedict/core.py +++ b/exca/cachedict/core.py @@ -139,6 +139,7 @@ def __init__( self._folder_modified = -1.0 self._jsonl_readers: dict[str, JsonlReader] = {} self._jsonl_reading_allowance = float("inf") + self._deleted_in_scope = False # DumpContext for this folder (load/delete; writes use per-thread _write_ctx) self._dumper: DumpContext | None = None if self.folder is not None: @@ -185,24 +186,31 @@ def keys(self) -> tp.Iterator[str]: keys = set(self._ram_data) | set(self._key_info) return iter(keys) - def _read_info_files(self, max_workers: int = 4) -> None: + def _read_info_files(self, max_workers: int = 4, force: bool = False) -> None: """Load current info files. Each writer appends to its own JSONL file, so concurrent writes of the same key produce duplicate entries across files. For duplicates, whichever file comes last in iterdir() order wins (non-deterministic); duplicates are kept so explicit deletion - clears every known copy.""" + clears every known copy. + + Parameters + ---------- + max_workers: + Maximum number of threads reading the info files. + force: + Read even if the folder is frozen or looks unmodified. + """ if self.folder is None or not self.folder.exists(): return readings = max((r.readings for r in self._jsonl_readers.values()), default=0) - if self._jsonl_reading_allowance <= readings: - # bypass reloading info files - return + if not force and self._jsonl_reading_allowance <= readings: + return # bypass reloading info files modified = self.folder.lstat().st_mtime nothing_new = self._folder_modified == modified self._folder_modified = modified - if nothing_new: + if nothing_new and not force: logger.debug("Nothing new to read from info files") return # nothing new! cpus = os.cpu_count() @@ -297,11 +305,12 @@ def _write_ctx(self, value: DumpContext | None) -> None: @contextlib.contextmanager def write(self) -> tp.Iterator["CacheDict[X]"]: - """Context manager for writing items to the cache.""" + """Context manager for writing to (and deleting from) the cache.""" if self._write_ctx is not None: raise RuntimeError("Cannot re-open an already open writer") if self.folder is not None: self._write_ctx = DumpContext(self.folder, permissions=self.permissions) + self._deleted_in_scope = False try: if self._write_ctx is not None: with self._write_ctx: @@ -312,6 +321,8 @@ def write(self) -> tp.Iterator["CacheDict[X]"]: self._write_ctx = None if self.folder is not None: utils.best_effort_utime(self.folder) + if self._deleted_in_scope: + self._read_info_files(force=True) # sweep emptied jsonl/data pairs @contextlib.contextmanager def writer(self) -> tp.Iterator["CacheDict[X]"]: @@ -359,6 +370,9 @@ def __delitem__(self, key: str) -> None: if self._dumper is None: del self._ram_data[key] return + if self._write_ctx is None: + raise RuntimeError("Cannot delete outside of a writer context") + self._deleted_in_scope = True if key not in self._key_info: _ = key in self # populate _key_info from disk self._ram_data.pop(key, None) diff --git a/exca/cachedict/test_cachedict.py b/exca/cachedict/test_cachedict.py index 2d0c23f7..85f07414 100644 --- a/exca/cachedict/test_cachedict.py +++ b/exca/cachedict/test_cachedict.py @@ -60,7 +60,10 @@ def test_array_cache(tmp_path: Path, in_ram: bool) -> None: cache2 = cd.CacheDict(folder=folder) assert isinstance(cache2["blublu"], np.ndarray) # del - del cache2["blublu"] + with pytest.raises(RuntimeError, match="writer context"): + del cache2["blublu"] + with cache2.write(): + del cache2["blublu"] assert set(cache2.keys()) == {"blabla"} # clear cache2.clear() @@ -115,6 +118,10 @@ def test_specialized_dump( cache_type = cache_type[:-2] memmap_cache_size = 0 proc = psutil.Process() + try: + proc.open_files() + except (psutil.AccessDenied, PermissionError) as e: + pytest.skip(f"psutil cannot list open files: {e}") cache: cd.CacheDict[tp.Any] = cd.CacheDict( folder=tmp_path, keep_in_ram=keep_in_ram, @@ -200,7 +207,8 @@ def test_info_jsonl_deletion(tmp_path: Path) -> None: assert out.startswith(b"{") and out.endswith(b"}\n") # remove one chosen = np.random.choice(keys) - del cache[chosen] + with cache.write(): + del cache[chosen] assert len(cache) == 2 cache = cd.CacheDict(folder=tmp_path, keep_in_ram=False) assert len(cache) == 2 @@ -215,7 +223,8 @@ def test_info_jsonl_deletion_removes_duplicate_entries(tmp_path: Path) -> None: cache = cd.CacheDict(folder=tmp_path, keep_in_ram=False) assert cache["x"] == 12 - del cache["x"] + with cache.write(): + del cache["x"] cache = cd.CacheDict(folder=tmp_path, keep_in_ram=False) assert "x" not in cache @@ -299,8 +308,11 @@ def test_clone_is_view_only(tmp_path: Path) -> None: assert revived["k"] == 7 +@pytest.mark.parametrize("read_before_delete", [False, True]) @pytest.mark.parametrize("cache_type", ["MemmapArrayFile", "String"]) -def test_orphaned_data_file_cleanup(tmp_path: Path, cache_type: str) -> None: +def test_orphaned_data_file_cleanup( + tmp_path: Path, cache_type: str, read_before_delete: bool +) -> None: """Test that orphaned data files are cleaned up when all items are deleted.""" data: tp.Any = np.random.rand(3, 12) if cache_type == "MemmapArrayFile" else "hello" cache: cd.CacheDict[tp.Any] = cd.CacheDict( @@ -311,13 +323,16 @@ def test_orphaned_data_file_cleanup(tmp_path: Path, cache_type: str) -> None: for c in "abc": ex.submit(_write_items, cache, [f"{c}1", f"{c}2"], data) assert len(list(tmp_path.glob("*-info.jsonl"))) == 3 - # Delete all items from one writer, files still exist (cleanup is lazy) - for key in ["a1", "a2", "c1", "b2"]: - del cache[key] - assert len(list(tmp_path.glob("*-info.jsonl"))) == 3 - # Trigger cleanup via keys() - orphaned pair should be deleted + if read_before_delete: + assert len(set(cache.keys())) == 6 + with cache.write(): + for key in ["a1", "a2", "c1", "b2"]: + del cache[key] + remaining = list(tmp_path.glob("*-info.jsonl")) + assert len(remaining) == 2, ( + f"leaving write() should drop the emptied pair {remaining}" + ) assert set(cache.keys()) == {"b1", "c2"} - assert len(list(tmp_path.glob("*-info.jsonl"))) == 2 @pytest.mark.parametrize( @@ -367,9 +382,8 @@ def test_jsonl_edge_cases(tmp_path: Path, content: str, should_delete: bool) -> # Write and delete an item to trigger reader initialization for our test file with cache.write(): cache["x"] = np.array([1]) - del cache["x"] - # Trigger cleanup - _ = list(cache.keys()) + with cache.write(): + del cache["x"] # Check result for fp in [jsonl, data_file]: if should_delete: diff --git a/exca/cachedict/test_legacy_format.py b/exca/cachedict/test_legacy_format.py index aa59885b..35b274a7 100644 --- a/exca/cachedict/test_legacy_format.py +++ b/exca/cachedict/test_legacy_format.py @@ -141,7 +141,8 @@ def test_legacy_external_static_roundtrip(tmp_path: Path) -> None: assert set(cache.keys()) == {"key1", "key2"} assert cache["key1"] == {"a": 1, "b": [2, 3]} assert cache["key2"] == "hello" - del cache["key1"] + with cache.write(): + del cache["key1"] assert set(cache.keys()) == {"key2"} cache2: cd.CacheDict[tp.Any] = cd.CacheDict(folder=tmp_path) assert set(cache2.keys()) == {"key2"} @@ -189,7 +190,8 @@ def test_mixed_old_and_new_format(tmp_path: Path) -> None: assert cache2["multiline"] == "line1\nline2\nline3" assert cache2["extra"] == "new value" # Delete an old item, verify new items survive - del cache2["hello"] + with cache2.write(): + del cache2["hello"] cache3: cd.CacheDict[tp.Any] = cd.CacheDict(folder=dst) assert set(cache3.keys()) == {"multiline", "extra"} diff --git a/exca/map.py b/exca/map.py index 98e104b3..adb363d0 100644 --- a/exca/map.py +++ b/exca/map.py @@ -273,8 +273,12 @@ def _find_missing(self, items: dict[str, tp.Any]) -> dict[str, tp.Any]: if to_remove: msg = "Clearing %s items for %s (infra.mode=%s)" logger.warning(msg, len(to_remove), self.uid(), self.mode) - for uid in to_remove: - del cache[uid] + assert isinstance(cache, CacheDict), ( + f"to_remove is only filled when caching (got {type(cache)})" + ) + with cache.write(): + for uid in to_remove: + del cache[uid] missing = {x: y for x, y in items.items() if x not in state.recomputed} if isinstance(cache, CacheDict): # dont record computed items if no cache diff --git a/exca/steps/backends.py b/exca/steps/backends.py index 62a92fa5..6796618b 100644 --- a/exca/steps/backends.py +++ b/exca/steps/backends.py @@ -367,7 +367,7 @@ def run_and_cache(self) -> None: "Clearing partial results after invalid _run_batch output: %s", self.paths.step_uid, ) - with self.cache_dict.frozen_cache_folder(): + with self.cache_dict.write(), self.cache_dict.frozen_cache_folder(): for uid in written_uids: if uid in self.cache_dict: del self.cache_dict[uid] @@ -592,7 +592,7 @@ def _clear_caches( logger.warning("Failed to cancel %s%s: %s", paths.step_uid, uids, e) # Success first → a mid-clear crash leaves a recoverable cached # error rather than a stale success (fail closed). - with cd.frozen_cache_folder(): + with cd.write(), cd.frozen_cache_folder(): for uid in uids: if uid in cd: del cd[uid] From 9df235c742baad53f54c5638f6b43cb1e52f3225 Mon Sep 17 00:00:00 2001 From: Jeremy Rapin Date: Thu, 17 Sep 2026 13:31:54 +0200 Subject: [PATCH 2/3] fix --- CHANGELOG.md | 6 ++++++ docs/infra/explanation.md | 2 +- exca/cachedict/core.py | 8 ++++++-- exca/map.py | 5 +---- exca/steps/backends.py | 8 ++++---- 5 files changed, 18 insertions(+), 11 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index eefdd229..bbc96266 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,12 @@ ## [Unreleased] +[breaking] + +- `CacheDict`: deletions require a `write()` context (like writes). [#326] + +[other] + - `DiscriminatedModel`: optimized look-up. [#313] - `steps`: fixed nested infra claim deadlock. [#323] diff --git a/docs/infra/explanation.md b/docs/infra/explanation.md index d70e2e42..81e4158c 100644 --- a/docs/infra/explanation.md +++ b/docs/infra/explanation.md @@ -137,7 +137,7 @@ cache = cachedict.CacheDict(folder=tmp_path, keep_in_ram=True) # the dictionary is empty: assert not cache -# writes require a context manager for efficiency with multiple writes +# writes and deletions require a context manager for efficiency with multiple writes x = np.random.rand(2, 12) with cache.write(): cache["blublu"] = x diff --git a/exca/cachedict/core.py b/exca/cachedict/core.py index b0c6b6f7..bfc72065 100644 --- a/exca/cachedict/core.py +++ b/exca/cachedict/core.py @@ -256,7 +256,7 @@ def _cleanup_orphaned_jsonl_files(self) -> None: except FileNotFoundError: self._jsonl_readers.pop(name, None) continue - logger.warning("Cleaning up orphaned files for %s", name) + logger.debug("Cleaning up orphaned files for %s", name) prefix = name.removesuffix("-info.jsonl") paths = [*self.folder.glob(f"{prefix}.*"), reader._fp] data_dir = self.folder / DumpContext.DATA_DIR @@ -322,7 +322,11 @@ def write(self) -> tp.Iterator["CacheDict[X]"]: if self.folder is not None: utils.best_effort_utime(self.folder) if self._deleted_in_scope: - self._read_info_files(force=True) # sweep emptied jsonl/data pairs + try: + self._read_info_files(force=True) # sweep emptied jsonl pairs + except Exception as e: + # reclaim is opportunistic: never mask the body's exception + logger.warning("Failed to sweep %s: %s", self.folder, e) @contextlib.contextmanager def writer(self) -> tp.Iterator["CacheDict[X]"]: diff --git a/exca/map.py b/exca/map.py index adb363d0..896d8e17 100644 --- a/exca/map.py +++ b/exca/map.py @@ -270,12 +270,9 @@ def _find_missing(self, items: dict[str, tp.Any]) -> dict[str, tp.Any]: self._check_configs(write=True) if self.mode == "force": to_remove = set(items) - set(missing) - state.recomputed - if to_remove: + if to_remove and isinstance(cache, CacheDict): msg = "Clearing %s items for %s (infra.mode=%s)" logger.warning(msg, len(to_remove), self.uid(), self.mode) - assert isinstance(cache, CacheDict), ( - f"to_remove is only filled when caching (got {type(cache)})" - ) with cache.write(): for uid in to_remove: del cache[uid] diff --git a/exca/steps/backends.py b/exca/steps/backends.py index 6796618b..72daf8b6 100644 --- a/exca/steps/backends.py +++ b/exca/steps/backends.py @@ -367,10 +367,10 @@ def run_and_cache(self) -> None: "Clearing partial results after invalid _run_batch output: %s", self.paths.step_uid, ) - with self.cache_dict.write(), self.cache_dict.frozen_cache_folder(): - for uid in written_uids: - if uid in self.cache_dict: - del self.cache_dict[uid] + with self.cache_dict.write(), self.cache_dict.frozen_cache_folder(): + for uid in written_uids: + if uid in self.cache_dict: + del self.cache_dict[uid] if folder is not None: e.add_note(f" -> cache may be invalid: {folder}") raise From edd03fc78fe1a41115c0497eac6aea05ca3c4091 Mon Sep 17 00:00:00 2001 From: Jeremy Rapin Date: Thu, 17 Sep 2026 16:09:18 +0200 Subject: [PATCH 3/3] fix --- CHANGELOG.md | 1 + exca/cachedict/core.py | 16 ++++++++-------- exca/cachedict/dumpcontext.py | 7 ++++--- exca/cachedict/handlers.py | 2 +- exca/cachedict/test_cachedict.py | 14 +++++++++++--- exca/cachedict/test_dumpcontext.py | 2 ++ 6 files changed, 27 insertions(+), 15 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index bbc96266..f32c9fc7 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,7 @@ [breaking] - `CacheDict`: deletions require a `write()` context (like writes). [#326] +- `DumpContext.shared_file`: content suffixes must start with `.`. [#326] [other] diff --git a/exca/cachedict/core.py b/exca/cachedict/core.py index bfc72065..4544a1b9 100644 --- a/exca/cachedict/core.py +++ b/exca/cachedict/core.py @@ -139,7 +139,6 @@ def __init__( self._folder_modified = -1.0 self._jsonl_readers: dict[str, JsonlReader] = {} self._jsonl_reading_allowance = float("inf") - self._deleted_in_scope = False # DumpContext for this folder (load/delete; writes use per-thread _write_ctx) self._dumper: DumpContext | None = None if self.folder is not None: @@ -162,6 +161,8 @@ def __reduce__(self) -> tp.Any: def clear(self) -> None: self._ram_data.clear() self._key_info.clear() + self._jsonl_readers.clear() + self._folder_modified = -1.0 if self.folder is None or not self.folder.exists(): return # let's remove content but not the folder to keep same permissions @@ -310,7 +311,7 @@ def write(self) -> tp.Iterator["CacheDict[X]"]: raise RuntimeError("Cannot re-open an already open writer") if self.folder is not None: self._write_ctx = DumpContext(self.folder, permissions=self.permissions) - self._deleted_in_scope = False + self._local.deleted_in_scope = False try: if self._write_ctx is not None: with self._write_ctx: @@ -321,11 +322,10 @@ def write(self) -> tp.Iterator["CacheDict[X]"]: self._write_ctx = None if self.folder is not None: utils.best_effort_utime(self.folder) - if self._deleted_in_scope: + if self._local.deleted_in_scope: try: self._read_info_files(force=True) # sweep emptied jsonl pairs - except Exception as e: - # reclaim is opportunistic: never mask the body's exception + except Exception as e: # must not mask the body's exception logger.warning("Failed to sweep %s: %s", self.folder, e) @contextlib.contextmanager @@ -343,7 +343,7 @@ def __setitem__(self, key: str, value: X) -> None: if not isinstance(key, str): raise TypeError(f"Non-string keys are not allowed (got {key!r})") if self.folder is not None and self._write_ctx is None: - raise RuntimeError("Cannot write outside of a writer context") + raise RuntimeError("Cannot write outside of a write() context") if self._folder_modified <= 0: _ = self.keys() if key in self._ram_data or key in self._key_info: @@ -375,8 +375,8 @@ def __delitem__(self, key: str) -> None: del self._ram_data[key] return if self._write_ctx is None: - raise RuntimeError("Cannot delete outside of a writer context") - self._deleted_in_scope = True + raise RuntimeError("Cannot delete outside of a write() context") + self._local.deleted_in_scope = True if key not in self._key_info: _ = key in self # populate _key_info from disk self._ram_data.pop(key, None) diff --git a/exca/cachedict/dumpcontext.py b/exca/cachedict/dumpcontext.py index 65165bbf..38e4305c 100644 --- a/exca/cachedict/dumpcontext.py +++ b/exca/cachedict/dumpcontext.py @@ -232,14 +232,15 @@ def shared_file(self, suffix: str) -> tuple[tp.IO[bytes], str]: """Open a shared file for appending. Returns (handle, relative_name). Content files go under DATA_DIR/; info files (-info.jsonl) stay in the root folder. Reused across calls with the same suffix.""" - if "." not in suffix: - raise ValueError(f"suffix must contain '.', got {suffix!r}") + is_info = suffix == self.INFO_SUFFIX + if not is_info and not suffix.startswith("."): + msg = f"suffix must start with '.' to be reclaimable, got {suffix!r}" + raise ValueError(msg) if self._stack is None: raise RuntimeError("DumpContext must be used as a context manager for writes") if threading.get_native_id() != self._thread_id: raise RuntimeError("DumpContext must not be shared across threads") basename = f"{self._prefix}{suffix}" - is_info = suffix == self.INFO_SUFFIX name = basename if is_info else f"{self.DATA_DIR}/{basename}" if name not in self._files: path = self.folder / name diff --git a/exca/cachedict/handlers.py b/exca/cachedict/handlers.py index 4020972b..5dbb1fe0 100644 --- a/exca/cachedict/handlers.py +++ b/exca/cachedict/handlers.py @@ -472,7 +472,7 @@ def __dump_info__(cls, ctx: DumpContext, value: tp.Any) -> dict[str, tp.Any]: ) from e if len(raw) <= cls.MAX_INLINE_SIZE: return {"content": value} - f, name = ctx.shared_file("-data.jsonl") + f, name = ctx.shared_file(".data.jsonl") offset = f.tell() f.write(raw + b"\n") return {"filename": name, "offset": offset, "length": len(raw)} diff --git a/exca/cachedict/test_cachedict.py b/exca/cachedict/test_cachedict.py index 85f07414..15c8324e 100644 --- a/exca/cachedict/test_cachedict.py +++ b/exca/cachedict/test_cachedict.py @@ -60,7 +60,7 @@ def test_array_cache(tmp_path: Path, in_ram: bool) -> None: cache2 = cd.CacheDict(folder=folder) assert isinstance(cache2["blublu"], np.ndarray) # del - with pytest.raises(RuntimeError, match="writer context"): + with pytest.raises(RuntimeError, match=r"write\(\) context"): del cache2["blublu"] with cache2.write(): del cache2["blublu"] @@ -309,12 +309,16 @@ def test_clone_is_view_only(tmp_path: Path) -> None: @pytest.mark.parametrize("read_before_delete", [False, True]) -@pytest.mark.parametrize("cache_type", ["MemmapArrayFile", "String"]) +@pytest.mark.parametrize("cache_type", ["MemmapArrayFile", "String", "Json"]) def test_orphaned_data_file_cleanup( tmp_path: Path, cache_type: str, read_before_delete: bool ) -> None: """Test that orphaned data files are cleaned up when all items are deleted.""" - data: tp.Any = np.random.rand(3, 12) if cache_type == "MemmapArrayFile" else "hello" + data: tp.Any = { + "MemmapArrayFile": np.random.rand(3, 12), + "String": "hello", + "Json": {"blob": "x" * 50_000}, # above MAX_INLINE_SIZE -> shared data file + }[cache_type] cache: cd.CacheDict[tp.Any] = cd.CacheDict( folder=tmp_path, keep_in_ram=False, cache_type=cache_type ) @@ -332,6 +336,10 @@ def test_orphaned_data_file_cleanup( assert len(remaining) == 2, ( f"leaving write() should drop the emptied pair {remaining}" ) + live = {p.name.removesuffix("-info.jsonl") for p in remaining} + data_files = (tmp_path / "data").glob("*") + stale = [p.name for p in data_files if p.name.split(".")[0] not in live] + assert not stale, f"data files outliving their info file {stale}" assert set(cache.keys()) == {"b1", "c2"} diff --git a/exca/cachedict/test_dumpcontext.py b/exca/cachedict/test_dumpcontext.py index ad9642a9..1fe63b34 100644 --- a/exca/cachedict/test_dumpcontext.py +++ b/exca/cachedict/test_dumpcontext.py @@ -110,6 +110,8 @@ def test_shared_file_lifecycle(tmp_path: Path) -> None: ctx = DumpContext(tmp_path) with pytest.raises(RuntimeError, match="context manager"): ctx.shared_file(".data") + with pytest.raises(ValueError, match="must start with"): + ctx.shared_file("-data.jsonl") with ctx: f1, name1 = ctx.shared_file(".data") f2, name2 = ctx.shared_file(".data")