All checks were successful
CI / test (push) Successful in 3m8s
架构审查候选④:ADR-0008 收敛调度域时为保测试兼容留下的过渡 wrapper 使命结束。 删除 8 个浅封装:task_queue 的 _is_slot_due / _calculate_session_deficit (零调用死函数)+ _calculate_priority / _get_attention_reason,decision 的 _parse_time_slot(零调用死函数)+ _get_current_slot / _count_sessions_in_slot / _has_high_severity_issues。调用方直接使用 domain 模块。 5 个隔着 wrapper 测 domain 行为的测试迁到新文件 test_intelligent_eval_domain.py,直接锁定 domain,覆盖零丢失。 删除测试通过:复杂度直接消失,时段/欠账/优先级知识只剩 domain 一处。 870 tests passed,零行为变化。
243 lines
7.4 KiB
Python
243 lines
7.4 KiB
Python
"""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
|