From ec46fe2915df23c79b7847533dffd2a141645d40 Mon Sep 17 00:00:00 2001 From: WhaleTech <304937387+ryo-whaletech@users.noreply.github.com> Date: Thu, 17 Sep 2026 19:37:00 +0900 Subject: [PATCH 1/3] Python: handle concurrent FileSystemAgentFileStore deletion --- .../agent_framework/_harness/_file_access.py | 25 ++++++++-- .../tests/core/test_harness_file_access.py | 49 +++++++++++++++++++ 2 files changed, 69 insertions(+), 5 deletions(-) diff --git a/python/packages/core/agent_framework/_harness/_file_access.py b/python/packages/core/agent_framework/_harness/_file_access.py index 7dc95779ebc..127bc294e7c 100644 --- a/python/packages/core/agent_framework/_harness/_file_access.py +++ b/python/packages/core/agent_framework/_harness/_file_access.py @@ -27,6 +27,7 @@ import logging import os import re +import threading from abc import ABC, abstractmethod from collections.abc import Awaitable, Mapping, MutableMapping from pathlib import Path @@ -1097,6 +1098,11 @@ class FileSystemAgentFileStore(AgentFileStore): hostile process that shares the root directory. """ + _DELETE_LOCK_STRIPE_COUNT: ClassVar[int] = 64 + _DELETE_LOCKS: ClassVar[tuple[threading.Lock, ...]] = tuple( + threading.Lock() for _ in range(_DELETE_LOCK_STRIPE_COUNT) + ) + def __init__(self, root_directory: str | os.PathLike[str]) -> None: """Initialize the file-system store. @@ -1284,11 +1290,20 @@ async def delete(self, path: str) -> bool: full_path = self._resolve_safe_path(path) return await asyncio.to_thread(self._delete_file_sync, full_path) - @staticmethod - def _delete_file_sync(full_path: Path) -> bool: - if not full_path.is_file(): - return False - full_path.unlink() + @classmethod + def _delete_lock(cls, full_path: Path) -> threading.Lock: + """Return the process-local deletion lock for a file.""" + return cls._DELETE_LOCKS[hash(full_path) % cls._DELETE_LOCK_STRIPE_COUNT] + + @classmethod + def _delete_file_sync(cls, full_path: Path) -> bool: + with cls._delete_lock(full_path): + if not full_path.is_file(): + return False + try: + full_path.unlink() + except FileNotFoundError: + return False return True async def list_children(self, directory: str = "") -> list[FileStoreEntry]: diff --git a/python/packages/core/tests/core/test_harness_file_access.py b/python/packages/core/tests/core/test_harness_file_access.py index da7b4255f40..2a939fc1349 100644 --- a/python/packages/core/tests/core/test_harness_file_access.py +++ b/python/packages/core/tests/core/test_harness_file_access.py @@ -7,9 +7,11 @@ import os import re import stat +import threading import time from pathlib import Path from types import SimpleNamespace +from typing import Any import pytest @@ -287,6 +289,53 @@ async def test_filesystem_store_round_trips_files(tmp_path: Path) -> None: assert await store.delete("nested/a.txt") is False +async def test_filesystem_store_concurrent_delete(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + """Concurrent deletion should report one deletion and one missing file.""" + store = FileSystemAgentFileStore(tmp_path) + other_store = FileSystemAgentFileStore(tmp_path) + await store.write("shared.txt", "content") + + file_path = (tmp_path / "shared.txt").resolve() + original_to_thread = asyncio.to_thread + original_unlink = Path.unlink + worker_barrier = threading.Barrier(2) + unlink_barrier = threading.Barrier(2) + unlink_guard = threading.Lock() + overlapping_delete_completed = False + + async def synchronized_to_thread(function: Any, /, *args: Any, **kwargs: Any) -> Any: + def synchronized_call() -> Any: + worker_barrier.wait(timeout=5) + return function(*args, **kwargs) + + return await original_to_thread(synchronized_call) + + def macos_style_unlink(path: Path, missing_ok: bool = False) -> None: + nonlocal overlapping_delete_completed + if path.resolve() != file_path: + original_unlink(path, missing_ok=missing_ok) + return + + try: + unlink_barrier.wait(timeout=0.5) + except threading.BrokenBarrierError: + original_unlink(path, missing_ok=missing_ok) + return + + # Model concurrent macOS unlinks, where both calls may report success even though only one removes the file. + with unlink_guard: + if not overlapping_delete_completed: + original_unlink(path, missing_ok=missing_ok) + overlapping_delete_completed = True + + monkeypatch.setattr(asyncio, "to_thread", synchronized_to_thread) + monkeypatch.setattr(Path, "unlink", macos_style_unlink) + + results = await asyncio.gather(store.delete("shared.txt"), other_store.delete("shared.txt")) + + assert sorted(results) == [False, True] + + async def test_filesystem_store_rejects_traversal_and_rooted_paths(tmp_path: Path) -> None: """The filesystem store should refuse paths that escape the configured root.""" store = FileSystemAgentFileStore(tmp_path) From 8a6eb8ebff539f7464640832e801a3591eb8c7f9 Mon Sep 17 00:00:00 2001 From: WhaleTech <304937387+ryo-whaletech@users.noreply.github.com> Date: Fri, 18 Sep 2026 21:12:56 +0900 Subject: [PATCH 2/3] Python: serialize FileSystemAgentFileStore deletion across path aliases --- .../agent_framework/_harness/_file_access.py | 12 +--- .../tests/core/test_harness_file_access.py | 65 ++++++++++++++++--- 2 files changed, 58 insertions(+), 19 deletions(-) diff --git a/python/packages/core/agent_framework/_harness/_file_access.py b/python/packages/core/agent_framework/_harness/_file_access.py index 127bc294e7c..9b9901fc55e 100644 --- a/python/packages/core/agent_framework/_harness/_file_access.py +++ b/python/packages/core/agent_framework/_harness/_file_access.py @@ -1098,10 +1098,7 @@ class FileSystemAgentFileStore(AgentFileStore): hostile process that shares the root directory. """ - _DELETE_LOCK_STRIPE_COUNT: ClassVar[int] = 64 - _DELETE_LOCKS: ClassVar[tuple[threading.Lock, ...]] = tuple( - threading.Lock() for _ in range(_DELETE_LOCK_STRIPE_COUNT) - ) + _DELETE_LOCK: ClassVar[threading.Lock] = threading.Lock() def __init__(self, root_directory: str | os.PathLike[str]) -> None: """Initialize the file-system store. @@ -1290,14 +1287,9 @@ async def delete(self, path: str) -> bool: full_path = self._resolve_safe_path(path) return await asyncio.to_thread(self._delete_file_sync, full_path) - @classmethod - def _delete_lock(cls, full_path: Path) -> threading.Lock: - """Return the process-local deletion lock for a file.""" - return cls._DELETE_LOCKS[hash(full_path) % cls._DELETE_LOCK_STRIPE_COUNT] - @classmethod def _delete_file_sync(cls, full_path: Path) -> bool: - with cls._delete_lock(full_path): + with cls._DELETE_LOCK: if not full_path.is_file(): return False try: diff --git a/python/packages/core/tests/core/test_harness_file_access.py b/python/packages/core/tests/core/test_harness_file_access.py index 2a939fc1349..030b80f151e 100644 --- a/python/packages/core/tests/core/test_harness_file_access.py +++ b/python/packages/core/tests/core/test_harness_file_access.py @@ -3,6 +3,7 @@ from __future__ import annotations import asyncio +import itertools import json import os import re @@ -289,13 +290,14 @@ async def test_filesystem_store_round_trips_files(tmp_path: Path) -> None: assert await store.delete("nested/a.txt") is False -async def test_filesystem_store_concurrent_delete(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: - """Concurrent deletion should report one deletion and one missing file.""" - store = FileSystemAgentFileStore(tmp_path) - other_store = FileSystemAgentFileStore(tmp_path) - await store.write("shared.txt", "content") - - file_path = (tmp_path / "shared.txt").resolve() +async def _run_deterministic_concurrent_deletes( + store: FileSystemAgentFileStore, + other_store: FileSystemAgentFileStore, + first_path: str, + second_path: str, + monkeypatch: pytest.MonkeyPatch, +) -> tuple[bool, bool]: + target_paths = {store._resolve_safe_path(first_path), other_store._resolve_safe_path(second_path)} original_to_thread = asyncio.to_thread original_unlink = Path.unlink worker_barrier = threading.Barrier(2) @@ -312,7 +314,7 @@ def synchronized_call() -> Any: def macos_style_unlink(path: Path, missing_ok: bool = False) -> None: nonlocal overlapping_delete_completed - if path.resolve() != file_path: + if path not in target_paths: original_unlink(path, missing_ok=missing_ok) return @@ -330,8 +332,53 @@ def macos_style_unlink(path: Path, missing_ok: bool = False) -> None: monkeypatch.setattr(asyncio, "to_thread", synchronized_to_thread) monkeypatch.setattr(Path, "unlink", macos_style_unlink) + return await asyncio.gather(store.delete(first_path), other_store.delete(second_path)) + + +async def test_filesystem_store_concurrent_delete(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + """Concurrent deletion should report one deletion and one missing file.""" + store = FileSystemAgentFileStore(tmp_path) + other_store = FileSystemAgentFileStore(tmp_path) + await store.write("shared.txt", "content") + + results = await _run_deterministic_concurrent_deletes(store, other_store, "shared.txt", "shared.txt", monkeypatch) + + assert sorted(results) == [False, True] + + +async def test_filesystem_store_concurrent_delete_case_aliases(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + """Concurrent deletion through case aliases should report one deletion.""" + store = FileSystemAgentFileStore(tmp_path) + other_store = FileSystemAgentFileStore(tmp_path) + original_name = "Shared.txt" + await store.write(original_name, "content") + + original_path = store._resolve_safe_path(original_name) + lower_path = tmp_path / original_name.lower() + if not await asyncio.to_thread(lower_path.exists) or not await asyncio.to_thread( + os.path.samefile, original_path, lower_path + ): + pytest.skip("filesystem does not treat ASCII case variants as aliases") + + # Choose a case alias with opposite hash parity so the regression + # deterministically exercises the former path-hash-keyed synchronization bug + # without depending on this interpreter's randomized hash values. + alias_name = next( + ( + candidate_name + for characters in itertools.product(*[ + (character.lower(), character.upper()) for character in original_name + ]) + if (candidate_name := "".join(characters)) != original_name + and hash(store._resolve_safe_path(candidate_name)) % 2 != hash(original_path) % 2 + ), + None, + ) + assert alias_name is not None + alias_path = other_store._resolve_safe_path(alias_name) + assert await asyncio.to_thread(os.path.samefile, original_path, alias_path) - results = await asyncio.gather(store.delete("shared.txt"), other_store.delete("shared.txt")) + results = await _run_deterministic_concurrent_deletes(store, other_store, original_name, alias_name, monkeypatch) assert sorted(results) == [False, True] From ca31bfefbd49e9e03a9b69f4d6fb5b086d3ea4d0 Mon Sep 17 00:00:00 2001 From: WhaleTech <304937387+ryo-whaletech@users.noreply.github.com> Date: Fri, 25 Sep 2026 10:51:08 +0900 Subject: [PATCH 3/3] Python: fix FileSystemAgentFileStore delete test on Windows Windows Path hashes case aliases the same way, so the native test cannot find an alias with opposite hash parity. Skip only when that precondition is unavailable and cover alias locking through delete() on every CI platform. Check that unlink-time FileNotFoundError returns False, PermissionError propagates, and deleting a directory keeps its existing result. --- .../agent_framework/_harness/_file_access.py | 1 + .../tests/core/test_harness_file_access.py | 90 ++++++++++++++++++- 2 files changed, 90 insertions(+), 1 deletion(-) diff --git a/python/packages/core/agent_framework/_harness/_file_access.py b/python/packages/core/agent_framework/_harness/_file_access.py index d8d2d48621e..b6e5caf098f 100644 --- a/python/packages/core/agent_framework/_harness/_file_access.py +++ b/python/packages/core/agent_framework/_harness/_file_access.py @@ -1219,6 +1219,7 @@ class FileSystemAgentFileStore(AgentFileStore): hostile process that shares the root directory. """ + # Case aliases can identify the same file while having different Path hashes. _DELETE_LOCK: ClassVar[threading.Lock] = threading.Lock() def __init__(self, root_directory: str | os.PathLike[str]) -> None: diff --git a/python/packages/core/tests/core/test_harness_file_access.py b/python/packages/core/tests/core/test_harness_file_access.py index 928505976d2..d3f60142e46 100644 --- a/python/packages/core/tests/core/test_harness_file_access.py +++ b/python/packages/core/tests/core/test_harness_file_access.py @@ -10,9 +10,10 @@ import stat import threading import time +from contextlib import suppress from pathlib import Path from types import SimpleNamespace -from typing import Any +from typing import Any, cast import pytest import regex @@ -327,6 +328,8 @@ async def test_filesystem_store_round_trips_files(tmp_path: Path) -> None: assert await store.delete("nested/a.txt") is True assert await store.delete("nested/a.txt") is False + assert await store.delete("nested") is False + assert (tmp_path / "nested").is_dir() async def _run_deterministic_concurrent_deletes( @@ -385,6 +388,89 @@ async def test_filesystem_store_concurrent_delete(tmp_path: Path, monkeypatch: p assert sorted(results) == [False, True] +async def test_filesystem_store_concurrent_delete_aliases_with_distinct_hashes( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + """Different hashes for aliases must not allow both deletes to report success.""" + exists = True + state_lock = threading.Lock() + worker_barrier = threading.Barrier(2) + probe_barrier = threading.Barrier(2) + + class AliasedPath: + def __init__(self, hash_value: int) -> None: + self.hash_value = hash_value + + def __hash__(self) -> int: + return self.hash_value + + def is_file(self) -> bool: + with state_lock: + present = exists + with suppress(threading.BrokenBarrierError): + probe_barrier.wait(timeout=1) + return present + + def unlink(self) -> None: + nonlocal exists + with state_lock: + exists = False + + store = FileSystemAgentFileStore(tmp_path) + other_store = FileSystemAgentFileStore(tmp_path) + monkeypatch.setattr(store, "_resolve_safe_path", lambda _path: cast(Path, AliasedPath(0))) + monkeypatch.setattr(other_store, "_resolve_safe_path", lambda _path: cast(Path, AliasedPath(1))) + + original_to_thread = asyncio.to_thread + + async def synchronized_to_thread(function: Any, /, *args: Any, **kwargs: Any) -> Any: + def synchronized_call() -> Any: + worker_barrier.wait(timeout=5) + return function(*args, **kwargs) + + return await original_to_thread(synchronized_call) + + monkeypatch.setattr(asyncio, "to_thread", synchronized_to_thread) + + results = await asyncio.gather( + store.delete("first.txt"), + other_store.delete("second.txt"), + ) + + assert sorted(results) == [False, True] + + +async def test_filesystem_store_delete_handles_only_missing_file_from_unlink( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + """A vanished file returns False, while an unrelated unlink failure propagates.""" + store = FileSystemAgentFileStore(tmp_path) + await store.write("gone.txt", "content") + gone_path = tmp_path / "gone.txt" + original_unlink = Path.unlink + + def disappear(path: Path, missing_ok: bool = False) -> None: + original_unlink(path, missing_ok=missing_ok) + raise FileNotFoundError(path) + + with monkeypatch.context() as patch: + patch.setattr(Path, "unlink", disappear) + assert await store.delete("gone.txt") is False + assert not gone_path.exists() + + await store.write("denied.txt", "content") + denied_path = tmp_path / "denied.txt" + + def deny_unlink(path: Path, missing_ok: bool = False) -> None: + raise PermissionError(f"Cannot unlink {path}") + + with monkeypatch.context() as patch: + patch.setattr(Path, "unlink", deny_unlink) + with pytest.raises(PermissionError, match="Cannot unlink"): + await store.delete("denied.txt") + assert denied_path.is_file() + + async def test_filesystem_store_concurrent_delete_case_aliases(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: """Concurrent deletion through case aliases should report one deletion.""" store = FileSystemAgentFileStore(tmp_path) @@ -413,6 +499,8 @@ async def test_filesystem_store_concurrent_delete_case_aliases(tmp_path: Path, m ), None, ) + if alias_name is None: + pytest.skip("no case alias has opposite Path hash parity on this platform") assert alias_name is not None alias_path = other_store._resolve_safe_path(alias_name) assert await asyncio.to_thread(os.path.samefile, original_path, alias_path)