fix(gateway): fail closed when transcript reads fail
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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 ──────────
|
||||
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
Reference in New Issue
Block a user