From 71d38ebef6c16ae70088ea98ba3e2d8d889a103a Mon Sep 17 00:00:00 2001 From: sinohqb Date: Tue, 18 Aug 2026 16:04:46 +0800 Subject: [PATCH] fix(intelligent-eval): settle stale tasks of finished evals MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 任务队列出现'评估已 completed 却有待认领/执行中任务'的残留:评估离开 executing 后,其 pending/assigned 任务无人清理(requeue_stale_assigned_tasks 只处理 executing 评估的 assigned,评估结束被跳过)。 - 新增 task_queue.settle_tasks_for_finished_evals:评估非 executing 时, 其 pending/assigned 任务回收为 completed;executing 评估的任务保留 - scan loop 每 60s 在 requeue+scan 后调用清理 测试:+1(已结束评估的 pending/assigned 回收、executing 保留、幂等),901 passed --- .../agenteval/intelligent_eval/task_queue.py | 38 ++++++++++++++ backend/agenteval/web/app.py | 5 ++ .../unit/test_intelligent_eval_task_queue.py | 51 +++++++++++++++++++ 3 files changed, 94 insertions(+) diff --git a/backend/agenteval/intelligent_eval/task_queue.py b/backend/agenteval/intelligent_eval/task_queue.py index 63bef84..f3125a8 100644 --- a/backend/agenteval/intelligent_eval/task_queue.py +++ b/backend/agenteval/intelligent_eval/task_queue.py @@ -277,6 +277,44 @@ def requeue_stale_assigned_tasks(session: Session) -> int: return requeued +def settle_tasks_for_finished_evals(session: Session) -> int: + """Mark pending/assigned tasks of finished (non-executing) evals as completed. + + 评估离开 executing(completed/failed/cancelled)后,其待认领/执行中任务不再 + 需要执行,应回收为 completed,否则会永久残留(requeue_stale_assigned_tasks + 只处理 executing 评估的 assigned 任务,评估结束后被跳过 → 队列里出现 + "已完成评估却有待认领/执行中任务")。 + + Returns: + 清理的任务数。 + """ + from sqlmodel import select + + from agenteval.intelligent_eval.models import IntelligentEvalStatus + from agenteval.storage.db import IntelligentEvalTaskQueueDB + + tasks = session.exec( + select(IntelligentEvalTaskQueueDB).where( + IntelligentEvalTaskQueueDB.status.in_(["pending", "assigned"]) + ) + ).all() + + settled = 0 + for task in tasks: + ev = session.get(IntelligentEvalDB, task.eval_id) + if ev is None or ev.status == IntelligentEvalStatus.EXECUTING.value: + continue # executing 评估的任务正常流转,不清理 + task.status = "completed" + task.completed_at = utc_now() + task.error = "评估已结束,任务不再需要执行" + task.updated_at = utc_now() + settled += 1 + + if settled: + session.commit() + return settled + + # --------------------------------------------------------------------------- # 任务队列监控(方案③可视化):列表 + 状态分布 # --------------------------------------------------------------------------- diff --git a/backend/agenteval/web/app.py b/backend/agenteval/web/app.py index 1d42c4e..2ec1c10 100644 --- a/backend/agenteval/web/app.py +++ b/backend/agenteval/web/app.py @@ -297,6 +297,7 @@ async def _intelligent_eval_scan_loop() -> None: from agenteval.intelligent_eval.task_queue import ( requeue_stale_assigned_tasks, scan_and_enqueue_tasks, + settle_tasks_for_finished_evals, ) r = requeue_stale_assigned_tasks(session) @@ -305,6 +306,10 @@ async def _intelligent_eval_scan_loop() -> None: n = scan_and_enqueue_tasks(session) if n: _logger.info("智能评估扫描:入队 %d 个 Worker 任务", n) + # 清理:评估已结束(非 executing)的 pending/assigned 任务回收为 completed + settled = settle_tasks_for_finished_evals(session) + if settled: + _logger.info("任务回收:%d 个已结束评估的任务标记完成", settled) # 审计兜底:agent 未上报决策日志时,平台按评估状态补录 added = _supplement_decision_logs(session) if added: diff --git a/tests/unit/test_intelligent_eval_task_queue.py b/tests/unit/test_intelligent_eval_task_queue.py index 363ea44..f5f7d4f 100644 --- a/tests/unit/test_intelligent_eval_task_queue.py +++ b/tests/unit/test_intelligent_eval_task_queue.py @@ -292,3 +292,54 @@ def test_complete_task(db_session: Session): db_session.refresh(task2) assert task2.status == "failed" assert task2.error == "test error" + + +def test_settle_tasks_for_finished_evals(db_session: Session): + """评估已结束(completed/failed/cancelled)时,其 pending/assigned 任务应回收为 completed。 + + 回归:t480 队列出现"评估已 completed 却有待认领/执行中任务"的残留。 + """ + # 已结束的评估(completed):pending + assigned 任务都应被清理 + done_eval = IntelligentEvalDB( + name="done", target_id="t1", + status=IntelligentEvalStatus.COMPLETED.value, started_at=utc_now(), + ) + db_session.add(done_eval) + db_session.commit() + + pending_task = IntelligentEvalTaskQueueDB( + eval_id=done_eval.id, status="pending", priority=1, reason="slot_due", + ) + assigned_task = IntelligentEvalTaskQueueDB( + eval_id=done_eval.id, status="assigned", priority=1, reason="slot_due", + assigned_cron_id="cron1", + ) + db_session.add_all([pending_task, assigned_task]) + db_session.commit() + + # 仍在执行中的评估:任务应保留(正常流转) + running_eval = IntelligentEvalDB( + name="running", target_id="t1", + status=IntelligentEvalStatus.EXECUTING.value, started_at=utc_now(), + ) + db_session.add(running_eval) + db_session.commit() + live_task = IntelligentEvalTaskQueueDB( + eval_id=running_eval.id, status="assigned", priority=1, reason="slot_due", + assigned_cron_id="cron2", + ) + db_session.add(live_task) + db_session.commit() + + settled = task_queue.settle_tasks_for_finished_evals(db_session) + assert settled == 2 # 只清理已结束评估的 2 个任务 + + db_session.refresh(pending_task) + db_session.refresh(assigned_task) + db_session.refresh(live_task) + assert pending_task.status == "completed" + assert assigned_task.status == "completed" + assert live_task.status == "assigned" # executing 评估的任务保留 + + # 幂等:再跑一次不再清理 + assert task_queue.settle_tasks_for_finished_evals(db_session) == 0