diff --git a/agent/deadline.py b/agent/deadline.py index 2be255c744..83d7f3cf2c 100644 --- a/agent/deadline.py +++ b/agent/deadline.py @@ -71,7 +71,7 @@ import sys import threading import time from dataclasses import dataclass -from typing import Any, Protocol, Awaitable, Callable, Optional +from typing import Any, Awaitable, Callable, Optional, Protocol logger = logging.getLogger(__name__) @@ -126,6 +126,12 @@ class SuspectableBackend(Protocol): instead of returning a poisoned handle to the cache. Consumers adopt incrementally (Phase 3b, one backend per PR), so the layer fails open: backends without the protocol are simply never marked. + + Adopter contract: ``mark_suspect`` MUST be cheap, non-blocking, and + must not acquire locks the guarded operation may hold. It runs inline — + on the event loop in the async flavor, and on the caller's thread in + the sync flavor while the wedged worker is still alive. Set a flag; + do the expensive health-check/recycle work in ``ensure_healthy``. """ def mark_suspect(self, reason: str) -> None: ... @@ -426,6 +432,11 @@ async def run_bounded_async( cleanup.add_done_callback(_consume_abandoned) # Phase 3a (#85125): the abandoned task may leave the backend # half-wedged; flag it so the owner recycles before reuse. + # Deliberately INLINE on the loop (adopter contract: mark_suspect is + # cheap and non-blocking). Running it synchronously guarantees the + # mark happens-before this BoundedResult returns AND before the + # ensure_future'd on_abandon cleanup can start (next loop tick) — an + # offloaded mark would race both. _mark_backend_suspect(backend, label, timeout_s) logger.warning( "[deadline] %r timed out after %.1fs; task abandoned", label, timeout_s @@ -507,8 +518,7 @@ def run_bounded_sync( # recycle/re-init in on_timeout never gets a stale flag on the healed # replacement. The sync flavor runs the mark inline — the protocol # contract requires mark_suspect to be cheap. - if backend is not None: - _mark_backend_suspect(backend, label, timeout_s) + _mark_backend_suspect(backend, label, timeout_s) if on_timeout is not None: try: on_timeout() diff --git a/tests/agent/test_deadline.py b/tests/agent/test_deadline.py index 6817bad16f..e66ac19b2e 100644 --- a/tests/agent/test_deadline.py +++ b/tests/agent/test_deadline.py @@ -575,7 +575,7 @@ def test_async_timeout_marks_backend_once(): backend = asyncio.run(drive()) assert len(backend.reasons) == 1 assert "phase3a" in backend.reasons[0] - assert "0.1" in backend.reasons[0] # rounded timeout present + assert "timed out after 0.1" in backend.reasons[0] def test_async_completion_never_marks_backend(): @@ -650,3 +650,68 @@ def test_mark_suspect_raising_never_corrupts_the_result(): result = asyncio.run(drive()) assert result.timed_out assert result.label == "phase3a-boom" + + +def test_sync_completion_never_marks_backend(): + from agent.deadline import run_bounded_sync + + backend = _RecordingBackend() + result = run_bounded_sync( + lambda: "ok", 5.0, label="phase3a-sync-ok", backend=backend + ) + assert not result.timed_out and result.value == "ok" + assert backend.reasons == [] + + +def test_sync_mark_happens_before_on_timeout(): + """The review-round ordering contract: mark BEFORE owner cleanup, so a + recycle in on_timeout never sees an unmarked backend (and a healed + replacement never inherits a stale flag).""" + from agent.deadline import run_bounded_sync + + backend = _RecordingBackend() + seen_at_cleanup: list[int] = [] + + def on_timeout(): + seen_at_cleanup.append(len(backend.reasons)) + + result = run_bounded_sync( + lambda: time.sleep(10), + 0.05, + label="phase3a-order-sync", + on_timeout=on_timeout, + backend=backend, + ) + assert result.timed_out + assert seen_at_cleanup == [1] # mark already applied when cleanup ran + + +def test_async_mark_happens_before_on_abandon_cleanup(): + """Pins the scheduling invariant the inline mark relies on: on_abandon + is ensure_future'd (can't start until the next loop tick), so the + synchronous mark always lands first. An offloaded (to_thread) mark + would break this — this test is the guard against that 'fix'.""" + from agent.deadline import run_bounded_async + + backend = _RecordingBackend() + seen_at_cleanup: list[int] = [] + cleaned = asyncio.Event() + + async def _cleanup(): + seen_at_cleanup.append(len(backend.reasons)) + cleaned.set() + + async def never(): + await asyncio.Event().wait() + + async def drive(): + result = await run_bounded_async( + never(), 0.05, label="phase3a-order-async", + on_abandon=_cleanup, backend=backend, + ) + await asyncio.wait_for(cleaned.wait(), timeout=5.0) + return result + + result = asyncio.run(drive()) + assert result.timed_out + assert seen_at_cleanup == [1] # mark already applied when cleanup started