All checks were successful
CI / test (push) Successful in 3m45s
方案③ worker 由平台触发 openclaw agent(cron=manual-run,非真实 cron), fault_tolerance 的 stuck 检测不适用——agent 中断/失败时任务永久卡 assigned, scan 只查 pending 不再入队(死锁)。 requeue_stale_assigned_tasks:assigned 超过 10 分钟且评估仍 executing 的 任务重置为 pending(清空认领),平台 scan 循环随后重新触发 worker 重试。 接入 scan loop,每轮先 requeue 再 scan。
273 lines
9.1 KiB
Python
273 lines
9.1 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, 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
|