diff --git a/gateway/run.py b/gateway/run.py index 52ae139c5c..21ab492bfa 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -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 diff --git a/gateway/turn_lease.py b/gateway/turn_lease.py index 802de0f178..9b4349f4a2 100644 --- a/gateway/turn_lease.py +++ b/gateway/turn_lease.py @@ -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() diff --git a/tests/gateway/test_turn_lease.py b/tests/gateway/test_turn_lease.py index 45fd546625..f9379f8740 100644 --- a/tests/gateway/test_turn_lease.py +++ b/tests/gateway/test_turn_lease.py @@ -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 # ---------------------------------------------------------------------------