From fafdc2344d822bf6eba621002b2c6e08bf84f110 Mon Sep 17 00:00:00 2001 From: liuhao1024 Date: Thu, 11 Jun 2026 23:19:41 +0800 Subject: [PATCH] fix(gateway): reset _last_flushed_db_idx on cached-agent reuse When the gateway reuses a cached AIAgent across turns, _init_cached_agent_for_turn() reset _api_call_count but left _last_flushed_db_idx stale from the previous turn. This caused _flush_messages_to_session_db() to compute flush_from from the stale cursor rather than the new turn's conversation_history length, silently skipping the assistant reply and producing consecutive-user transcript rows that corrupt context replay on subsequent turns. Reset _last_flushed_db_idx to 0 at the start of every cached-agent turn (regardless of interrupt depth) so the flush cursor is always aligned with the current turn's history boundary. Fixes #44327 --- gateway/run.py | 8 ++++++ tests/gateway/test_agent_cache.py | 48 +++++++++++++++++++++++++++++++ 2 files changed, 56 insertions(+) diff --git a/gateway/run.py b/gateway/run.py index 897cb85f65232..885cfa7fa406d 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -12433,6 +12433,14 @@ def _init_cached_agent_for_turn(agent: Any, interrupt_depth: int) -> None: agent._last_activity_ts = time.time() agent._last_activity_desc = "starting new turn (cached)" agent._api_call_count = 0 + # Reset the SessionDB flush cursor so _flush_messages_to_session_db + # recomputes flush_from from the new turn's conversation_history + # length instead of carrying over a stale offset from the previous + # turn. Without this, a cached agent with _last_flushed_db_idx=N + # skips persisting the assistant reply when the new turn's history + # length is â‰ĪN, producing consecutive user rows in the transcript + # and context-corrupting replay on subsequent turns (#44327). + agent._last_flushed_db_idx = 0 def _release_evicted_agent_soft(self, agent: Any) -> None: """Soft cleanup for cache-evicted agents — preserves session tool state. diff --git a/tests/gateway/test_agent_cache.py b/tests/gateway/test_agent_cache.py index 37f8b51a458d4..6a239d9996203 100644 --- a/tests/gateway/test_agent_cache.py +++ b/tests/gateway/test_agent_cache.py @@ -1448,6 +1448,54 @@ def test_watchdog_accumulation_across_recursive_turns(self): ) + def test_flush_cursor_reset_on_fresh_turn(self): + """_last_flushed_db_idx must be reset to 0 on every cached-agent + turn so _flush_messages_to_session_db recomputes flush_from from + the new turn's conversation_history length (#44327).""" + from gateway.run import GatewayRunner + + agent = self._fake_agent() + agent._last_flushed_db_idx = 42 # stale from previous turn + + with patch("gateway.run.time") as mock_time: + mock_time.time.return_value = _FAKE_NOW + GatewayRunner._init_cached_agent_for_turn(agent, interrupt_depth=0) + + assert agent._last_flushed_db_idx == 0, ( + "_last_flushed_db_idx must be reset on cached-agent reuse; " + "stale cursor causes _flush_messages_to_session_db to skip " + "assistant replies, producing consecutive-user transcript rows" + ) + + def test_flush_cursor_reset_on_interrupt_turn(self): + """_last_flushed_db_idx must be reset even on interrupt-recursive + turns (depth>0) — the flush cursor is turn-scoped, not + watchdog-scoped.""" + from gateway.run import GatewayRunner + + agent = self._fake_agent() + agent._last_flushed_db_idx = 30 + + GatewayRunner._init_cached_agent_for_turn(agent, interrupt_depth=1) + + assert agent._last_flushed_db_idx == 0, ( + "_last_flushed_db_idx must reset at any interrupt depth" + ) + + def test_flush_cursor_reset_when_already_zero(self): + """Reset is idempotent — no-op when cursor is already 0.""" + from gateway.run import GatewayRunner + + agent = self._fake_agent() + agent._last_flushed_db_idx = 0 + + with patch("gateway.run.time") as mock_time: + mock_time.time.return_value = _FAKE_NOW + GatewayRunner._init_cached_agent_for_turn(agent, interrupt_depth=0) + + assert agent._last_flushed_db_idx == 0 + + class TestAgentConfigSignatureUserId: """Shared-thread cache must not reuse an agent across users.