AgentEvalTool/backend/agenteval/intelligent_eval/task_queue.py
sinohqb b4f9c887f4
All checks were successful
CI / test (push) Successful in 4m1s
feat(intelligent-eval): task queue monitor (方案③可视化)
方案③的定时触发(scan loop 每 60s 入队 + 触发 worker)此前只有 Worker
消费端 API,无可查看的列表。新增:
- GET /api/intelligent-evals/tasks:任务明细(含评估名/状态)+ 状态分布统计
  (注册在 /{eval_id} 之前避免被捕获为 eval_id="tasks")
- 前端 TaskQueueMonitor 组件 + 智能评估页任务队列入口:5s 轮询
  (usePolling),状态卡 + 状态筛选 + 明细表
测试:+3(列表/筛选/不被 {eval_id} 遮蔽),892 passed,tsc 通过
2026-08-17 13:57:00 +08:00

338 lines
12 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""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 agentcron=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 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 消费端 APInext/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}