mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-09-09 21:49:37 +00:00
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
This commit is contained in:
parent
9b2c03b429
commit
d7afdbf9a3
@ -75,6 +75,39 @@ async def _reclaim_prepare_and_cleanup(prepare: asyncio.Future[tuple[list[Worksp
|
|||||||
await _remove_text_cache_dir(orphaned)
|
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(
|
async def capture_workspace_snapshot(
|
||||||
thread_id: str,
|
thread_id: str,
|
||||||
*,
|
*,
|
||||||
@ -104,8 +137,9 @@ async def capture_workspace_snapshot(
|
|||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
pass
|
pass
|
||||||
raise
|
raise
|
||||||
try:
|
|
||||||
return await asyncio.to_thread(
|
scan = asyncio.ensure_future(
|
||||||
|
asyncio.to_thread(
|
||||||
scan_workspace_roots,
|
scan_workspace_roots,
|
||||||
roots,
|
roots,
|
||||||
limits=limits,
|
limits=limits,
|
||||||
@ -113,6 +147,16 @@ async def capture_workspace_snapshot(
|
|||||||
text_cache_dir=text_cache_dir,
|
text_cache_dir=text_cache_dir,
|
||||||
extra_excluded_dir_names=extra_excluded_dir_names,
|
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:
|
except Exception:
|
||||||
if text_cache_dir is not None:
|
if text_cache_dir is not None:
|
||||||
await _remove_text_cache_dir(text_cache_dir)
|
await _remove_text_cache_dir(text_cache_dir)
|
||||||
|
|||||||
@ -16,8 +16,10 @@ rmtree this anchor exists to guard.
|
|||||||
|
|
||||||
Because ``mkdtemp`` must be offloaded, its worker handoff is also a cancellation
|
Because ``mkdtemp`` must be offloaded, its worker handoff is also a cancellation
|
||||||
hazard: a run cancelled after ``mkdtemp`` but before the coroutine receives the
|
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
|
path would orphan the dir. The prepare-handoff tests pin the shield+reclaim guard;
|
||||||
removes such a dir instead of leaking it.
|
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
|
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.
|
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-*")))
|
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}"
|
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}"
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user