diff --git a/src/cachier/cores/pickle.py b/src/cachier/cores/pickle.py index fda3b308..ba214edd 100644 --- a/src/cachier/cores/pickle.py +++ b/src/cachier/cores/pickle.py @@ -180,6 +180,9 @@ def _clear_all_cache_files(self) -> None: try: os.remove(fpath) break + except FileNotFoundError: + # A concurrent clear already removed this file. + break except PermissionError: if attempt < 2: time.sleep(0.1 * (attempt + 1)) diff --git a/tests/pickle_tests/test_pickle_core.py b/tests/pickle_tests/test_pickle_core.py index 451363da..04c0c486 100644 --- a/tests/pickle_tests/test_pickle_core.py +++ b/tests/pickle_tests/test_pickle_core.py @@ -1475,6 +1475,93 @@ def mock_func(): core._clear_all_cache_files() +def _clear_race_target(x): + return x + + +def _clear_race_unrelated(x): + return x + + +@pytest.mark.pickle +@pytest.mark.parametrize("shared_wrapper", [True, False], ids=["shared_wrapper", "separate_wrappers"]) +def test_concurrent_clear_cache_separate_files(tmp_path, monkeypatch, shared_wrapper): + """Test concurrent clear_cache() calls tolerate files already removed by another caller.""" + wrapper_a = _get_decorated_func(_clear_race_target, cache_dir=tmp_path, separate_files=True) + wrapper_b = ( + wrapper_a + if shared_wrapper + else _get_decorated_func(_clear_race_target, cache_dir=tmp_path, separate_files=True) + ) + unrelated = _get_decorated_func(_clear_race_unrelated, cache_dir=tmp_path, separate_files=True) + for i in range(3): + wrapper_a(i) + unrelated(0) + other_file = tmp_path / "other_file.txt" + other_file.write_text("other content") + + target_prefix = f".{__name__}.{_clear_race_target.__qualname__}_" + target_files = {p.name for p in tmp_path.iterdir() if p.name.startswith(target_prefix)} + assert len(target_files) == 3 + untouched = {p.name for p in tmp_path.iterdir()} - target_files + assert len(untouched) == 2 + + real_listdir = os.listdir + real_remove = os.remove + # Both callers must list the directory before either starts deleting. + listed = threading.Barrier(2, timeout=10) + claims_lock = threading.Lock() + removed: dict[str, threading.Event] = {} + remove_calls: list[str] = [] + + def listdir_after_both_listed(path): + entries = real_listdir(path) + if os.path.realpath(path) == os.path.realpath(tmp_path): + listed.wait() + return entries + + def remove_in_order(path): + # The first caller to reach a file deletes it; the second waits for + # that deletion to finish before attempting its own real os.remove(). + with claims_lock: + remove_calls.append(os.path.basename(path)) + first = path not in removed + if first: + removed[path] = threading.Event() + if first: + try: + real_remove(path) + finally: + removed[path].set() + else: + assert removed[path].wait(timeout=10) + real_remove(path) + + monkeypatch.setattr("cachier.cores.pickle.os.listdir", listdir_after_both_listed) + monkeypatch.setattr("cachier.cores.pickle.os.remove", remove_in_order) + + errors = [] + + def clear(wrapper): + try: + wrapper.clear_cache() + except Exception as exc: + errors.append(exc) + + threads = [threading.Thread(target=clear, args=(w,)) for w in (wrapper_a, wrapper_b)] + for t in threads: + t.start() + for t in threads: + t.join(timeout=20) + monkeypatch.undo() + + assert not any(t.is_alive() for t in threads) + assert errors == [] + # Both callers attempted to delete every matching file. + assert sorted(remove_calls) == sorted([*target_files, *target_files]) + assert {p.name for p in tmp_path.iterdir()} == untouched + + # Redis core static method tests @pytest.mark.parametrize( ("test_input", "expected"),