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).
This commit is contained in:
Teknium
2026-09-01 22:17:09 -07:00
parent 625bcbd697
commit 98eb6ebdef
2 changed files with 111 additions and 77 deletions

View File

@@ -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",

View File

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