fix(gateway): fail closed when session turn lease times out

This commit is contained in:
texasich
2026-08-06 08:49:57 -05:00
committed by kshitij
parent 2ef294f020
commit 29af112cd4
3 changed files with 151 additions and 29 deletions
+24 -10
View File
@@ -2366,7 +2366,7 @@ from gateway.delivery import (
looks_like_telegram_private_chat_id,
resolve_delivery_transport,
)
from gateway.turn_lease import SessionTurnLeaseRegistry
from gateway.turn_lease import SessionTurnLeaseRegistry, TurnLeaseTimeoutError
from gateway.session_state import (
SERVICE_TIER_UNSET as _SERVICE_TIER_UNSET,
SessionState,
@@ -16624,19 +16624,33 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
# stale history base and interleaving transcript writes. Same-key
# messages never reach this point mid-turn (adapter + runner guards
# hold them), so the lock is uncontended outside the alias-key route.
# Fail-open: on timeout the token comes back degraded and the turn
# proceeds unserialized (never a wedged session). Released in
# _handle_message's finally via _release_turn_lease — granted per
# Fail-closed on timeout: never enter the transcript region without a
# lease. A bounded retry response is safer than recreating the exact
# concurrent-turn corruption this lease exists to prevent. Released
# in _handle_message's finally via _release_turn_lease — granted per
# (routing key, run generation) so a stale unwind can't release a
# newer turn's lease.
_lease_registry = getattr(self, "_turn_leases", None)
if _lease_registry is not None:
_lease_token = await _lease_registry.acquire(
session_entry.session_id,
owner_key=_quick_key,
generation=run_generation,
timeout=_float_env("HERMES_AGENT_TIMEOUT", 1800),
)
try:
_lease_token = await _lease_registry.acquire(
session_entry.session_id,
owner_key=_quick_key,
generation=run_generation,
timeout=_float_env("HERMES_AGENT_TIMEOUT", 1800),
)
except TurnLeaseTimeoutError:
logger.error(
"Deferring turn for routing key %s on session %s after "
"turn-lease timeout; transcript load was not started",
_quick_key,
session_entry.session_id,
)
return (
"⏳ Another turn is still running on this session. To "
"protect the transcript, this message was not processed. "
"Wait for the active turn to finish, then resend it."
)
if _lease_token is not None:
_lease_state = self._session_state(_quick_key).turn
_lease_state.lease_token = _lease_token
+47 -16
View File
@@ -31,9 +31,10 @@ Safety properties:
exact token is the current holder — a stale unwind can never release a
newer turn's lease (the #28686 ownership lesson applied). Release is
idempotent.
- **Fail-open on timeout.** A stuck holder degrades to today's unserialized
behavior with a loud ERROR after the configured wait — never a wedged
session. A degraded token holds nothing and releases nothing.
- **Fail-closed on timeout.** A timed-out waiter raises
:class:`TurnLeaseTimeoutError` and must be deferred by the dispatch layer.
It never runs concurrently against the still-live holder and therefore
cannot defeat the serialization invariant this lease exists to enforce.
- **Bounded registry.** The per-session lease map is size-capped; eviction
only ever removes idle (unheld, uncontended) entries, never a live lease.
@@ -61,17 +62,43 @@ logger = logging.getLogger(__name__)
DEFAULT_MAX_LEASES = 512
# Fallback wait (seconds) when the caller passes no positive timeout. Matches
# the gateway's default agent inactivity timeout so a stuck holder fails open
# on the same clock the turn itself would be declared stuck on.
# the gateway's default agent inactivity timeout. A caller that reaches this
# bound must defer the turn rather than run it concurrently with the holder.
DEFAULT_LEASE_WAIT = 1800.0
class TurnLeaseTimeoutError(TimeoutError):
"""The session lease stayed held for the caller's full wait budget.
This is a fail-closed signal: the caller did not acquire the lease and
must not enter the transcript load/run/flush region for this turn.
"""
def __init__(
self,
session_id: str,
*,
owner_key: str,
generation: int,
wait_seconds: float,
) -> None:
self.session_id = session_id
self.owner_key = owner_key
self.generation = generation
self.wait_seconds = wait_seconds
super().__init__(
f"turn lease wait timed out after {wait_seconds:.0f}s on session "
f"{session_id} for routing key {owner_key} (gen {generation})"
)
class TurnLeaseToken:
"""Handle returned by :meth:`SessionTurnLeaseRegistry.acquire`.
``degraded`` means the acquire timed out and the turn is proceeding
UNSERIALIZED (fail-open); such a token holds nothing and its release is a
no-op. ``released`` makes release idempotent.
``degraded`` is retained for compatibility with older callers and test
doubles, but :meth:`acquire` no longer returns degraded tokens: a timeout
raises :class:`TurnLeaseTimeoutError` instead. ``released`` makes release
idempotent.
"""
__slots__ = ("session_id", "owner_key", "generation", "degraded", "released")
@@ -161,9 +188,10 @@ class SessionTurnLeaseRegistry:
) -> Optional[TurnLeaseToken]:
"""Acquire the turn lease for ``session_id``, waiting if held.
Returns a :class:`TurnLeaseToken` — degraded when the wait timed out
(fail-open: caller proceeds unserialized). Returns ``None`` for a
falsy ``session_id``.
Returns a held :class:`TurnLeaseToken`. Raises
:class:`TurnLeaseTimeoutError` when the wait budget expires; the caller
must defer rather than enter the serialized region. Returns ``None``
for a falsy ``session_id``.
"""
if not session_id:
return None
@@ -194,9 +222,8 @@ class SessionTurnLeaseRegistry:
logger.error(
"turn lease wait timed out after %.0fs on session %s "
"(waiter: routing key %s gen %s; holder: routing key %s "
"gen %s) — failing open: this turn runs UNSERIALIZED against "
"the stuck holder rather than wedging the session; transcript "
"writes may interleave",
"gen %s) — failing closed: refusing to run this turn "
"UNSERIALIZED against the still-live holder",
wait,
session_id,
owner_key,
@@ -204,8 +231,12 @@ class SessionTurnLeaseRegistry:
holder.owner_key if holder else "?",
holder.generation if holder else "?",
)
token.degraded = True
return token
raise TurnLeaseTimeoutError(
session_id,
owner_key=owner_key,
generation=generation,
wait_seconds=wait,
) from None
lease.holder = token
lease.acquired_at = time.time()
+80 -3
View File
@@ -11,8 +11,8 @@ Covers:
- distinct sessions do not contend
- generation-scoped, idempotent release: a stale unwind can never free a
newer turn's lease; double-release is a no-op
- timeout fail-open: a stuck holder degrades to unserialized with a degraded
token, never a wedged session, and the degraded token releases nothing
- timeout fail-closed: a timed-out waiter never enters the transcript region,
and the gateway returns a safe retry response before loading history
- registry stays bounded; live leases are never evicted
- GatewayRunner._release_turn_lease wiring (bare-runner safe, token-scoped)
"""
@@ -89,10 +89,87 @@ def test_distinct_sessions_do_not_contend():
# ---------------------------------------------------------------------------
# Fail-open on timeout
# Timeout safety
# ---------------------------------------------------------------------------
def test_timeout_fails_closed_instead_of_authorizing_an_unserialized_turn():
"""A timed-out waiter must never run against the still-live holder.
Returning a degraded token used to authorize exactly that unsafe path.
The two turns could then load the same history base and interleave their
transcript writes, defeating the serialization invariant this lease owns.
"""
from gateway.turn_lease import TurnLeaseTimeoutError
async def scenario():
registry = SessionTurnLeaseRegistry()
holder = await registry.acquire(
"sess-timeout", owner_key="key-a", generation=1, timeout=1
)
assert holder is not None
with pytest.raises(TurnLeaseTimeoutError):
await registry.acquire(
"sess-timeout", owner_key="key-b", generation=1, timeout=0.02
)
# The timeout neither steals nor releases the live holder's lease.
assert registry._leases["sess-timeout"].holder is holder
assert registry.release(holder) is True
# Once the holder releases, a later turn can acquire normally.
successor = await registry.acquire(
"sess-timeout", owner_key="key-b", generation=2, timeout=1
)
assert successor is not None
assert registry.release(successor) is True
_run(scenario())
@pytest.mark.asyncio
async def test_gateway_defers_timed_out_lease_before_loading_transcript(
monkeypatch, tmp_path
):
"""The dispatch layer turns a lease timeout into a safe retry response.
Most importantly, transcript loading and agent execution must not start:
both would operate without the per-session serialization guarantee.
"""
from tests.gateway.test_42039_duplicate_user_message import (
_bootstrap,
_event,
_source,
)
runner = _bootstrap(monkeypatch, tmp_path)
runner._turn_leases = SessionTurnLeaseRegistry()
holder = await runner._turn_leases.acquire(
"sess-dedup", owner_key="holder-key", generation=1, timeout=1
)
assert holder is not None
monkeypatch.setenv("HERMES_AGENT_TIMEOUT", "0.02")
runner.session_store.load_transcript.side_effect = AssertionError(
"transcript must not load after a turn-lease timeout"
)
runner._run_agent = pytest.fail
try:
response = await runner._handle_message_with_agent(
_event(), _source(), "agent:main:telegram:group:-1001:12345", 1
)
finally:
assert runner._turn_leases.release(holder) is True
assert isinstance(response, str)
assert "still running" in response.lower()
assert "not processed" in response.lower()
assert "resend" in response.lower()
runner.session_store.load_transcript.assert_not_called()
# ---------------------------------------------------------------------------
# Bounded registry
# ---------------------------------------------------------------------------