Files
roboco/tests/unit/runtime/test_stale_claim_reaper.py
T
Renn F c12aad3005 fix(orchestrator): consolidate stale-heartbeat config + drop dead _task_svc slot
I1: claim_heartbeat_ttl_seconds (300s) overlapped semantically with the
pre-existing claim_stale_seconds (180s). Between 180-300s of silence,
trigger_filter queued duplicate spawns while the reaper hadn't yet
released the claim — exactly the dispatcher churn the reaper was
supposed to close. Collapse to one field (claim_stale_seconds, 180s);
reaper now consumes the same setting trigger_filter uses, so both
agree on 'stale' on the same tick and the reaper runs first.

I2: _task_svc injection slot on AgentOrchestrator.__init__ was
production-dead (always None) and only used by __new__-based test
instances. Drop the __init__ slot + the production branch in
_reap_stale_claims that read it. Tests still pre-bind on __new__
instances; the attribute exists per-instance, not per-class.
2026-05-03 04:57:25 +02:00

95 lines
3.2 KiB
Python

"""Reaper releases tasks whose last_heartbeat_at exceeds the TTL.
The orchestrator dispatcher periodically calls `_reap_stale_claims` to
release tasks whose holder has gone silent past the heartbeat TTL. The
schema-level column `last_heartbeat_at` (DateTime(timezone=True)) has
existed since migration 006; this test covers the runtime decision that
turns a stale heartbeat into a freed claim.
Datetimes used here are timezone-aware UTC because the underlying column
is tz-aware — comparing naive vs aware would raise TypeError in production
even though it would silently work against an in-memory mock.
"""
from __future__ import annotations
from datetime import UTC, datetime, timedelta
from unittest.mock import AsyncMock
from uuid import uuid4
import pytest
from roboco.runtime.orchestrator import AgentOrchestrator
@pytest.mark.asyncio
async def test_reap_stale_claims_releases_dead_holders() -> None:
"""A task past TTL is unclaimed; a fresh one is left alone."""
stale_id = uuid4()
fresh_id = uuid4()
now = datetime.now(UTC)
stale_task = type(
"T",
(),
{"id": stale_id, "last_heartbeat_at": now - timedelta(seconds=600)},
)()
fresh_task = type(
"T",
(),
{"id": fresh_id, "last_heartbeat_at": now - timedelta(seconds=10)},
)()
orch = AgentOrchestrator.__new__(AgentOrchestrator) # bypass __init__
orch._claim_heartbeat_ttl = 300
svc = AsyncMock()
svc.list_in_progress_or_claimed.return_value = [stale_task, fresh_task]
svc.unclaim_for_reaper = AsyncMock()
await orch._reap_with_service(svc)
svc.unclaim_for_reaper.assert_awaited_once_with(stale_id)
@pytest.mark.asyncio
async def test_reap_stale_claims_releases_holders_with_null_heartbeat() -> None:
"""A claimed task that never heartbeated (NULL column) is treated as stale."""
null_id = uuid4()
null_task = type("T", (), {"id": null_id, "last_heartbeat_at": None})()
orch = AgentOrchestrator.__new__(AgentOrchestrator)
orch._claim_heartbeat_ttl = 300
svc = AsyncMock()
svc.list_in_progress_or_claimed.return_value = [null_task]
svc.unclaim_for_reaper = AsyncMock()
await orch._reap_with_service(svc)
svc.unclaim_for_reaper.assert_awaited_once_with(null_id)
@pytest.mark.asyncio
async def test_reap_stale_claims_swallows_unclaim_errors() -> None:
"""An unclaim_for_reaper failure must not abort the reap loop."""
stale_a = uuid4()
stale_b = uuid4()
now = datetime.now(UTC)
task_a = type(
"T", (), {"id": stale_a, "last_heartbeat_at": now - timedelta(seconds=600)}
)()
task_b = type(
"T", (), {"id": stale_b, "last_heartbeat_at": now - timedelta(seconds=900)}
)()
orch = AgentOrchestrator.__new__(AgentOrchestrator)
orch._claim_heartbeat_ttl = 300
svc = AsyncMock()
svc.list_in_progress_or_claimed.return_value = [task_a, task_b]
svc.unclaim_for_reaper = AsyncMock(
side_effect=[RuntimeError("transient"), None]
)
await orch._reap_with_service(svc)
# Both stale tasks attempted; second succeeded despite first raising.
expected_attempts = 2
assert svc.unclaim_for_reaper.await_count == expected_attempts