"""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, }, }, }