From d7afdbf9a314e72980f64bc159434cbeb4f7cc41 Mon Sep 17 00:00:00 2001 From: Jun <84921700+Amazingjun-j@users.noreply.github.com> Date: Sun, 6 Sep 2026 16:39:52 +0800 Subject: [PATCH] fix(workspace-changes): drain snapshot scan before cancellation cleanup (#5232) * fix(workspace-changes): drain cancelled snapshot scans * test(workspace-changes): cover cancellation during snapshot scan * fix(workspace-changes): consume drained scan outcome * test(workspace-changes): keep scan cancellation regression focused * test(workspace-changes): pin cleanup ownership under recancel --- .../deerflow/workspace_changes/recorder.py | 48 +++++- .../test_workspace_changes_recorder.py | 151 +++++++++++++++++- 2 files changed, 195 insertions(+), 4 deletions(-) diff --git a/backend/packages/harness/deerflow/workspace_changes/recorder.py b/backend/packages/harness/deerflow/workspace_changes/recorder.py index 892d63ff0..67ca4573d 100644 --- a/backend/packages/harness/deerflow/workspace_changes/recorder.py +++ b/backend/packages/harness/deerflow/workspace_changes/recorder.py @@ -75,6 +75,39 @@ async def _reclaim_prepare_and_cleanup(prepare: asyncio.Future[tuple[list[Worksp await _remove_text_cache_dir(orphaned) +async def _drain_scan_and_cleanup( + scan: asyncio.Future[WorkspaceSnapshot], + text_cache_dir: Path | None, +) -> None: + """Let a cancelled scan finish before removing the cache it may still use.""" + while not scan.done(): + try: + await asyncio.shield(scan) + except asyncio.CancelledError: + pass + except Exception: + break + + # Cancellation remains the caller-visible outcome, but consume any late scan + # failure so the drained task cannot emit an un-retrieved exception warning. + try: + scan.result() + except asyncio.CancelledError: + pass + except Exception: + pass + + if text_cache_dir is None: + return + + cleanup = asyncio.create_task(_remove_text_cache_dir(text_cache_dir)) + while not cleanup.done(): + try: + await asyncio.shield(cleanup) + except asyncio.CancelledError: + pass + + async def capture_workspace_snapshot( thread_id: str, *, @@ -104,8 +137,9 @@ async def capture_workspace_snapshot( except asyncio.CancelledError: pass raise - try: - return await asyncio.to_thread( + + scan = asyncio.ensure_future( + asyncio.to_thread( scan_workspace_roots, roots, limits=limits, @@ -113,6 +147,16 @@ async def capture_workspace_snapshot( text_cache_dir=text_cache_dir, extra_excluded_dir_names=extra_excluded_dir_names, ) + ) + try: + return await asyncio.shield(scan) + except asyncio.CancelledError: + # The scan runs in a worker thread and cannot be stopped by cancelling + # this coroutine. Keep the cache alive until that worker has finished, + # then remove it before propagating cancellation. Repeated cancellation + # must not abandon either phase. + await _drain_scan_and_cleanup(scan, text_cache_dir) + raise except Exception: if text_cache_dir is not None: await _remove_text_cache_dir(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 index cb0f6e2b5..b8bdc3808 100644 --- a/backend/tests/blocking_io/test_workspace_changes_recorder.py +++ b/backend/tests/blocking_io/test_workspace_changes_recorder.py @@ -16,8 +16,10 @@ 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. +path would orphan the dir. The prepare-handoff tests pin the shield+reclaim guard; +the scan-cancellation tests pin the symmetric requirement to keep the cache alive +until the already-running scan worker drains, then remove it before cancellation +propagates. 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. @@ -207,3 +209,148 @@ async def test_capture_workspace_snapshot_repeated_cancellation_leaks_no_text_ca 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}" + + +async def test_capture_workspace_snapshot_cancelled_scan_drains_before_cleanup(tmp_path: Path, monkeypatch) -> None: + """A cancelled scan keeps its text cache until its worker is finished.""" + 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() + + def _blocking_scan(*_args: Any, text_cache_dir: str | Path | None = None, **_kwargs: Any) -> WorkspaceSnapshot: + assert text_cache_dir is not None + cache_dir = Path(text_cache_dir) + assert cache_dir.exists() + entered.set() + release.wait(timeout=5) + assert cache_dir.exists(), "scan cache was removed while the worker was still running" + return WorkspaceSnapshot(files={}, truncated=False, text_cache_dir=str(cache_dir)) + + monkeypatch.setattr(recorder, "scan_workspace_roots", _blocking_scan) + + task = asyncio.create_task(recorder.capture_workspace_snapshot("t1", include_text=True)) + assert await asyncio.to_thread(entered.wait, 5), "scan worker did not start" + + task.cancel() + for _ in range(5): + await asyncio.sleep(0) + assert not task.done(), "cancelled capture abandoned the still-running scan worker" + + release.set() + with pytest.raises(asyncio.CancelledError): + await task + + leftovers = await asyncio.to_thread(lambda: sorted(cache_root.glob("deerflow-workspace-changes-*"))) + assert leftovers == [], f"cancelled scan leaked a text cache dir: {leftovers}" + + +async def test_capture_workspace_snapshot_repeated_cancel_during_scan_still_cleans_up( + tmp_path: Path, + monkeypatch, +) -> None: + """Repeated cancellation cannot abandon scan draining or cache cleanup.""" + 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() + + def _blocking_scan(*_args: Any, text_cache_dir: str | Path | None = None, **_kwargs: Any) -> WorkspaceSnapshot: + assert text_cache_dir is not None + cache_dir = Path(text_cache_dir) + entered.set() + release.wait(timeout=5) + assert cache_dir.exists(), "scan cache was removed before the worker drained" + return WorkspaceSnapshot(files={}, truncated=False, text_cache_dir=str(cache_dir)) + + monkeypatch.setattr(recorder, "scan_workspace_roots", _blocking_scan) + + task = asyncio.create_task(recorder.capture_workspace_snapshot("t1", include_text=True)) + assert await asyncio.to_thread(entered.wait, 5), "scan worker did not start" + + task.cancel() + for _ in range(5): + await asyncio.sleep(0) + task.cancel() + for _ in range(5): + await asyncio.sleep(0) + assert not task.done(), "second cancellation abandoned the still-running scan worker" + + release.set() + 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 scan leaked a text cache dir: {leftovers}" + + +async def test_capture_workspace_snapshot_repeated_cancel_during_cleanup_still_cleans_up( + tmp_path: Path, + monkeypatch, +) -> None: + """A second cancellation cannot abandon cleanup after the scan has drained.""" + 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)) + + scan_entered = threading.Event() + scan_release = threading.Event() + cleanup_entered = asyncio.Event() + cleanup_release = asyncio.Event() + real_remove_text_cache_dir = recorder._remove_text_cache_dir + + def _blocking_scan(*_args: Any, text_cache_dir: str | Path | None = None, **_kwargs: Any) -> WorkspaceSnapshot: + assert text_cache_dir is not None + cache_dir = Path(text_cache_dir) + scan_entered.set() + scan_release.wait(timeout=5) + assert cache_dir.exists(), "scan cache was removed before the worker drained" + return WorkspaceSnapshot(files={}, truncated=False, text_cache_dir=str(cache_dir)) + + async def _blocking_cleanup(text_cache_dir: str | Path) -> None: + cleanup_entered.set() + await cleanup_release.wait() + await real_remove_text_cache_dir(text_cache_dir) + + monkeypatch.setattr(recorder, "scan_workspace_roots", _blocking_scan) + monkeypatch.setattr(recorder, "_remove_text_cache_dir", _blocking_cleanup) + + task = asyncio.create_task(recorder.capture_workspace_snapshot("t1", include_text=True)) + assert await asyncio.to_thread(scan_entered.wait, 5), "scan worker did not start" + + task.cancel() + scan_release.set() + await asyncio.wait_for(cleanup_entered.wait(), timeout=5) + + task.cancel() + for _ in range(5): + await asyncio.sleep(0) + assert not task.done(), "second cancellation abandoned the in-progress cache cleanup" + parked = await asyncio.to_thread(lambda: sorted(cache_root.glob("deerflow-workspace-changes-*"))) + assert parked, "cache should remain until the owned cleanup task is released" + + cleanup_release.set() + with pytest.raises(asyncio.CancelledError): + await task + + leftovers = await asyncio.to_thread(lambda: sorted(cache_root.glob("deerflow-workspace-changes-*"))) + assert leftovers == [], f"repeated cancellation during cleanup leaked a text cache dir: {leftovers}"