From a0e1d82ef45cee1c4b0ccc9e2cab891b522dac03 Mon Sep 17 00:00:00 2001 From: Aari Date: Sat, 18 Jul 2026 17:30:14 +0800 Subject: [PATCH] fix(workspace-changes): offload blocking filesystem IO in text-cache lifecycle (#4268) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(workspace-changes): offload blocking filesystem IO in text-cache lifecycle capture_workspace_snapshot and record_workspace_changes offload their scans via asyncio.to_thread, but ran the snapshot text cache's whole lifecycle on the event loop: roots resolution (os.path.abspath), tempfile.mkdtemp, and shutil.rmtree on both the capture-failure branch and record_workspace_changes' finally. That finally runs on every agent run, including abort paths, so each run removed up to max_files cached texts on the loop. Offload the roots + mkdtemp prep through one _prepare_capture worker hop, and route both rmtree call sites through _remove_text_cache_dir. Cleanup stays best-effort: it swallows and logs, so a failing cleanup cannot replace the exception or result already in flight. asyncio.shield is deliberately not used -- to_thread submits to the pool immediately, so cancelling the future does not stop the running thread and the cache is still removed under single and repeated cancellation. Externally observable behavior is unchanged. Found via `make detect-blocking-io`: these were the last 2 HIGH findings in the repo, which now reports none. The roots resolution is invisible to that scanner (sync helper, cross-file call) but blocks the same async path, and the anchor cannot reach the rmtree without it. Add tests/blocking_io/test_workspace_changes_recorder.py, driving the capture-failure branch and the record finally. Teeth verified per clause under the strict Blockbuster gate: reverting each offload alone reddens its own call (os.path.samestat for rmtree, os.path.abspath for roots/mkdtemp). * fix(workspace-changes): make text-cache prepare handoff cancellation-safe _prepare_capture creates the mkdtemp cache dir inside the to_thread worker, so a run cancelled after mkdtemp but before the coroutine receives the path orphaned the dir. Shield the prepare future and, on cancellation, reclaim its result to remove the dir before re-raising. mkdtemp stays offloaded (the blocking-io gate flags os.mkdir from deerflow code). Adds a deterministic cancellation regression. * fix(workspace-changes): drain repeated cancellation in text-cache reclaim The mkdtemp handoff guard reclaimed the shielded worker's result on the first CancelledError, but the reclaim await was itself cancellable: a second cancel landed there, slipped past `except Exception` (CancelledError is BaseException), and skipped the reclaim while the shielded worker still finished — orphaning the deerflow-workspace-changes-* dir. Repeated cancelled runs accumulate leaks. Move reclaim+remove into a task the caller cannot abandon and drain repeated cancellation until it completes, then restore the cancellation. A repeat cancel interrupts the await, not the task, so the dir is never abandoned; the loop exits only once cleanup has finished, leaving no pending task. Non-cancel paths are unchanged. Adds test_capture_workspace_snapshot_repeated_cancellation_leaks_no_text_cache (double-cancel regression). make test-blocking-io: 43 passed. --- backend/AGENTS.md | 7 +- .../deerflow/workspace_changes/recorder.py | 69 +++++- .../test_workspace_changes_recorder.py | 209 ++++++++++++++++++ 3 files changed, 276 insertions(+), 9 deletions(-) create mode 100644 backend/tests/blocking_io/test_workspace_changes_recorder.py diff --git a/backend/AGENTS.md b/backend/AGENTS.md index c83c652f3..079c5ac63 100644 --- a/backend/AGENTS.md +++ b/backend/AGENTS.md @@ -142,10 +142,13 @@ Blocking-IO runtime gate (`tests/blocking_io/`): `test_jsonl_run_event_store.py` (locks `JsonlRunEventStore`'s async API offloading its file IO via `asyncio.to_thread`); `test_uploads_middleware.py` (locks `UploadsMiddleware.abefore_agent` - offloading the uploads-directory scan off the event loop); and + offloading the uploads-directory scan off the event loop); `test_uploads_router.py` (locks Gateway upload/list/delete endpoints offloading upload directory creation, staged writes, chmod/cleanup, - directory scans/deletes, and remote sandbox sync off the event loop). + directory scans/deletes, and remote sandbox sync off the event loop); and + `test_workspace_changes_recorder.py` (locks the offload around the snapshot + text cache lifecycle — roots resolution, `mkdtemp`, and the `shutil.rmtree` + on both the capture-failure branch and `record_workspace_changes`' `finally`). - `test_gate_smoke.py` is a meta-test asserting the gate actually catches unoffloaded blocking IO and that the `@pytest.mark.allow_blocking_io` opt-out works. diff --git a/backend/packages/harness/deerflow/workspace_changes/recorder.py b/backend/packages/harness/deerflow/workspace_changes/recorder.py index b1fc893b3..a56186799 100644 --- a/backend/packages/harness/deerflow/workspace_changes/recorder.py +++ b/backend/packages/harness/deerflow/workspace_changes/recorder.py @@ -38,6 +38,42 @@ def build_thread_workspace_roots(thread_id: str, *, user_id: str | None = None) ] +def _prepare_capture(thread_id: str, *, user_id: str | None, include_text: bool) -> tuple[list[WorkspaceRoot], Path | None]: + # Worker thread: resolving the sandbox roots hits the filesystem, and mkdtemp + # creates the text cache directory — both blocking IO that must stay off the + # event loop. + roots = build_thread_workspace_roots(thread_id, user_id=user_id) + text_cache_dir = Path(tempfile.mkdtemp(prefix="deerflow-workspace-changes-")) if include_text else None + return roots, text_cache_dir + + +async def _remove_text_cache_dir(text_cache_dir: str | Path) -> None: + """Remove a snapshot's text cache off the event loop. + + Best-effort by contract: every caller is a failure or teardown path, so a + cleanup error must never replace the exception or result already in flight. + """ + try: + await asyncio.to_thread(shutil.rmtree, text_cache_dir, ignore_errors=True) + except Exception: + logger.warning("Failed to remove workspace text cache %s", text_cache_dir, exc_info=True) + + +async def _reclaim_prepare_and_cleanup(prepare: asyncio.Future[tuple[list[WorkspaceRoot], Path | None]]) -> None: + """Await a cancelled prepare handoff and remove any dir it created. + + Owned by its own task so that caller cancellation during reclaim can interrupt + the *await* but never abandon a just-created text cache dir. Best-effort, + mirroring `_remove_text_cache_dir`. + """ + try: + _, orphaned = await prepare + except Exception: + return # prepare failed before creating a dir; nothing to reclaim + if orphaned is not None: + await _remove_text_cache_dir(orphaned) + + async def capture_workspace_snapshot( thread_id: str, *, @@ -45,8 +81,27 @@ async def capture_workspace_snapshot( limits: WorkspaceChangeLimits | None = None, include_text: bool = True, ) -> WorkspaceSnapshot: - roots = build_thread_workspace_roots(thread_id, user_id=user_id) - text_cache_dir = Path(tempfile.mkdtemp(prefix="deerflow-workspace-changes-")) if include_text else None + # `_prepare_capture` creates the text cache dir inside the worker, so the + # handoff must be cancellation-safe: if the run is cancelled after mkdtemp + # but before we receive the path, the shielded worker still finishes and we + # reclaim its result to remove the orphaned dir before re-raising. + prepare = asyncio.ensure_future(asyncio.to_thread(_prepare_capture, thread_id, user_id=user_id, include_text=include_text)) + try: + roots, text_cache_dir = await asyncio.shield(prepare) + except asyncio.CancelledError: + # `prepare` is shielded, so it keeps running and may still create the dir + # after this cancel. Own the reclaim+remove in a task the caller cannot + # abandon: a repeat cancel can interrupt our await but not the task, so we + # drain repeated cancellation until cleanup finishes, then restore it. A + # second `shield()` on the await alone would let a re-cancel skip the + # reclaim and orphan the dir. + cleanup = asyncio.ensure_future(_reclaim_prepare_and_cleanup(prepare)) + while not cleanup.done(): + try: + await asyncio.shield(cleanup) + except asyncio.CancelledError: + pass + raise try: return await asyncio.to_thread( scan_workspace_roots, @@ -57,7 +112,7 @@ async def capture_workspace_snapshot( ) except Exception: if text_cache_dir is not None: - shutil.rmtree(text_cache_dir, ignore_errors=True) + await _remove_text_cache_dir(text_cache_dir) raise @@ -71,7 +126,7 @@ async def record_workspace_changes( limits: WorkspaceChangeLimits | None = None, ) -> dict | None: try: - roots = build_thread_workspace_roots(thread_id, user_id=user_id) + roots = await asyncio.to_thread(build_thread_workspace_roots, thread_id, user_id=user_id) after_metadata = await asyncio.to_thread( scan_workspace_roots, roots, @@ -103,9 +158,9 @@ async def record_workspace_changes( metadata={WORKSPACE_CHANGES_METADATA_KEY: payload}, ) finally: - _cleanup_snapshot_text_cache(before) + await _cleanup_snapshot_text_cache(before) -def _cleanup_snapshot_text_cache(snapshot: WorkspaceSnapshot) -> None: +async def _cleanup_snapshot_text_cache(snapshot: WorkspaceSnapshot) -> None: if snapshot.text_cache_dir: - shutil.rmtree(snapshot.text_cache_dir, ignore_errors=True) + await _remove_text_cache_dir(snapshot.text_cache_dir) diff --git a/backend/tests/blocking_io/test_workspace_changes_recorder.py b/backend/tests/blocking_io/test_workspace_changes_recorder.py new file mode 100644 index 000000000..cb0f6e2b5 --- /dev/null +++ b/backend/tests/blocking_io/test_workspace_changes_recorder.py @@ -0,0 +1,209 @@ +"""Regression anchor: workspace snapshot text-cache cleanup must not block the event loop. + +``capture_workspace_snapshot`` offloads the scan itself via ``asyncio.to_thread`` +but owns a ``tempfile.mkdtemp`` text cache whose lifecycle runs on the async +path: the directory is created up front, removed on the scan-failure branch, and +removed again by ``record_workspace_changes``' ``finally`` after every run. Those +create/delete calls are blocking filesystem IO (``shutil.rmtree`` walks and +unlinks up to ``max_files`` cached texts). If any of them regresses back onto the +event loop, the strict Blockbuster gate raises ``BlockingError`` and these tests +fail. + +Both cleanup branches are driven explicitly — the failure branch of +``capture_workspace_snapshot`` and the always-run ``finally`` of +``record_workspace_changes`` — because the happy path alone never reaches the +rmtree this anchor exists to guard. + +Because ``mkdtemp`` must be offloaded, its worker handoff is also a cancellation +hazard: a run cancelled after ``mkdtemp`` but before the coroutine receives the +path would orphan the dir. The last test pins the shield+reclaim guard that +removes such a dir instead of leaking it. + +Imports are kept at module top so any import-time IO runs at collection (outside +the gate); the surface under test runs on the event loop inside the gated test. +""" + +from __future__ import annotations + +import asyncio +import tempfile +import threading +from pathlib import Path +from typing import Any + +import pytest + +from deerflow.workspace_changes import recorder +from deerflow.workspace_changes.types import WorkspaceSnapshot + +pytestmark = pytest.mark.asyncio + + +class _RecordingEventStore: + """Stand-in for the real event store: the external boundary, not the offload.""" + + def __init__(self) -> None: + self.puts: list[dict[str, Any]] = [] + + async def put(self, **kwargs: Any) -> dict[str, Any]: + self.puts.append(kwargs) + return kwargs + + +def _seed_workspace(tmp_path: Path) -> None: + """Create a real on-disk workspace so the scan has something to walk.""" + work = tmp_path / "work" + work.mkdir(parents=True, exist_ok=True) + (work / "note.txt").write_text("hello\n", encoding="utf-8") + + +async def test_capture_workspace_snapshot_cleanup_does_not_block_event_loop(tmp_path: Path, monkeypatch) -> None: + """The scan-failure branch removes the text cache; that rmtree must be offloaded.""" + monkeypatch.setenv("DEER_FLOW_HOME", str(tmp_path)) + import deerflow.config.paths as paths_mod + + monkeypatch.setattr(paths_mod, "_paths", None) + + # Pin mkdtemp's parent so the assertion below sees this test's cache dir and + # nothing else (the platform temp root is shared and macOS is not /tmp). + cache_root = tmp_path / "tmp" + cache_root.mkdir() + monkeypatch.setattr(tempfile, "tempdir", str(cache_root)) + + # Force the failure branch. This mocks the scan (a separate, already-offloaded + # call), never the text-cache cleanup this anchor guards. + def _boom(*args: Any, **kwargs: Any) -> None: + raise RuntimeError("scan failed") + + monkeypatch.setattr(recorder, "scan_workspace_roots", _boom) + + with pytest.raises(RuntimeError, match="scan failed"): + await recorder.capture_workspace_snapshot("t1", include_text=True) + + # The cache dir was really created, then really removed — cleanup still runs, + # it merely moved off the loop. + leftovers = await asyncio.to_thread(lambda: sorted(cache_root.glob("deerflow-workspace-changes-*"))) + assert leftovers == [], f"text cache dir leaked on the failure branch: {leftovers}" + + +async def test_record_workspace_changes_cleanup_does_not_block_event_loop(tmp_path: Path, monkeypatch) -> None: + """``record_workspace_changes`` rmtrees the snapshot text cache in its ``finally``.""" + monkeypatch.setenv("DEER_FLOW_HOME", str(tmp_path)) + import deerflow.config.paths as paths_mod + + monkeypatch.setattr(paths_mod, "_paths", None) + + _seed_workspace(tmp_path) + + # A real text cache dir holding real files, so rmtree does real filesystem work. + cache_dir = tmp_path / "text-cache" + cache_dir.mkdir() + for i in range(5): + (cache_dir / f"cached_{i}.txt").write_text("cached\n", encoding="utf-8") + + before = WorkspaceSnapshot(files={}, truncated=False, text_cache_dir=str(cache_dir)) + + await recorder.record_workspace_changes( + _RecordingEventStore(), + "t1", + "r1", + before, + ) + + still_there = await asyncio.to_thread(cache_dir.exists) + assert not still_there, "record_workspace_changes must remove the snapshot text cache" + + +async def test_capture_workspace_snapshot_cancelled_handoff_leaks_no_text_cache(tmp_path: Path, monkeypatch) -> None: + """A run cancelled during the mkdtemp handoff must not orphan the text cache. + + ``mkdtemp`` runs in the ``_prepare_capture`` worker, so if the run is + cancelled after the dir is created but before the coroutine receives the + path, nothing downstream owns it. The shield+reclaim guard waits for the + worker and removes the dir; without it the dir leaks into the temp root. + """ + monkeypatch.setenv("DEER_FLOW_HOME", str(tmp_path)) + import deerflow.config.paths as paths_mod + + monkeypatch.setattr(paths_mod, "_paths", None) + + cache_root = tmp_path / "tmp" + cache_root.mkdir() + monkeypatch.setattr(tempfile, "tempdir", str(cache_root)) + + entered = threading.Event() + release = threading.Event() + real_mkdtemp = tempfile.mkdtemp + + def _blocking_mkdtemp(*args: Any, **kwargs: Any) -> str: + created = real_mkdtemp(*args, **kwargs) # the dir really exists now + entered.set() + release.wait(timeout=5) # park the worker mid-handoff, holding the result + return created + + monkeypatch.setattr(recorder.tempfile, "mkdtemp", _blocking_mkdtemp) + + task = asyncio.ensure_future(recorder.capture_workspace_snapshot("t1", include_text=True)) + await asyncio.to_thread(entered.wait, 5) # mkdtemp created the dir; worker is parked + parked = await asyncio.to_thread(lambda: sorted(cache_root.glob("deerflow-workspace-changes-*"))) + assert parked, "text cache dir should exist while the worker is parked mid-handoff" + + task.cancel() + release.set() # let the worker finish and hand its result to the reclaim path + with pytest.raises(asyncio.CancelledError): + await task + + leftovers = await asyncio.to_thread(lambda: sorted(cache_root.glob("deerflow-workspace-changes-*"))) + assert leftovers == [], f"cancelled capture leaked a text cache dir: {leftovers}" + + +async def test_capture_workspace_snapshot_repeated_cancellation_leaks_no_text_cache(tmp_path: Path, monkeypatch) -> None: + """A *second* cancellation during the reclaim await must not orphan the cache. + + After the first cancel enters the reclaim path, the coroutine awaits the + shielded worker's result. A second cancel lands on that await: because the + reclaim+remove is owned by a task the caller cannot abandon, the guard drains + the repeated cancellation until the dir is removed, then restores the + cancellation. A plain re-await would let the second ``CancelledError`` skip + the reclaim (``except Exception`` does not catch it) while the shielded worker + still finishes and leaks its dir. + """ + monkeypatch.setenv("DEER_FLOW_HOME", str(tmp_path)) + import deerflow.config.paths as paths_mod + + monkeypatch.setattr(paths_mod, "_paths", None) + + cache_root = tmp_path / "tmp" + cache_root.mkdir() + monkeypatch.setattr(tempfile, "tempdir", str(cache_root)) + + entered = threading.Event() + release = threading.Event() + real_mkdtemp = tempfile.mkdtemp + + def _blocking_mkdtemp(*args: Any, **kwargs: Any) -> str: + created = real_mkdtemp(*args, **kwargs) # the dir really exists now + entered.set() + release.wait(timeout=5) # park the worker mid-handoff, holding the result + return created + + monkeypatch.setattr(recorder.tempfile, "mkdtemp", _blocking_mkdtemp) + + task = asyncio.ensure_future(recorder.capture_workspace_snapshot("t1", include_text=True)) + await asyncio.to_thread(entered.wait, 5) # mkdtemp created the dir; worker is parked + parked = await asyncio.to_thread(lambda: sorted(cache_root.glob("deerflow-workspace-changes-*"))) + assert parked, "text cache dir should exist while the worker is parked mid-handoff" + + task.cancel() # cancel #1 -> enters reclaim, awaits the shielded cleanup task + for _ in range(5): + await asyncio.sleep(0) # let the reclaim path reach its await while the worker is still parked + task.cancel() # cancel #2 -> lands on the reclaim await + for _ in range(5): + await asyncio.sleep(0) + + release.set() # worker finishes; the drained cleanup reclaims and removes the dir + with pytest.raises(asyncio.CancelledError): + await task + + leftovers = await asyncio.to_thread(lambda: sorted(cache_root.glob("deerflow-workspace-changes-*"))) + assert leftovers == [], f"repeated-cancel capture leaked a text cache dir: {leftovers}"