From f1f4af99bbf1cffb83a80744bf6e4a44ee8f24be Mon Sep 17 00:00:00 2001 From: Nefelibata <124799179+MeiSiristhebest@users.noreply.github.com> Date: Sat, 22 Aug 2026 17:24:20 +0800 Subject: [PATCH] fix(tools): retain strong reference to deferred subagent cleanup tasks (#4928) * fix(tools): retain strong reference to deferred subagent cleanup tasks * chore(tests): organize imports and format test_task_tool_core_logic.py --- .../deerflow/tools/builtins/task_tool.py | 8 ++++- backend/tests/test_task_tool_core_logic.py | 30 +++++++++++++++++++ 2 files changed, 37 insertions(+), 1 deletion(-) diff --git a/backend/packages/harness/deerflow/tools/builtins/task_tool.py b/backend/packages/harness/deerflow/tools/builtins/task_tool.py index a7136d2d5..97788df40 100644 --- a/backend/packages/harness/deerflow/tools/builtins/task_tool.py +++ b/backend/packages/harness/deerflow/tools/builtins/task_tool.py @@ -84,10 +84,16 @@ def _log_cleanup_failure(cleanup_task: asyncio.Task[None], *, trace_id: str, exe logger.error(f"[trace={trace_id}] Deferred cleanup failed for execution {execution_id}: {exc}") -def _schedule_deferred_subagent_cleanup(execution_id: str, trace_id: str, max_polls: int) -> None: +_deferred_cleanup_tasks: set[asyncio.Task[None]] = set() + + +def _schedule_deferred_subagent_cleanup(execution_id: str, trace_id: str, max_polls: int) -> asyncio.Task[None]: logger.debug(f"[trace={trace_id}] Scheduling deferred cleanup for cancelled execution {execution_id}") cleanup_task = asyncio.create_task(_deferred_cleanup_subagent_task(execution_id, trace_id, max_polls)) + _deferred_cleanup_tasks.add(cleanup_task) + cleanup_task.add_done_callback(_deferred_cleanup_tasks.discard) cleanup_task.add_done_callback(lambda task: _log_cleanup_failure(task, trace_id=trace_id, execution_id=execution_id)) + return cleanup_task def _find_usage_recorder(runtime: Any) -> Any | None: diff --git a/backend/tests/test_task_tool_core_logic.py b/backend/tests/test_task_tool_core_logic.py index a220d1c11..df8370aba 100644 --- a/backend/tests/test_task_tool_core_logic.py +++ b/backend/tests/test_task_tool_core_logic.py @@ -1,8 +1,10 @@ """Core behavior tests for task tool orchestration.""" import asyncio +import gc import importlib import inspect +import weakref from enum import Enum from types import SimpleNamespace from unittest.mock import MagicMock @@ -1853,3 +1855,31 @@ def test_terminal_event_usage_none_when_no_records(monkeypatch): completed = [e for e in events if e["type"] == "task_completed"] assert len(completed) == 1 assert completed[0]["usage"] is None + + +@pytest.mark.asyncio +async def test_deferred_cleanup_task_retained_and_survives_gc(monkeypatch): + """Verify deferred cleanup task is retained in _deferred_cleanup_tasks and completes after GC.""" + cleaned = [] + orig_sleep = asyncio.sleep + + monkeypatch.setattr(task_tool_module, "SubagentStatus", FakeSubagentStatus) + monkeypatch.setattr(task_tool_module, "get_background_task_result", lambda _: _make_result(FakeSubagentStatus.COMPLETED, result="ok")) + monkeypatch.setattr(task_tool_module, "cleanup_background_task", cleaned.append) + monkeypatch.setattr(task_tool_module.asyncio, "sleep", lambda _: orig_sleep(0)) + + task = task_tool_module._schedule_deferred_subagent_cleanup("exec-gc", "trace-gc", 5) + assert task in task_tool_module._deferred_cleanup_tasks + weak_task = weakref.ref(task) + del task + gc.collect() + + assert weak_task() is not None and weak_task() in task_tool_module._deferred_cleanup_tasks + for _ in range(10): + if cleaned: + break + await orig_sleep(0.01) + await orig_sleep(0.01) + + assert cleaned == ["exec-gc"] + assert weak_task() not in task_tool_module._deferred_cleanup_tasks