diff --git a/tests/integration/test_worker_task_resilience.py b/tests/integration/test_worker_task_resilience.py new file mode 100644 index 0000000..d6a694d --- /dev/null +++ b/tests/integration/test_worker_task_resilience.py @@ -0,0 +1,161 @@ +"""Worker task resilience tests (Gitea issue #3 / P0). + +T1 — concurrent assign: only one worker wins; complete×stuck does not lose task. +T10 — fault loop: heartbeat busy → assigned → cron stuck → handler fails task + creates retry. +""" +from datetime import timedelta +import pytest +from unittest.mock import AsyncMock, MagicMock + +from sqlmodel import Session, create_engine, select + +from agenteval.intelligent_eval import cron_pool +from agenteval.intelligent_eval.models import IntelligentEvalStatus +from agenteval.intelligent_eval.task_queue import assign_task, complete_task, get_next_task +from agenteval.storage.db import ( + IntelligentEvalDB, + IntelligentEvalTaskQueueDB, + OpenClawCronPoolDB, + utc_now, +) + + +# T1 — concurrent assign: only one worker wins + + +def test_concurrent_assign_only_one_succeeds(db_session): + task = IntelligentEvalTaskQueueDB( + eval_id="eval1", status="pending", priority=1, reason="slot_due" + ) + db_session.add(task) + db_session.commit() + + ok1 = assign_task(task.id, "cron-A", db_session) + ok2 = assign_task(task.id, "cron-B", db_session) + + assert ok1 is True + assert ok2 is False + + db_session.refresh(task) + assert task.status == "assigned" + assert task.assigned_cron_id == "cron-A" + + +@pytest.mark.xfail( + reason=( + "Known race: assign_task lacks atomic CAS — two sessions can both read " + "status=pending and both commit. Tracked in .scratch/v111-architecture-scan.md " + "discoveries. Fix requires atomic UPDATE ... WHERE status=pending in assign_task." + ), + strict=False, +) +def test_concurrent_assign_via_two_sessions(tmp_db_path, db_session): + task = IntelligentEvalTaskQueueDB( + eval_id="eval1", status="pending", priority=1, reason="slot_due" + ) + db_session.add(task) + db_session.commit() + + engine2 = create_engine( + f"sqlite:///{tmp_db_path}", + connect_args={"check_same_thread": False}, + ) + session2 = Session(engine2) + try: + t_main = get_next_task(db_session) + t_other = get_next_task(session2) + assert t_main is not None and t_other is not None + assert t_main.id == t_other.id + + results = ( + assign_task(t_main.id, "cron-main", db_session), + assign_task(t_other.id, "cron-other", session2), + ) + assert sorted(results) == [False, True] + + session2.expire_all() + final = db_session.get(IntelligentEvalTaskQueueDB, t_main.id) + assert final.status == "assigned" + assert final.assigned_cron_id in {"cron-main", "cron-other"} + finally: + session2.close() + engine2.dispose() + + +# T1 — complete × stuck-handler convergence: task is never lost + + +def test_complete_and_stuck_converge_to_terminal_state(db_session): + task = IntelligentEvalTaskQueueDB( + eval_id="eval1", + status="assigned", + priority=1, + reason="slot_due", + assigned_cron_id="cron-1", + assigned_at=utc_now(), + ) + db_session.add(task) + db_session.commit() + + complete_task(task.id, True, None, db_session) + complete_task(task.id, False, "Cron stuck", db_session) + + db_session.refresh(task) + assert task.status in {"completed", "failed"} + + +# T10 — heartbeat busy → kill → requeue + + +def test_heartbeat_busy_kill_requeue_creates_retry_task(db_session): + import asyncio + + eval_db = IntelligentEvalDB( + name="requeue-test", + target_id="t1", + status=IntelligentEvalStatus.EXECUTING.value, + ) + db_session.add(eval_db) + db_session.commit() + + cron = OpenClawCronPoolDB( + openclaw_cron_id="cron-stuck", + status="busy", + last_active_at=utc_now() - timedelta(minutes=cron_pool.STUCK_THRESHOLD_MINUTES + 5), + current_eval_id=eval_db.id, + ) + db_session.add(cron) + + original = IntelligentEvalTaskQueueDB( + eval_id=eval_db.id, + status="assigned", + priority=1, + reason="slot_due", + assigned_cron_id="cron-stuck", + assigned_at=utc_now(), + ) + db_session.add(original) + db_session.commit() + + client = MagicMock() + client.delete_cron = AsyncMock(return_value=None) + client.create_cron = AsyncMock(return_value="cron-new") + + stuck = cron_pool.detect_stuck_crons(db_session) + assert any(c.openclaw_cron_id == "cron-stuck" for c in stuck) + + asyncio.run(cron_pool.handle_stuck_cron(cron, db_session, client)) + + db_session.refresh(original) + assert original.status == "failed" + assert "stuck" in (original.error or "").lower() + + retry = db_session.exec( + select(IntelligentEvalTaskQueueDB).where( + IntelligentEvalTaskQueueDB.eval_id == eval_db.id, + IntelligentEvalTaskQueueDB.reason == "cron_stuck_retry", + ) + ).first() + assert retry is not None + assert retry.status == "pending" + assert retry.assigned_cron_id is None