From e67a6911324a9ef34cb70585a6a92e9b63ac49b2 Mon Sep 17 00:00:00 2001 From: redcomet168 Date: Sat, 18 Jul 2026 23:06:28 -0700 Subject: [PATCH] fix: per-job EventLoopThread + timeout kill for parallel tools All parallel jobs shared a single EventLoopThread singleton (THREAD_BACKGROUND), so any blocking I/O in one worker froze all concurrent jobs until timeout. Timed-out jobs were marked but never killed, leaking threads and FDs. Changes (helpers/parallel_tools.py only): - start_parallel_jobs: use per-job thread name instead of shared singleton - await_parallel_jobs: kill timed-out jobs with terminate_thread=True - cleanup_parallel_job: use terminate_thread=True for thread cleanup - _cancel_job: use terminate_thread=True for thread cleanup No changes to defer.py, agent.py, or initialization. Fixes #1485, #1011, #1248 --- helpers/parallel_tools.py | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/helpers/parallel_tools.py b/helpers/parallel_tools.py index 0b9c9f72a6..48a45f8316 100644 --- a/helpers/parallel_tools.py +++ b/helpers/parallel_tools.py @@ -216,7 +216,7 @@ async def start_parallel_jobs( try: job.state = "running" job.started_at = time.time() - task = DeferredTask(thread_name=THREAD_BACKGROUND) + task = DeferredTask(thread_name=f"parallel-{job.id}") job.deferred_task = task task.start_task(_run_parallel_job, context.id, job.id) except Exception as exc: @@ -252,6 +252,10 @@ async def await_parallel_jobs( if time.time() >= deadline: wait_timed_out_job_ids = {job.id for job in active} + for job in active: + if job.deferred_task: + job.deferred_task.kill(terminate_thread=True) + _finish_job(job, "timeout", error="Job exceeded timeout and was killed.") break await asyncio.sleep(POLL_INTERVAL_SECONDS) @@ -316,7 +320,7 @@ async def refresh_parallel_jobs(agent: "Agent") -> list[ParallelJob]: async def cleanup_parallel_job(agent: "Agent", job: ParallelJob) -> None: if job.deferred_task and job.deferred_task.is_alive(): - job.deferred_task.kill() + job.deferred_task.kill(terminate_thread=True) if job.kind == "tool": await _remove_context(job.worker_context_id) @@ -604,7 +608,7 @@ async def _cancel_job( message: str = "Parallel job was cancelled.", ) -> None: if job.deferred_task and job.deferred_task.is_alive(): - job.deferred_task.kill() + job.deferred_task.kill(terminate_thread=True) _finish_job(job, state, error=message)