AgentEvalTool/backend/agenteval/intelligent_eval/task_queue.py
sinohqb 182b0e59cb
All checks were successful
CI / test (push) Successful in 3m9s
refactor(intelligent-eval): extract scheduler runtime + internalize decision-log dedup
架构深化两则(架构审查候选①②):

① scheduler 抽取:web/app.py 约 400 行触发式执行编排(60s 扫描循环、
docker exec 触发、失败落账)沉入 intelligent_eval/scheduler.py,runtime
单例 start()/stop()/scan_once() 与 campaign_runtime 惯例一致;worker/planner
两处重复触发代码合并为一个触发原语;_supplement_decision_logs 归入
decision_logs.py。测试改为直接驱动 scan_once(interface 即测试面)。

② 决策日志去重内化:create_decision_log 的去重只服务 agent 上报路径;
新增 append_decision_log(平台落账纯追加)与 count_decisions(计数原语),
lifecycle/task_queue 全部平台落账切换,调用方不再需要塞 attempt 骗去重。

零行为变化:提示词、60s 节拍、编排顺序、闸门语义原样保留,866 tests passed。
2026-08-21 03:20:14 +08:00

375 lines
13 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/assigned task (去重)."""
existing = session.exec(
select(IntelligentEvalTaskQueueDB).where(
IntelligentEvalTaskQueueDB.eval_id == eval_id,
IntelligentEvalTaskQueueDB.status.in_(["pending", "assigned"]),
)
).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
# ---------------------------------------------------------------------------
# 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
# ADR-0011卡死任务的重试预算。超出后任务置 failed而非无限重入队
# 评估级收尾由 executing 兜底 watchdog 负责。
MAX_TASK_ATTEMPTS = 3
def requeue_stale_assigned_tasks(session: Session) -> int:
"""Requeue tasks that stayed `assigned` too long without completing.
方案③ worker 由平台触发 openclaw agentcron=manual-run-...,非真实 cron
若 agent 中断/失败,任务会永久卡在 `assigned`scan 只查 pending 不会再入队。
这里把「assigned 超过 STALE_ASSIGNED_MINUTES 且对应评估仍 executing」的任务
重置为 pending清空认领信息平台 scan 循环随后会重新触发 worker 重试。
ADR-0011每次 requeue 累计 attempts达到 MAX_TASK_ATTEMPTS 后改置 failed
并补录决策日志——确定性失败的任务不再无限烧触发。
Returns:
重新入队的任务数(不含被判 failed 的)。
"""
from agenteval.intelligent_eval.decision_logs import append_decision_log
threshold = utc_now() - timedelta(minutes=STALE_ASSIGNED_MINUTES)
stale = session.exec(
select(IntelligentEvalTaskQueueDB).where(
IntelligentEvalTaskQueueDB.status == "assigned",
IntelligentEvalTaskQueueDB.assigned_at < threshold,
)
).all()
requeued = 0
abandoned = 0
for task in stale:
ev = session.get(IntelligentEvalDB, task.eval_id)
if ev is None or ev.status != IntelligentEvalStatus.EXECUTING.value:
continue
task.attempts += 1
if task.attempts >= MAX_TASK_ATTEMPTS:
task.status = "failed"
task.error = f"卡死重试 {task.attempts} 次仍未完成,按 ADR-0011 放弃"
task.completed_at = utc_now()
task.updated_at = utc_now()
abandoned += 1
append_decision_log(
task.eval_id,
"task_abandoned",
f"平台兜底:任务卡死重试 {task.attempts} 次仍未完成,置为失败",
"platform",
{"platform_supplemented": True, "task_id": task.id, "attempts": task.attempts},
session,
)
continue
task.status = "pending"
task.assigned_cron_id = None
task.assigned_at = None
task.updated_at = utc_now()
requeued += 1
if requeued or abandoned:
session.commit()
return requeued
def settle_tasks_for_finished_evals(session: Session) -> int:
"""Settle pending/assigned tasks of finished (non-executing) evals.
评估离开 executing 后,其待认领/执行中任务不再需要执行,按评估终态回收
ADR-0011 语义诚实completed → 任务 completedcancelled/failed →
任务 failederror="评估已终止,任务回收")。不回收会永久残留
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 评估的任务正常流转,不清理
if ev.status == IntelligentEvalStatus.COMPLETED.value:
task.status = "completed"
task.error = "评估已结束,任务不再需要执行"
else: # cancelled / failed
task.status = "failed"
task.error = "评估已终止,任务回收"
task.completed_at = utc_now()
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,
eval_id: Optional[str] = None,
) -> dict:
"""List task-queue entries with their eval names, newest first.
方案③的"定时触发"scan loop 每分钟扫描入队 + 触发 worker此前只有
Worker 消费端 APInext/assign/complete没有可查看的列表。这里提供
任务明细 + 状态分布统计,供前端任务队列监控页展示。``eval_id`` 把明细
限定到单个评估(执行过程视图的活动流);统计始终保持全局口径。
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)
if eval_id:
stmt = stmt.where(IntelligentEvalTaskQueueDB.eval_id == eval_id)
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}