From 98eb6ebdefc7992e02f1cd0df511f73dab7c8679 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Tue, 1 Sep 2026 22:17:09 -0700 Subject: [PATCH] fix(gateway): rescued FIFO orphan runs exactly once, chain stays in order (#99882) Follow-up to the salvaged #99912 rescue. The original helper left the rescued orphan IN the adapter slot while the caller also swapped it in as the current turn, so the post-turn _dequeue_pending_event ran the same follow-up a second time (live repro: TURNS=['Sent','C','C','D']). The helper now pops the oldest orphan and returns it to run as this turn, stages the NEXT orphan in the slot so the drain continues the chain in arrival order, and the call site parks the incoming message behind the chain via _enqueue_fifo (slot when free, overflow otherwise) instead of always appending to overflow. The rescued event's own source drives the turn so reply anchors point at the message actually being answered. Tests: contract updated for the new return type; added the 2-orphan chain case and the single-orphan-then-new-message slot case (both fail against the original helper shape). --- gateway/run.py | 75 +++++++------- tests/gateway/test_fifo_overflow_rescue.py | 113 +++++++++++++-------- 2 files changed, 111 insertions(+), 77 deletions(-) diff --git a/gateway/run.py b/gateway/run.py index f6dc23d3c8..ced859dd7e 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -9765,8 +9765,10 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew depth += 1 return depth - def _rescue_orphaned_overflow(self, session_key: str, adapter: Any) -> int: - """Stage any orphaned FIFO overflow into the pending slot (#99882). + def _rescue_orphaned_overflow( + self, session_key: str, adapter: Any + ) -> Optional["MessageEvent"]: + """Pop the oldest orphaned FIFO overflow event for an idle session (#99882). The FIFO overflow (``queued_events``) drains only at the post-turn promotion site (``_promote_queued_event`` inside the ``_run_agent`` @@ -9781,29 +9783,36 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew This rescue runs at the point where a NEW event arrives for a session that is NOT busy (the idle entry in ``_process_message_priority``). If the session went idle with a - populated overflow, the orphaned events are re-staged in FIFO - order ahead of the incoming event's own enqueue, so arrival order - (#28503) is preserved: the orphaned follow-ups run first, then the - new message. The slot must be empty at this point (the session is - idle), so staging is a plain slot assignment. + populated overflow, the oldest orphan is returned so the caller runs + it as THIS turn, and the next orphan (if any) is staged into the + slot so the post-turn drain continues the chain in arrival order + (#28503). The caller then enqueues the incoming event behind the + chain via ``_enqueue_fifo``. - Returns the number of orphaned events re-staged (0 when none). + The returned event is REMOVED from both stores: leaving it in the + slot while it also runs as the current turn would make the post-turn + ``_dequeue_pending_event`` run it a second time. + + Returns the orphaned event to run now, or ``None`` when there is + nothing to rescue (no overflow, slot occupied, or no slot storage). """ try: _q_state = self._peek_session_state(session_key) overflow = _q_state.conversation.queued_events if _q_state else None if not overflow: - return 0 + return None pending_slot = getattr(adapter, "_pending_messages", None) if not isinstance(pending_slot, dict) or pending_slot.get(session_key): # Slot occupied (busy) or no slot storage — promotion owns # this; do not fight it from the idle path. - return 0 - # Only stage ONE orphan into the slot — the remaining overflow - # stays queued and will drain via the normal - # _promote_queued_event post-turn promotion. Staging more than - # one would clobber the slot (single-slot design). - pending_slot[session_key] = overflow.pop(0) + return None + head = overflow.pop(0) + # Keep the slot occupied for the rest of the chain so the drain + # promotes in order and any mid-chain arrival routes to overflow + # instead of jumping the queue (same invariant as the drain's + # own _promote_queued_event). Only ONE event fits the slot. + if overflow: + pending_slot[session_key] = overflow.pop(0) logger.warning( "Rescued orphaned FIFO overflow event for idle session " "%s — it was queued during a busy window but the post-turn " @@ -9817,10 +9826,10 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew len(overflow), session_key, ) - return 1 + return head except Exception: logger.debug("FIFO overflow rescue failed for %s", session_key, exc_info=True) - return 0 + return None @staticmethod def _is_goal_continuation_event(event_or_text: Any) -> bool: @@ -19787,27 +19796,25 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew _orphan_adapter is not None and not bool(getattr(event, "internal", False)) and not event.get_command() - and self._queue_depth( - _quick_key, adapter=_orphan_adapter - ) >= 1 - and not _orphan_adapter._pending_messages.get(_quick_key) ): _rescued = self._rescue_orphaned_overflow( _quick_key, _orphan_adapter ) - if _rescued: - # Orphans staged in the slot; park the incoming event - # behind them in the overflow so it runs AFTER the - # rescued chain (FIFO). - _head = _orphan_adapter._pending_messages.get(_quick_key) - if _head is not None: - # The slot head runs as this turn via the loop - # below; enqueue the incoming event into overflow. - self._session_state(_quick_key).conversation.queued_events.append( - event - ) - # Swap: this turn now processes the oldest orphan. - event = _head + if _rescued is not None: + # The oldest orphan runs as THIS turn. Park the + # incoming event behind the rest of the chain: into the + # slot when the chain was a single orphan (so the + # post-turn drain picks it up), otherwise into overflow + # behind the already-staged next orphan (FIFO). + self._enqueue_fifo(_quick_key, event, _orphan_adapter) + event = _rescued + # Same session key by construction; carry the orphan's + # own source so reply anchors / thread metadata point + # at the message that is actually being answered. + _rescued_source = getattr(_rescued, "source", None) + if _rescued_source is not None: + source = _rescued_source + is_internal = bool(getattr(_rescued, "internal", False)) except Exception: logger.debug( "FIFO orphan rescue pre-claim failed for %s", diff --git a/tests/gateway/test_fifo_overflow_rescue.py b/tests/gateway/test_fifo_overflow_rescue.py index 932a7a735a..e1f3efd110 100644 --- a/tests/gateway/test_fifo_overflow_rescue.py +++ b/tests/gateway/test_fifo_overflow_rescue.py @@ -5,19 +5,17 @@ it lands in SessionState.conversation.queued_events (overflow) with the current turn's event occupying adapter._pending_messages[session_key] (slot). After the slot's turn completes, _promote_queued_event moves the overflow head into the slot. When that drain never runs — the -compression window ended through an exit that skipped the promotion -site — the overflow is silently orphaned: never dispatched, never -persisted, never logged. +busy window ended through an exit that skipped the promotion site +(/stop, turn exception, generation bump) — the overflow is silently +orphaned: never dispatched, never persisted, never logged. -The rescue in GatewayRunner._rescue_orphaned_overflow stages one orphan -into the slot on the next idle arrival, so FIFO order (#28503) holds. +The rescue in GatewayRunner._rescue_orphaned_overflow pops the oldest +orphan for the caller to run as the current turn and stages the next +orphan in the slot, so FIFO order (#28503) holds and nothing runs twice. """ -import asyncio from unittest.mock import MagicMock -import pytest - from gateway.platforms.base import ( BasePlatformAdapter, MessageEvent, @@ -56,33 +54,47 @@ def _text_event(text: str, msg_id: str) -> MessageEvent: ) +def _runner() -> GatewayRunner: + runner = GatewayRunner.__new__(GatewayRunner) + runner._queued_events = {} + return runner + + class TestRescueOrphanedOverflow: - def test_moves_overflow_head_to_empty_slot(self): - runner = GatewayRunner.__new__(GatewayRunner) - runner._queued_events = {} - # Minimal session_state with queued_events + def test_single_orphan_is_returned_and_removed_from_both_stores(self): + runner = _runner() adapter = _StubAdapter() session_key = "telegram:user:1" - # Two overflow items orphaned after slot turn completed - runner._session_state(session_key).conversation.queued_events.extend( - [_text_event("orphan-1", "o1"), _text_event("orphan-2", "o2")] + runner._session_state(session_key).conversation.queued_events.append( + _text_event("orphan-1", "o1") ) - # Slot empty (session went idle) assert session_key not in adapter._pending_messages rescued = runner._rescue_orphaned_overflow(session_key, adapter) - assert rescued == 1 - # Slot now holds the oldest orphan - assert adapter._pending_messages[session_key].text == "orphan-1" - # Remaining orphan stays in overflow - overflow = runner._session_state(session_key).conversation.queued_events - assert len(overflow) == 1 - assert overflow[0].text == "orphan-2" + assert rescued is not None and rescued.text == "orphan-1" + # The rescued event runs as the current turn, so it must NOT also + # sit in the slot — the post-turn drain would run it a second time. + assert session_key not in adapter._pending_messages + assert runner._session_state(session_key).conversation.queued_events == [] + + def test_two_orphans_return_oldest_and_stage_next_in_slot(self): + runner = _runner() + adapter = _StubAdapter() + session_key = "telegram:user:1b" + runner._session_state(session_key).conversation.queued_events.extend( + [_text_event("orphan-1", "o1"), _text_event("orphan-2", "o2")] + ) + + rescued = runner._rescue_orphaned_overflow(session_key, adapter) + + assert rescued is not None and rescued.text == "orphan-1" + # Slot now holds the NEXT orphan so the drain continues the chain. + assert adapter._pending_messages[session_key].text == "orphan-2" + assert runner._session_state(session_key).conversation.queued_events == [] def test_noop_when_slot_occupied(self): - runner = GatewayRunner.__new__(GatewayRunner) - runner._queued_events = {} + runner = _runner() adapter = _StubAdapter() session_key = "telegram:user:2" runner._session_state(session_key).conversation.queued_events.append( @@ -92,41 +104,56 @@ class TestRescueOrphanedOverflow: rescued = runner._rescue_orphaned_overflow(session_key, adapter) - assert rescued == 0 + assert rescued is None assert adapter._pending_messages[session_key].text == "busy-slot" assert len(runner._session_state(session_key).conversation.queued_events) == 1 def test_noop_when_no_overflow(self): - runner = GatewayRunner.__new__(GatewayRunner) - runner._queued_events = {} + runner = _runner() adapter = _StubAdapter() session_key = "telegram:user:3" rescued = runner._rescue_orphaned_overflow(session_key, adapter) - assert rescued == 0 + assert rescued is None assert session_key not in adapter._pending_messages def test_fifo_order_preserved_across_rescue_and_new_message(self): - """Oldest orphan runs first, new arrival last — FIFO (#28503).""" - runner = GatewayRunner.__new__(GatewayRunner) - runner._queued_events = {} + """Oldest orphan runs first, new arrival last — FIFO (#28503). + + Mirrors the idle-arrival call site: rescue → _enqueue_fifo(new). + """ + runner = _runner() adapter = _StubAdapter() session_key = "telegram:user:4" - - # Two orphans from the lost window runner._session_state(session_key).conversation.queued_events.extend( [_text_event("orphan-1", "o1"), _text_event("orphan-2", "o2")] ) - # New message arrives for idle session — rescue stages orphan-1 rescued = runner._rescue_orphaned_overflow(session_key, adapter) - assert rescued == 1 - # Simulate the caller enqueueing the new message behind the rescued chain - new_event = _text_event("new-msg", "new1") - runner._session_state(session_key).conversation.queued_events.append(new_event) + assert rescued is not None and rescued.text == "orphan-1" + runner._enqueue_fifo(session_key, _text_event("new-msg", "new1"), adapter) - # Drain order: slot (orphan-1), then overflow[0] (orphan-2), then new-msg - assert adapter._pending_messages[session_key].text == "orphan-1" - overflow_texts = [e.text for e in runner._session_state(session_key).conversation.queued_events] - assert overflow_texts == ["orphan-2", "new-msg"] + # Drain order after this turn: slot (orphan-2), then overflow (new-msg) + assert adapter._pending_messages[session_key].text == "orphan-2" + overflow_texts = [ + e.text for e in runner._session_state(session_key).conversation.queued_events + ] + assert overflow_texts == ["new-msg"] + + def test_single_orphan_then_new_message_lands_in_slot(self): + """With one orphan the slot is free after rescue, so the incoming + message must go to the slot (not overflow) or the drain never sees it.""" + runner = _runner() + adapter = _StubAdapter() + session_key = "telegram:user:5" + runner._session_state(session_key).conversation.queued_events.append( + _text_event("orphan-1", "o1") + ) + + rescued = runner._rescue_orphaned_overflow(session_key, adapter) + assert rescued is not None and rescued.text == "orphan-1" + runner._enqueue_fifo(session_key, _text_event("new-msg", "new1"), adapter) + + assert adapter._pending_messages[session_key].text == "new-msg" + assert runner._session_state(session_key).conversation.queued_events == []