All checks were successful
CI / test (push) Successful in 3m57s
任务队列出现'评估已 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
376 lines
13 KiB
Python
376 lines
13 KiB
Python
"""Task queue for intelligent evaluations (任务队列).
|
||
|
||
Platform scans executing evals every minute and enqueues tasks for
|
||
OpenClaw workers to pick up. Tasks are prioritized by:
|
||
1. Time slot due (时段到期)
|
||
2. Session deficit (欠账多)
|
||
3. Wait time (等待时间长)
|
||
"""
|
||
|
||
from datetime import timedelta
|
||
from typing import Optional
|
||
|
||
from sqlmodel import Session, func, select, update
|
||
|
||
from agenteval.intelligent_eval.models import IntelligentEvalStatus
|
||
from agenteval.storage.db import (
|
||
IntelligentEvalDB,
|
||
IntelligentEvalTaskQueueDB,
|
||
utc_now,
|
||
)
|
||
|
||
|
||
def _is_slot_due(slot: dict, current_offset: timedelta) -> bool:
|
||
"""Thin wrapper delegating to :func:`agenteval.intelligent_eval.domain.is_slot_due`."""
|
||
from agenteval.intelligent_eval.domain import is_slot_due as _impl
|
||
|
||
return _impl(slot, current_offset)
|
||
|
||
|
||
def _calculate_session_deficit(eval_db: IntelligentEvalDB, session: Session) -> int:
|
||
"""Thin wrapper delegating to :func:`agenteval.intelligent_eval.domain.calculate_session_deficit`."""
|
||
from agenteval.intelligent_eval.domain import calculate_session_deficit as _impl
|
||
|
||
return _impl(eval_db, session)
|
||
|
||
|
||
def _calculate_priority(eval_db: IntelligentEvalDB, session: Session) -> int:
|
||
"""Thin wrapper delegating to :func:`agenteval.intelligent_eval.domain.calculate_priority`."""
|
||
from agenteval.intelligent_eval.domain import calculate_priority as _impl
|
||
|
||
return _impl(eval_db, session)
|
||
|
||
|
||
def _get_attention_reason(eval_db: IntelligentEvalDB, session: Session) -> Optional[str]:
|
||
"""Thin wrapper delegating to :func:`agenteval.intelligent_eval.domain.get_attention_reason`."""
|
||
from agenteval.intelligent_eval.domain import get_attention_reason as _impl
|
||
|
||
return _impl(eval_db, session)
|
||
|
||
|
||
def _has_pending_task(eval_id: str, session: Session) -> bool:
|
||
"""Check if eval already has a pending task (去重)."""
|
||
existing = session.exec(
|
||
select(IntelligentEvalTaskQueueDB).where(
|
||
IntelligentEvalTaskQueueDB.eval_id == eval_id,
|
||
IntelligentEvalTaskQueueDB.status == "pending",
|
||
)
|
||
).first()
|
||
return existing is not None
|
||
|
||
|
||
def scan_and_enqueue_tasks(session: Session) -> int:
|
||
"""Scan all executing evals and enqueue tasks.
|
||
|
||
Returns:
|
||
Number of tasks enqueued
|
||
"""
|
||
# Get all executing evals
|
||
evals = session.exec(
|
||
select(IntelligentEvalDB).where(IntelligentEvalDB.status == IntelligentEvalStatus.EXECUTING.value)
|
||
).all()
|
||
|
||
enqueued = 0
|
||
|
||
for eval_db in evals:
|
||
# Check if eval needs attention
|
||
reason = _get_attention_reason(eval_db, session)
|
||
if reason is None:
|
||
continue
|
||
|
||
# Check if already has pending task (去重)
|
||
if _has_pending_task(eval_db.id, session):
|
||
continue
|
||
|
||
# Calculate priority
|
||
priority = _calculate_priority(eval_db, session)
|
||
|
||
# Create task
|
||
task = IntelligentEvalTaskQueueDB(
|
||
eval_id=eval_db.id,
|
||
status="pending",
|
||
priority=priority,
|
||
reason=reason,
|
||
created_at=utc_now(),
|
||
updated_at=utc_now(),
|
||
)
|
||
session.add(task)
|
||
enqueued += 1
|
||
|
||
session.commit()
|
||
return enqueued
|
||
|
||
|
||
def get_next_task(session: Session) -> Optional[IntelligentEvalTaskQueueDB]:
|
||
"""Get next pending task (highest priority).
|
||
|
||
Returns:
|
||
Task with lowest priority value (highest priority), or None
|
||
"""
|
||
task = session.exec(
|
||
select(IntelligentEvalTaskQueueDB)
|
||
.where(IntelligentEvalTaskQueueDB.status == "pending")
|
||
.order_by(IntelligentEvalTaskQueueDB.priority, IntelligentEvalTaskQueueDB.created_at)
|
||
.limit(1)
|
||
).first()
|
||
return task
|
||
|
||
|
||
def assign_task(task_id: str, cron_id: str, session: Session) -> bool:
|
||
"""Atomically assign a pending task to a cron (CAS on status).
|
||
|
||
P1 真问题修复(§6.1): use ``UPDATE ... WHERE status='pending'`` and decide
|
||
on ``rowcount`` so two concurrent workers cannot both claim the same task.
|
||
The previous read-check-write left a race because SQLite + two sessions
|
||
could each read ``status=pending`` and each commit.
|
||
"""
|
||
stmt = (
|
||
update(IntelligentEvalTaskQueueDB)
|
||
.where(IntelligentEvalTaskQueueDB.id == task_id)
|
||
.where(IntelligentEvalTaskQueueDB.status == "pending")
|
||
.values(
|
||
status="assigned",
|
||
assigned_cron_id=cron_id,
|
||
assigned_at=utc_now(),
|
||
updated_at=utc_now(),
|
||
)
|
||
)
|
||
result = session.exec(stmt)
|
||
session.commit()
|
||
return result.rowcount > 0
|
||
|
||
|
||
def complete_task(task_id: str, success: bool, error: Optional[str], session: Session) -> bool:
|
||
"""Atomically mark a task as completed/failed (CAS on status).
|
||
|
||
P1 真问题修复(§6.1 审计): guard with ``status='assigned'`` so a
|
||
double-complete from worker + stuck-handler leaves the DB in one state.
|
||
"""
|
||
terminal = "completed" if success else "failed"
|
||
stmt = (
|
||
update(IntelligentEvalTaskQueueDB)
|
||
.where(IntelligentEvalTaskQueueDB.id == task_id)
|
||
.where(IntelligentEvalTaskQueueDB.status == "assigned")
|
||
.values(
|
||
status=terminal,
|
||
completed_at=utc_now(),
|
||
error=error,
|
||
updated_at=utc_now(),
|
||
)
|
||
)
|
||
result = session.exec(stmt)
|
||
session.commit()
|
||
return result.rowcount > 0
|
||
|
||
|
||
def requeue_stuck_task(eval_id: str, cron_id: str, session: Session) -> bool:
|
||
"""Mark the cron-stuck task as failed and enqueue a retry task.
|
||
|
||
P1 deepening (S4): the stuck-task settlement logic that previously lived
|
||
inside ``cron_pool.handle_stuck_cron`` (with a runtime import) now lives
|
||
here as a first-class operation. ``cron_pool`` only calls this.
|
||
|
||
Returns:
|
||
True if a stuck task was found and requeued, False otherwise.
|
||
"""
|
||
from agenteval.storage.db import IntelligentEvalTaskQueueDB
|
||
|
||
task = session.exec(
|
||
select(IntelligentEvalTaskQueueDB).where(
|
||
IntelligentEvalTaskQueueDB.eval_id == eval_id,
|
||
IntelligentEvalTaskQueueDB.status == "assigned",
|
||
IntelligentEvalTaskQueueDB.assigned_cron_id == cron_id,
|
||
)
|
||
).first()
|
||
if task is None:
|
||
return False
|
||
|
||
complete_task(task.id, False, "Cron stuck", session)
|
||
|
||
new_task = IntelligentEvalTaskQueueDB(
|
||
eval_id=eval_id,
|
||
status="pending",
|
||
priority=1,
|
||
reason="cron_stuck_retry",
|
||
)
|
||
session.add(new_task)
|
||
return True
|
||
|
||
|
||
# ---------------------------------------------------------------------------
|
||
# P3 deepening (S2) — get_next_task with embedded eval info
|
||
# ---------------------------------------------------------------------------
|
||
|
||
|
||
def get_next_task_with_eval(session: Session) -> Optional[dict]:
|
||
"""Return the next pending task with its eval details, or None.
|
||
|
||
P3 deepening (S2): the eval-loading + dict-building that previously lived in
|
||
``web/routers/intelligent_evals.py::get_next_task`` now lives here.
|
||
"""
|
||
task = get_next_task(session)
|
||
if task is None:
|
||
return None
|
||
|
||
eval_db = session.get(IntelligentEvalDB, task.eval_id)
|
||
if eval_db is None:
|
||
return None
|
||
|
||
return {
|
||
"task": {
|
||
"id": task.id,
|
||
"eval_id": task.eval_id,
|
||
"priority": task.priority,
|
||
"reason": task.reason,
|
||
"eval": {
|
||
"id": eval_db.id,
|
||
"name": eval_db.name,
|
||
"status": eval_db.status,
|
||
"plan": eval_db.get_plan(),
|
||
"started_at": eval_db.started_at.isoformat() if eval_db.started_at else None,
|
||
},
|
||
},
|
||
}
|
||
|
||
|
||
# ---------------------------------------------------------------------------
|
||
# 卡死恢复(方案③遗留):assigned 超时重新入队
|
||
# ---------------------------------------------------------------------------
|
||
|
||
STALE_ASSIGNED_MINUTES = 10
|
||
|
||
|
||
def requeue_stale_assigned_tasks(session: Session) -> int:
|
||
"""Requeue tasks that stayed `assigned` too long without completing.
|
||
|
||
方案③ worker 由平台触发 openclaw agent(cron=manual-run-...,非真实 cron),
|
||
fault_tolerance 的 stuck 检测(基于 cron pool busy + last_active_at)不适用。
|
||
若 agent 中断/失败,任务会永久卡在 `assigned`,scan 只查 pending 不会再入队。
|
||
|
||
这里把「assigned 超过 STALE_ASSIGNED_MINUTES 且对应评估仍 executing」的任务
|
||
重置为 pending(清空认领信息),平台 scan 循环随后会重新触发 worker 重试。
|
||
|
||
Returns:
|
||
重新入队的任务数。
|
||
"""
|
||
threshold = utc_now() - timedelta(minutes=STALE_ASSIGNED_MINUTES)
|
||
stale = session.exec(
|
||
select(IntelligentEvalTaskQueueDB).where(
|
||
IntelligentEvalTaskQueueDB.status == "assigned",
|
||
IntelligentEvalTaskQueueDB.assigned_at < threshold,
|
||
)
|
||
).all()
|
||
|
||
requeued = 0
|
||
for task in stale:
|
||
ev = session.get(IntelligentEvalDB, task.eval_id)
|
||
if ev is None or ev.status != IntelligentEvalStatus.EXECUTING.value:
|
||
continue
|
||
task.status = "pending"
|
||
task.assigned_cron_id = None
|
||
task.assigned_at = None
|
||
task.updated_at = utc_now()
|
||
requeued += 1
|
||
|
||
if requeued:
|
||
session.commit()
|
||
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
|
||
|
||
|
||
# ---------------------------------------------------------------------------
|
||
# 任务队列监控(方案③可视化):列表 + 状态分布
|
||
# ---------------------------------------------------------------------------
|
||
|
||
|
||
def list_tasks(
|
||
session: Session,
|
||
status: Optional[str] = None,
|
||
limit: int = 100,
|
||
) -> dict:
|
||
"""List task-queue entries with their eval names, newest first.
|
||
|
||
方案③的"定时触发"(scan loop 每分钟扫描入队 + 触发 worker)此前只有
|
||
Worker 消费端 API(next/assign/complete),没有可查看的列表。这里提供
|
||
任务明细 + 状态分布统计,供前端任务队列监控页展示。
|
||
|
||
Returns:
|
||
{"tasks": [...], "stats": {pending, assigned, completed, failed, unresolved}}
|
||
"""
|
||
stmt = select(IntelligentEvalTaskQueueDB).order_by(IntelligentEvalTaskQueueDB.created_at.desc())
|
||
if status:
|
||
stmt = stmt.where(IntelligentEvalTaskQueueDB.status == status)
|
||
tasks = session.exec(stmt.limit(max(1, min(limit, 500)))).all()
|
||
|
||
# 状态分布(全量统计,不受 limit 影响)
|
||
stats = {"pending": 0, "assigned": 0, "completed": 0, "failed": 0, "unresolved": 0}
|
||
rows = session.exec(
|
||
select(
|
||
IntelligentEvalTaskQueueDB.status,
|
||
func.count(IntelligentEvalTaskQueueDB.id),
|
||
).group_by(IntelligentEvalTaskQueueDB.status)
|
||
).all()
|
||
for status_val, cnt in rows:
|
||
if status_val in stats:
|
||
stats[status_val] = cnt
|
||
stats["unresolved"] = stats["pending"] + stats["assigned"]
|
||
|
||
result = []
|
||
for t in tasks:
|
||
ev = session.get(IntelligentEvalDB, t.eval_id)
|
||
result.append(
|
||
{
|
||
"id": t.id,
|
||
"eval_id": t.eval_id,
|
||
"eval_name": ev.name if ev else None,
|
||
"eval_status": ev.status if ev else None,
|
||
"status": t.status,
|
||
"priority": t.priority,
|
||
"reason": t.reason,
|
||
"assigned_cron_id": t.assigned_cron_id,
|
||
"assigned_at": t.assigned_at.isoformat() if t.assigned_at else None,
|
||
"completed_at": t.completed_at.isoformat() if t.completed_at else None,
|
||
"created_at": t.created_at.isoformat() if t.created_at else None,
|
||
"updated_at": t.updated_at.isoformat() if t.updated_at else None,
|
||
"error": t.error,
|
||
}
|
||
)
|
||
return {"tasks": result, "stats": stats}
|