"""Unit tests for intelligent eval task queue. 时段/欠账/优先级等调度域的测试已随 ADR-0008 过渡 wrapper 删除迁移到 test_intelligent_eval_domain.py,直接锁定 domain 模块。 """ from datetime import timedelta from agenteval.intelligent_eval import task_queue from agenteval.intelligent_eval.models import IntelligentEvalStatus from agenteval.storage.db import ( IntelligentEvalDB, IntelligentEvalTaskQueueDB, utc_now, ) from sqlmodel import Session, select def test_has_pending_task(db_session: Session): """Test pending task detection (去重).""" eval_db = IntelligentEvalDB( name="test", target_id="target1", status=IntelligentEvalStatus.EXECUTING.value, ) db_session.add(eval_db) db_session.commit() # No pending task initially assert task_queue._has_pending_task(eval_db.id, db_session) is False # Add pending task task = IntelligentEvalTaskQueueDB( eval_id=eval_db.id, status="pending", priority=1, reason="slot_due", ) db_session.add(task) db_session.commit() # Should detect pending task assert task_queue._has_pending_task(eval_db.id, db_session) is True def test_scan_and_enqueue_tasks(db_session: Session): """Test task scanning and enqueueing.""" # Create eval that needs attention eval_db = IntelligentEvalDB( name="test", target_id="target1", status=IntelligentEvalStatus.EXECUTING.value, started_at=utc_now() - timedelta(hours=9), ) eval_db.set_plan({ "time_distribution": [{"time_slot": "8-10h", "sessions": 2}], "estimated_sessions": 2, }) db_session.add(eval_db) db_session.commit() # Scan and enqueue enqueued = task_queue.scan_and_enqueue_tasks(db_session) assert enqueued == 1 # Verify task created task = db_session.exec( select(IntelligentEvalTaskQueueDB).where(IntelligentEvalTaskQueueDB.eval_id == eval_db.id) ).first() assert task is not None assert task.status == "pending" assert task.reason == "slot_due" assert task.priority < 100 # Should have reduced priority # Scan again, should not create duplicate enqueued = task_queue.scan_and_enqueue_tasks(db_session) assert enqueued == 0 def test_get_next_task(db_session: Session): """Test getting next task (highest priority).""" # Create tasks with different priorities task1 = IntelligentEvalTaskQueueDB( eval_id="eval1", status="pending", priority=50, reason="slot_due", ) task2 = IntelligentEvalTaskQueueDB( eval_id="eval2", status="pending", priority=10, # Higher priority (lower number) reason="slot_due", ) task3 = IntelligentEvalTaskQueueDB( eval_id="eval3", status="assigned", # Not pending priority=1, reason="slot_due", ) db_session.add_all([task1, task2, task3]) db_session.commit() # Should get task2 (priority 10) next_task = task_queue.get_next_task(db_session) assert next_task is not None assert next_task.id == task2.id def test_assign_task(db_session: Session): """Test task assignment.""" task = IntelligentEvalTaskQueueDB( eval_id="eval1", status="pending", priority=1, reason="slot_due", ) db_session.add(task) db_session.commit() # Assign task success = task_queue.assign_task(task.id, "cron1", db_session) assert success is True # Verify assignment db_session.refresh(task) assert task.status == "assigned" assert task.assigned_cron_id == "cron1" assert task.assigned_at is not None # Try to assign again (should fail) success = task_queue.assign_task(task.id, "cron2", db_session) assert success is False def test_complete_task(db_session: Session): """Test task completion.""" task = IntelligentEvalTaskQueueDB( eval_id="eval1", status="assigned", priority=1, reason="slot_due", assigned_cron_id="cron1", ) db_session.add(task) db_session.commit() # Complete task successfully success = task_queue.complete_task(task.id, True, None, db_session) assert success is True # Verify completion db_session.refresh(task) assert task.status == "completed" assert task.completed_at is not None assert task.error is None # Complete task with error task2 = IntelligentEvalTaskQueueDB( eval_id="eval2", status="assigned", priority=1, reason="slot_due", assigned_cron_id="cron1", ) db_session.add(task2) db_session.commit() success = task_queue.complete_task(task2.id, False, "test error", db_session) assert success is True db_session.refresh(task2) assert task2.status == "failed" assert task2.error == "test error" def test_settle_tasks_for_finished_evals(db_session: Session): """评估已结束时按终态回收其 pending/assigned 任务(ADR-0011)。 回归:t480 队列出现"评估已 completed 却有待认领/执行中任务"的残留。 completed → 任务 completed;cancelled/failed → 任务 failed。 """ # 已结束的评估(completed):pending + assigned 任务都应被清理 done_eval = IntelligentEvalDB( name="done", target_id="t1", status=IntelligentEvalStatus.COMPLETED.value, started_at=utc_now(), ) db_session.add(done_eval) db_session.commit() pending_task = IntelligentEvalTaskQueueDB( eval_id=done_eval.id, status="pending", priority=1, reason="slot_due", ) assigned_task = IntelligentEvalTaskQueueDB( eval_id=done_eval.id, status="assigned", priority=1, reason="slot_due", assigned_cron_id="cron1", ) db_session.add_all([pending_task, assigned_task]) db_session.commit() # 终止的评估(failed):任务应回收为 failed failed_eval = IntelligentEvalDB( name="failed", target_id="t1", status=IntelligentEvalStatus.FAILED.value, started_at=utc_now(), ) db_session.add(failed_eval) db_session.commit() failed_task = IntelligentEvalTaskQueueDB( eval_id=failed_eval.id, status="pending", priority=1, reason="slot_due", ) db_session.add(failed_task) db_session.commit() # 仍在执行中的评估:任务应保留(正常流转) running_eval = IntelligentEvalDB( name="running", target_id="t1", status=IntelligentEvalStatus.EXECUTING.value, started_at=utc_now(), ) db_session.add(running_eval) db_session.commit() live_task = IntelligentEvalTaskQueueDB( eval_id=running_eval.id, status="assigned", priority=1, reason="slot_due", assigned_cron_id="cron2", ) db_session.add(live_task) db_session.commit() settled = task_queue.settle_tasks_for_finished_evals(db_session) assert settled == 3 # 只清理已结束评估的 3 个任务 db_session.refresh(pending_task) db_session.refresh(assigned_task) db_session.refresh(failed_task) db_session.refresh(live_task) assert pending_task.status == "completed" assert assigned_task.status == "completed" assert failed_task.status == "failed" assert failed_task.error == "评估已终止,任务回收" assert live_task.status == "assigned" # executing 评估的任务保留 # 幂等:再跑一次不再清理 assert task_queue.settle_tasks_for_finished_evals(db_session) == 0