AgentEvalTool/backend/agenteval/intelligent_eval/task_queue.py
sinohqb 38e3817433
Some checks failed
CI / test (push) Has been cancelled
fix(intelligent-eval): atomic CAS in assign_task and complete_task (resolves §6.1)
Replace read-check-write in task_queue.assign_task with UPDATE...WHERE
status='pending' and decide on rowcount so two concurrent workers
cannot both claim the same task. Also harden complete_task with the
same CAS pattern (status='assigned') so a worker + stuck-handler
double-complete leaves the DB in one state.

The xfail guard in test_worker_task_resilience now passes (4/4).
2026-08-14 14:52:34 +08:00

229 lines
7.5 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, 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,
},
},
}