fix(workspace-changes): offload blocking filesystem IO in text-cache lifecycle (#4268)

* 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.
This commit is contained in:
Aari
2026-07-18 17:30:14 +08:00
committed by GitHub
parent 208df4931d
commit a0e1d82ef4
3 changed files with 276 additions and 9 deletions
+5 -2
View File
@@ -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.
@@ -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)
@@ -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}"