fix(deadline): document the inline-mark contract; pin the ordering invariants (Phase 3a salvage round)
Record correction: the previous commit's message says the async flavor offloads mark_suspect via asyncio.to_thread — it does NOT (and must not). The mark is deliberately inline on the event loop: running it synchronously guarantees mark-happens-before-BoundedResult-return and mark-before-on_abandon-cleanup (cleanup is ensure_future'd and cannot start until the next loop tick). An offloaded mark would race both. The trade-off is that a slow adopter mark_suspect would block the loop (measured: a 2s mark stalls every coroutine for 2.003s), so the adopter contract is now explicit in the Protocol docstring and at the async call site: mark_suspect must be cheap, non-blocking, lock-free; expensive recycle work belongs in ensure_healthy. New pins so the negotiated semantics can't silently regress: - test_sync_mark_happens_before_on_timeout (the review-round ordering) - test_async_mark_happens_before_on_abandon_cleanup (the scheduling invariant an offloaded mark would break) - test_sync_completion_never_marks_backend (sync counterpart of the async completion test)
This commit is contained in:
@@ -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()
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user