diff --git a/gateway/run.py b/gateway/run.py index 93ef09b73d..c35593ab8b 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -2986,6 +2986,7 @@ from gateway.session import ( SessionStore, SessionSource, SessionContext, + TranscriptReadError, build_session_context, build_session_context_prompt, build_channel_continuity_note, @@ -20867,8 +20868,23 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew # began processing if the gateway died while it was still waiting. await self._mark_durable_active_turn(event, session_entry.session_key) - # Load conversation history from transcript - history = await self.async_session_store.load_transcript(session_entry.session_id) + # Load conversation history from transcript. An unreadable canonical + # store is not an empty conversation: stop before the agent can invent + # continuity from a plausible-looking []. This return happens before + # the broad cleanup finally below, so restore task-local context here; + # the outer dispatch still clears the durable marker and turn lease. + try: + history = await self.async_session_store.load_transcript( + session_entry.session_id + ) + except TranscriptReadError: + self._clear_session_env(_session_env_tokens) + return ( + "⚠️ This session's history is temporarily unavailable, so " + "this message was not processed. Ask the operator to inspect " + "state.db, then resend after it is healthy. Use /reset only " + "if you intentionally want to start a new conversation." + ) # ----------------------------------------------------------------- # Session hygiene: auto-compress pathologically large transcripts diff --git a/gateway/session.py b/gateway/session.py index 568322d83f..a6c76bb3d4 100644 --- a/gateway/session.py +++ b/gateway/session.py @@ -23,6 +23,14 @@ from typing import Dict, List, Optional, Any logger = logging.getLogger(__name__) +class TranscriptReadError(RuntimeError): + """Raised when persisted history cannot be read safely.""" + + def __init__(self, session_id: str) -> None: + self.session_id = session_id + super().__init__(f"transcript read failed for session {session_id}") + + def _now() -> datetime: """Return the current local time.""" return datetime.now() @@ -4400,15 +4408,17 @@ class SessionStore: session_id, repair_alternation=True ) except Exception as e: - # A failed read must be distinguishable from an empty transcript: - # downstream guards treat [] as "nothing persisted" and may make - # routing decisions on it (#82616). WARNING, not DEBUG. - logger.warning( - "Transcript read failed for session %s (returning empty; " - "downstream must not treat this as data loss): %s", - session_id, e, + # Empty history is valid data; a failed canonical read is not. + # Preserve that distinction so live-replay callers can fail closed + # instead of starting the model with a plausible-looking []. + logger.error( + "Transcript read failed for session %s; refusing to treat the " + "conversation as empty: %s", + session_id, + e, + exc_info=True, ) - return [] + raise TranscriptReadError(session_id) from e def rewind_session( self, diff --git a/tests/gateway/test_42039_duplicate_user_message.py b/tests/gateway/test_42039_duplicate_user_message.py index 13a73181f6..3ddc30d90e 100644 --- a/tests/gateway/test_42039_duplicate_user_message.py +++ b/tests/gateway/test_42039_duplicate_user_message.py @@ -24,7 +24,7 @@ import pytest import gateway.run as gateway_run from gateway.config import GatewayConfig, Platform from gateway.platforms.base import MessageEvent -from gateway.session import SessionEntry, SessionSource +from gateway.session import SessionEntry, SessionSource, TranscriptReadError def _bootstrap(monkeypatch, tmp_path): @@ -185,6 +185,24 @@ async def test_not_new_messages_skip_db_when_agent_has_session_db( ) +@pytest.mark.asyncio +async def test_transcript_read_failure_stops_turn_before_agent_or_append( + monkeypatch, tmp_path +): + runner = _bootstrap(monkeypatch, tmp_path) + runner.session_store.load_transcript.side_effect = TranscriptReadError("sess-dedup") + runner._run_agent = AsyncMock() + + response = await runner._handle_message_with_agent( + _event(), _source(), "agent:main:telegram:group:-1001:12345", 1 + ) + + assert "history is temporarily unavailable" in response + assert "not processed" in response + runner._run_agent.assert_not_awaited() + runner.session_store.append_to_transcript.assert_not_called() + + # ── Post-stream MEDIA delivery keeps prior-turn deduplication ────────── diff --git a/tests/gateway/test_session_continuity_82616.py b/tests/gateway/test_session_continuity_82616.py index d99bccfa77..7f9a498dc6 100644 --- a/tests/gateway/test_session_continuity_82616.py +++ b/tests/gateway/test_session_continuity_82616.py @@ -201,6 +201,28 @@ class TestPeerResolutionRecency: class TestLoadTranscriptReroutes: + def test_load_transcript_raises_when_message_read_fails(self, tmp_path, monkeypatch): + from gateway.session import SessionStore, TranscriptReadError + + from gateway.config import GatewayConfig + + store = SessionStore(sessions_dir=tmp_path / "gw-failed-read", config=GatewayConfig()) + db = store._db + assert db is not None + monkeypatch.setattr(db, "get_compression_tip", lambda _session_id: None) + + def _malformed(_session_id, *, repair_alternation): + assert repair_alternation is True + raise RuntimeError("database disk image is malformed") + + monkeypatch.setattr(db, "get_messages_as_conversation", _malformed) + + with pytest.raises(TranscriptReadError) as exc_info: + store.load_transcript("existing-session") + + assert exc_info.value.session_id == "existing-session" + assert isinstance(exc_info.value.__cause__, RuntimeError) + def test_load_transcript_follows_reroute_chain(self, tmp_path): from gateway.session import SessionStore