Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
48 changes: 48 additions & 0 deletions tests/gateway/test_agent_cache.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
Loading