test(intelligent-eval): add worker task resilience tests (#3)
All checks were successful
CI / test (push) Successful in 4m21s
All checks were successful
CI / test (push) Successful in 4m21s
T1 concurrent assign: only one worker wins (sequential + dual-session) T1 complete×stuck convergence: task reaches terminal state, never lost T10 fault loop: heartbeat busy → cron stuck → handler fails task + creates retry 发现真 bug(assign_task 缺原子 CAS,详见 .scratch/v111-architecture-scan.md §6.1): 两个独立 session 并发 assign 同一 task 双认领。 以 xfail(strict=False) 留守卫,CI 不阻塞;修复后移除 xfail 即转绿。 该 bug 修复不在本 #3 范围,待开独立 issue。
This commit is contained in:
parent
9535a18a44
commit
3025dcd255
161
tests/integration/test_worker_task_resilience.py
Normal file
161
tests/integration/test_worker_task_resilience.py
Normal file
@ -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
|
||||||
Loading…
Reference in New Issue
Block a user