Skip to content
Open
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
3 changes: 3 additions & 0 deletions src/cachier/cores/pickle.py
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down
87 changes: 87 additions & 0 deletions tests/pickle_tests/test_pickle_core.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"),
Expand Down