fix(pheby): restore paged history across session resets
This commit is contained in:
+7
-1
@@ -99,6 +99,7 @@ Limits: chat text ≤ 64,000 chars; inbound WS frame ≤ 2 MiB (violations get
|
|||||||
"messages": [
|
"messages": [
|
||||||
{ "message_id": "m12", "role": "user", "text": "hey", "ts": "…|null" },
|
{ "message_id": "m12", "role": "user", "text": "hey", "ts": "…|null" },
|
||||||
{ "message_id": "m13", "role": "assistant", "text": "hi!", "ts": "…|null" } ],
|
{ "message_id": "m13", "role": "assistant", "text": "hi!", "ts": "…|null" } ],
|
||||||
|
"has_more": true,
|
||||||
"attachments": [],
|
"attachments": [],
|
||||||
"run": null,
|
"run": null,
|
||||||
"tools": [],
|
"tools": [],
|
||||||
@@ -107,7 +108,12 @@ Limits: chat text ≤ 64,000 chars; inbound WS frame ≤ 2 MiB (violations get
|
|||||||
```
|
```
|
||||||
|
|
||||||
History is the authoritative Hermes transcript (`role` is always `user` or
|
History is the authoritative Hermes transcript (`role` is always `user` or
|
||||||
`assistant`). The other fields form a recoverable snapshot: unexpired
|
`assistant`). It spans all durable reset and compression generations sharing
|
||||||
|
the Pheby conversation key. `limit` is capped at 500; the response returns the
|
||||||
|
newest page in chronological order. While `has_more` is true, request earlier
|
||||||
|
pages with `"before_message_id": "m12"`, using the first message ID from the
|
||||||
|
previous page as an exclusive cursor. This is display history only; model
|
||||||
|
context can still be compressed. The other fields form a recoverable snapshot: unexpired
|
||||||
attachments, the active run (if any), latest structured tool states, and
|
attachments, the active run (if any), latest structured tool states, and
|
||||||
pending approval/clarification requests. On reconnect, replace local state
|
pending approval/clarification requests. On reconnect, replace local state
|
||||||
with this snapshot, then consume new live events. Unknown conversation →
|
with this snapshot, then consume new live events. Unknown conversation →
|
||||||
|
|||||||
@@ -235,22 +235,26 @@ def _display_user_text(text: str) -> str:
|
|||||||
return body.split(end, 1)[0] if end in body else text
|
return body.split(end, 1)[0] if end in body else text
|
||||||
|
|
||||||
|
|
||||||
async def conversation_history(conversation_id: str, limit: int
|
async def conversation_history(conversation_id: str, limit: int,
|
||||||
) -> Tuple[List[Dict[str, Any]], bool]:
|
before_id: Optional[int] = None
|
||||||
"""Load transcript rows for a conversation from Hermes state.db.
|
) -> Tuple[List[Dict[str, Any]], bool, bool]:
|
||||||
|
"""Page the user-facing transcript across every reset generation.
|
||||||
|
|
||||||
Returns ``(messages, found)``. ``found`` is False when neither the
|
A gateway reset starts a new session ID for the same Pheby session key;
|
||||||
session store nor the session DB knows the conversation.
|
compression can archive rows within any generation. Both are display
|
||||||
|
history, even though neither belongs in the model's active context.
|
||||||
|
Returns (messages, found, has_more) in ascending row-ID order.
|
||||||
"""
|
"""
|
||||||
messages: List[Dict[str, Any]] = []
|
messages: List[Dict[str, Any]] = []
|
||||||
found = False
|
found = False
|
||||||
|
has_more = False
|
||||||
|
|
||||||
store = _session_store()
|
store = _session_store()
|
||||||
session_id: Optional[str] = None
|
session_id: Optional[str] = None
|
||||||
|
session_key = _session_key_for(conversation_id)
|
||||||
if store is not None:
|
if store is not None:
|
||||||
try:
|
try:
|
||||||
entry = await asyncio.to_thread(store.peek_session_id,
|
entry = await asyncio.to_thread(store.peek_session_id, session_key)
|
||||||
_session_key_for(conversation_id))
|
|
||||||
if entry:
|
if entry:
|
||||||
session_id = str(entry)
|
session_id = str(entry)
|
||||||
found = True
|
found = True
|
||||||
@@ -258,43 +262,67 @@ async def conversation_history(conversation_id: str, limit: int
|
|||||||
logger.debug("[pheby] peek_session_id failed", exc_info=True)
|
logger.debug("[pheby] peek_session_id failed", exc_info=True)
|
||||||
|
|
||||||
db = _session_db()
|
db = _session_db()
|
||||||
if db is not None and session_id:
|
if db is not None:
|
||||||
try:
|
try:
|
||||||
rows = await asyncio.to_thread(
|
# The routing index only knows the current generation. Query the
|
||||||
db.get_messages_as_conversation, session_id,
|
# persisted key to include resets even when parent_session_id is
|
||||||
include_row_ids=True)
|
# absent. Hermes's own decoder handles archived/compacted rows.
|
||||||
for row in rows[-limit:]:
|
ids = await asyncio.to_thread(
|
||||||
|
db._read_all,
|
||||||
|
"SELECT id FROM sessions WHERE session_key = ? AND source = ? ORDER BY started_at, id",
|
||||||
|
(session_key, "pheby"),
|
||||||
|
)
|
||||||
|
session_ids = [str(row["id"]) for row in ids]
|
||||||
|
if session_id and session_id not in session_ids:
|
||||||
|
session_ids.append(session_id)
|
||||||
|
if session_ids:
|
||||||
|
found = True
|
||||||
|
rows = []
|
||||||
|
for sid in session_ids:
|
||||||
|
rows.extend(await asyncio.to_thread(
|
||||||
|
db.get_messages_as_conversation, sid,
|
||||||
|
include_row_ids=True, include_compacted=True))
|
||||||
|
rows.sort(key=lambda row: int(row.get("_row_id") or 0))
|
||||||
|
visible = []
|
||||||
|
seen_replays = set()
|
||||||
|
for row in rows:
|
||||||
|
rid = row.get("_row_id")
|
||||||
role = row.get("role")
|
role = row.get("role")
|
||||||
if role not in ("user", "assistant"):
|
if not isinstance(rid, int) or role not in ("user", "assistant"):
|
||||||
continue
|
continue
|
||||||
content = row.get("content")
|
content = row.get("content")
|
||||||
text = content if isinstance(content, str) else str(content or "")
|
text = content if isinstance(content, str) else str(content or "")
|
||||||
if role == "user":
|
if role == "user":
|
||||||
|
# Reset/compression can carry the same user turn into a
|
||||||
|
# new generation; suppress exact replayed copies across
|
||||||
|
# the entire transcript, before applying the page cursor.
|
||||||
|
replay = (row.get("timestamp"), text)
|
||||||
|
if replay in seen_replays:
|
||||||
|
continue
|
||||||
|
seen_replays.add(replay)
|
||||||
text = _display_user_text(text)
|
text = _display_user_text(text)
|
||||||
# Tool-call rows can surface as assistant rows with empty
|
if before_id is not None and rid >= before_id:
|
||||||
# content; skip empties so the client transcript stays clean.
|
continue
|
||||||
if not text.strip() and role == "assistant":
|
if not text.strip() and role == "assistant":
|
||||||
continue
|
continue
|
||||||
messages.append({
|
visible.append({
|
||||||
"message_id": f"m{row.get('_row_id')}"
|
"message_id": f"m{rid}", "role": role, "text": text,
|
||||||
if isinstance(row.get("_row_id"), (int, str)) else None,
|
|
||||||
"role": role,
|
|
||||||
"text": text,
|
|
||||||
"ts": _iso(row.get("timestamp")),
|
"ts": _iso(row.get("timestamp")),
|
||||||
})
|
})
|
||||||
found = True
|
has_more = len(visible) > limit
|
||||||
except Exception:
|
messages = visible[-limit:]
|
||||||
logger.debug("[pheby] transcript load failed", exc_info=True)
|
except Exception as exc:
|
||||||
|
logger.warning("[pheby] transcript load failed", exc_info=True)
|
||||||
|
raise RuntimeError("Unable to load conversation transcript") from exc
|
||||||
|
|
||||||
# A router-known conversation with no messages yet is still "found" so a
|
# A router-known conversation with no messages yet is still "found".
|
||||||
# fresh client can open it as an empty chat.
|
|
||||||
if not found:
|
if not found:
|
||||||
server = _current_server()
|
server = _current_server()
|
||||||
if server is not None:
|
if server is not None:
|
||||||
name = await server.router.get_name(conversation_id)
|
name = await server.router.get_name(conversation_id)
|
||||||
if name is not None:
|
if name is not None:
|
||||||
found = True
|
found = True
|
||||||
return messages, found
|
return messages, found, has_more
|
||||||
|
|
||||||
|
|
||||||
async def rename_conversation(conversation_id: str, name: str) -> bool:
|
async def rename_conversation(conversation_id: str, name: str) -> bool:
|
||||||
@@ -353,11 +381,20 @@ async def delete_conversation(conversation_id: str) -> bool:
|
|||||||
if not router_known and not session_id:
|
if not router_known and not session_id:
|
||||||
return False
|
return False
|
||||||
|
|
||||||
if session_id and db is not None:
|
if db is not None:
|
||||||
try:
|
try:
|
||||||
deleted = await asyncio.to_thread(db.delete_session, session_id)
|
rows = await asyncio.to_thread(
|
||||||
if deleted is False:
|
db._read_all,
|
||||||
return False
|
"SELECT id FROM sessions WHERE session_key = ? AND source = ? ORDER BY started_at, id",
|
||||||
|
(session_key, "pheby"),
|
||||||
|
)
|
||||||
|
ids = [str(row["id"]) for row in rows]
|
||||||
|
if session_id and session_id not in ids:
|
||||||
|
ids.append(str(session_id))
|
||||||
|
if ids:
|
||||||
|
deleted = await asyncio.to_thread(db.delete_sessions, ids)
|
||||||
|
if deleted != len(ids):
|
||||||
|
return False
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.error("[pheby] session db delete failed", exc_info=True)
|
logger.error("[pheby] session db delete failed", exc_info=True)
|
||||||
return False
|
return False
|
||||||
|
|||||||
+12
-2
@@ -426,8 +426,17 @@ class PhebyServer:
|
|||||||
limit = max(1, min(int(limit), proto.MAX_HISTORY_MESSAGES))
|
limit = max(1, min(int(limit), proto.MAX_HISTORY_MESSAGES))
|
||||||
except (TypeError, ValueError):
|
except (TypeError, ValueError):
|
||||||
limit = proto.MAX_HISTORY_MESSAGES
|
limit = proto.MAX_HISTORY_MESSAGES
|
||||||
history, found = await self.bridge.conversation_history(
|
cursor = message.get("before_message_id")
|
||||||
conversation_id, limit)
|
before_id = None
|
||||||
|
if cursor is not None:
|
||||||
|
if (not isinstance(cursor, str) or not cursor.startswith("m")
|
||||||
|
or not cursor[1:].isdigit() or len(cursor) > 20):
|
||||||
|
await client.send_json(proto.error_event(
|
||||||
|
proto.ERR_BAD_REQUEST, "Invalid before_message_id", request_id))
|
||||||
|
return
|
||||||
|
before_id = int(cursor[1:])
|
||||||
|
history, found, has_more = await self.bridge.conversation_history(
|
||||||
|
conversation_id, limit, before_id=before_id)
|
||||||
if not found:
|
if not found:
|
||||||
await client.send_json(proto.error_event(
|
await client.send_json(proto.error_event(
|
||||||
proto.ERR_CONVERSATION_NOT_FOUND,
|
proto.ERR_CONVERSATION_NOT_FOUND,
|
||||||
@@ -438,6 +447,7 @@ class PhebyServer:
|
|||||||
"type": proto.S_CONVERSATION_HISTORY,
|
"type": proto.S_CONVERSATION_HISTORY,
|
||||||
"conversation_id": conversation_id,
|
"conversation_id": conversation_id,
|
||||||
"messages": history,
|
"messages": history,
|
||||||
|
"has_more": has_more,
|
||||||
"attachments": self.store.list_for_conversation(conversation_id),
|
"attachments": self.store.list_for_conversation(conversation_id),
|
||||||
**runtime,
|
**runtime,
|
||||||
**({"request_id": request_id} if request_id else {}),
|
**({"request_id": request_id} if request_id else {}),
|
||||||
|
|||||||
+69
-6
@@ -334,6 +334,16 @@ class TestConversations:
|
|||||||
ev = client.ws.events()[-1]
|
ev = client.ws.events()[-1]
|
||||||
assert ev["error"]["code"] == proto.ERR_BAD_REQUEST
|
assert ev["error"]["code"] == proto.ERR_BAD_REQUEST
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_invalid_history_cursor_rejected(self, tmp_path):
|
||||||
|
server = make_server(tmp_path)
|
||||||
|
client = FakeClientConnection()
|
||||||
|
client.authenticated = True
|
||||||
|
await server._handle_conversation_open(client, {
|
||||||
|
"type": proto.C_CONVERSATION_OPEN,
|
||||||
|
"conversation_id": "conv", "before_message_id": "m0 OR 1=1"}, "r1")
|
||||||
|
assert client.ws.events()[-1]["error"]["code"] == proto.ERR_BAD_REQUEST
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_ids_survive_router_reload(self, tmp_path):
|
async def test_ids_survive_router_reload(self, tmp_path):
|
||||||
server = make_server(tmp_path)
|
server = make_server(tmp_path)
|
||||||
@@ -443,8 +453,12 @@ class TestBridge:
|
|||||||
return "session-1"
|
return "session-1"
|
||||||
|
|
||||||
class DB:
|
class DB:
|
||||||
def get_messages_as_conversation(self, _session_id, include_row_ids=False):
|
def _read_all(self, sql, params):
|
||||||
|
return [{"id": "session-1"}]
|
||||||
|
|
||||||
|
def get_messages_as_conversation(self, _session_id, include_row_ids=False, include_compacted=False):
|
||||||
assert include_row_ids is True
|
assert include_row_ids is True
|
||||||
|
assert include_compacted is True
|
||||||
return [
|
return [
|
||||||
{"_row_id": 10, "role": "user", "content": "hello", "timestamp": 1_789_000_000.25},
|
{"_row_id": 10, "role": "user", "content": "hello", "timestamp": 1_789_000_000.25},
|
||||||
{"_row_id": 11, "role": "assistant", "content": "reply", "timestamp": 1_789_000_001.5},
|
{"_row_id": 11, "role": "assistant", "content": "reply", "timestamp": 1_789_000_001.5},
|
||||||
@@ -453,13 +467,58 @@ class TestBridge:
|
|||||||
monkeypatch.setattr(hb, "_session_store", lambda: Store())
|
monkeypatch.setattr(hb, "_session_store", lambda: Store())
|
||||||
monkeypatch.setattr(hb, "_session_db", lambda: DB())
|
monkeypatch.setattr(hb, "_session_db", lambda: DB())
|
||||||
|
|
||||||
history, found = await hb.conversation_history("conv", 20)
|
history, found, older = await hb.conversation_history("conv", 20)
|
||||||
|
|
||||||
assert found is True
|
assert found is True
|
||||||
|
assert older is False
|
||||||
assert [message["message_id"] for message in history] == ["m10", "m11"]
|
assert [message["message_id"] for message in history] == ["m10", "m11"]
|
||||||
assert history[0]["ts"] == "2026-09-10T00:26:40.250000+00:00"
|
assert history[0]["ts"] == "2026-09-10T00:26:40.250000+00:00"
|
||||||
assert history[1]["ts"] == "2026-09-10T00:26:41.500000+00:00"
|
assert history[1]["ts"] == "2026-09-10T00:26:41.500000+00:00"
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_history_spans_reset_sessions_and_pages_without_duplicates(self, monkeypatch):
|
||||||
|
from pheby import hermes_bridge as hb
|
||||||
|
|
||||||
|
class Store:
|
||||||
|
def peek_session_id(self, _key):
|
||||||
|
return "new"
|
||||||
|
|
||||||
|
class DB:
|
||||||
|
def _read_all(self, sql, params):
|
||||||
|
assert "session_key" in sql and params[1] == "pheby"
|
||||||
|
return [{"id": "old"}, {"id": "new"}]
|
||||||
|
|
||||||
|
def get_messages_as_conversation(self, sid, **kwargs):
|
||||||
|
assert kwargs == {"include_row_ids": True, "include_compacted": True}
|
||||||
|
start = 1 if sid == "old" else 4
|
||||||
|
return [{"_row_id": n, "role": "user", "content": f"turn {n}", "timestamp": float(n)}
|
||||||
|
for n in range(start, start + 3)]
|
||||||
|
|
||||||
|
monkeypatch.setattr(hb, "_session_store", lambda: Store())
|
||||||
|
monkeypatch.setattr(hb, "_session_db", lambda: DB())
|
||||||
|
latest, found, older = await hb.conversation_history("conv", 2)
|
||||||
|
assert found and older
|
||||||
|
assert [m["message_id"] for m in latest] == ["m5", "m6"]
|
||||||
|
previous, found, older = await hb.conversation_history("conv", 2, before_id=5)
|
||||||
|
assert found and older
|
||||||
|
assert [m["message_id"] for m in previous] == ["m3", "m4"]
|
||||||
|
first, found, older = await hb.conversation_history("conv", 2, before_id=3)
|
||||||
|
assert found and not older
|
||||||
|
assert [m["message_id"] for m in first] == ["m1", "m2"]
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_history_read_failure_is_not_reported_as_empty_chat(self, monkeypatch):
|
||||||
|
from pheby import hermes_bridge as hb
|
||||||
|
|
||||||
|
class DB:
|
||||||
|
def _read_all(self, _sql, _params):
|
||||||
|
raise OSError("database offline")
|
||||||
|
|
||||||
|
monkeypatch.setattr(hb, "_session_store", lambda: None)
|
||||||
|
monkeypatch.setattr(hb, "_session_db", lambda: DB())
|
||||||
|
with pytest.raises(RuntimeError, match="transcript"):
|
||||||
|
await hb.conversation_history("conv", 20)
|
||||||
|
|
||||||
def test_iso_parses_by_instant_for_naive_local_datetimes(self):
|
def test_iso_parses_by_instant_for_naive_local_datetimes(self):
|
||||||
# gateway session-store timestamps are naive LOCAL datetimes
|
# gateway session-store timestamps are naive LOCAL datetimes
|
||||||
# (gateway.session_lifecycle._now). last_active must come back with a
|
# (gateway.session_lifecycle._now). last_active must come back with a
|
||||||
@@ -776,9 +835,13 @@ class TestBridge:
|
|||||||
class DB:
|
class DB:
|
||||||
deleted = None
|
deleted = None
|
||||||
|
|
||||||
def delete_session(self, session_id):
|
def _read_all(self, sql, params):
|
||||||
self.deleted = session_id
|
assert "session_key" in sql and params == (session_key, "pheby")
|
||||||
return True
|
return [{"id": "session-old"}, {"id": "session-delete"}]
|
||||||
|
|
||||||
|
def delete_sessions(self, session_ids):
|
||||||
|
self.deleted = list(session_ids)
|
||||||
|
return len(session_ids)
|
||||||
|
|
||||||
store, db = Store(), DB()
|
store, db = Store(), DB()
|
||||||
from pheby import hermes_bridge as hb
|
from pheby import hermes_bridge as hb
|
||||||
@@ -786,7 +849,7 @@ class TestBridge:
|
|||||||
monkeypatch.setattr(hb, "_session_db", lambda: db)
|
monkeypatch.setattr(hb, "_session_db", lambda: db)
|
||||||
|
|
||||||
assert await hb.delete_conversation(cid) is True
|
assert await hb.delete_conversation(cid) is True
|
||||||
assert db.deleted == "session-delete"
|
assert db.deleted == ["session-old", "session-delete"]
|
||||||
assert session_key not in store._entries
|
assert session_key not in store._entries
|
||||||
assert store.saved is True
|
assert store.saved is True
|
||||||
assert await server.router.get_name(cid) is None
|
assert await server.router.get_name(cid) is None
|
||||||
|
|||||||
Reference in New Issue
Block a user