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 == []