diff --git a/docs/PROTOCOL.md b/docs/PROTOCOL.md index fbf04dd..85a26fe 100644 --- a/docs/PROTOCOL.md +++ b/docs/PROTOCOL.md @@ -99,6 +99,7 @@ Limits: chat text ≤ 64,000 chars; inbound WS frame ≤ 2 MiB (violations get "messages": [ { "message_id": "m12", "role": "user", "text": "hey", "ts": "…|null" }, { "message_id": "m13", "role": "assistant", "text": "hi!", "ts": "…|null" } ], + "has_more": true, "attachments": [], "run": null, "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 -`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 pending approval/clarification requests. On reconnect, replace local state with this snapshot, then consume new live events. Unknown conversation → diff --git a/plugin/pheby/hermes_bridge.py b/plugin/pheby/hermes_bridge.py index 2146e74..26388ed 100644 --- a/plugin/pheby/hermes_bridge.py +++ b/plugin/pheby/hermes_bridge.py @@ -235,22 +235,26 @@ def _display_user_text(text: str) -> str: return body.split(end, 1)[0] if end in body else text -async def conversation_history(conversation_id: str, limit: int - ) -> Tuple[List[Dict[str, Any]], bool]: - """Load transcript rows for a conversation from Hermes state.db. +async def conversation_history(conversation_id: str, limit: int, + before_id: Optional[int] = None + ) -> 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 - session store nor the session DB knows the conversation. + A gateway reset starts a new session ID for the same Pheby session key; + 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]] = [] found = False + has_more = False store = _session_store() session_id: Optional[str] = None + session_key = _session_key_for(conversation_id) if store is not None: try: - entry = await asyncio.to_thread(store.peek_session_id, - _session_key_for(conversation_id)) + entry = await asyncio.to_thread(store.peek_session_id, session_key) if entry: session_id = str(entry) 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) db = _session_db() - if db is not None and session_id: + if db is not None: try: - rows = await asyncio.to_thread( - db.get_messages_as_conversation, session_id, - include_row_ids=True) - for row in rows[-limit:]: + # The routing index only knows the current generation. Query the + # persisted key to include resets even when parent_session_id is + # absent. Hermes's own decoder handles archived/compacted rows. + 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") - if role not in ("user", "assistant"): + if not isinstance(rid, int) or role not in ("user", "assistant"): continue content = row.get("content") text = content if isinstance(content, str) else str(content or "") 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) - # Tool-call rows can surface as assistant rows with empty - # content; skip empties so the client transcript stays clean. + if before_id is not None and rid >= before_id: + continue if not text.strip() and role == "assistant": continue - messages.append({ - "message_id": f"m{row.get('_row_id')}" - if isinstance(row.get("_row_id"), (int, str)) else None, - "role": role, - "text": text, + visible.append({ + "message_id": f"m{rid}", "role": role, "text": text, "ts": _iso(row.get("timestamp")), }) - found = True - except Exception: - logger.debug("[pheby] transcript load failed", exc_info=True) + has_more = len(visible) > limit + messages = visible[-limit:] + 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 - # fresh client can open it as an empty chat. + # A router-known conversation with no messages yet is still "found". if not found: server = _current_server() if server is not None: name = await server.router.get_name(conversation_id) if name is not None: found = True - return messages, found + return messages, found, has_more 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: return False - if session_id and db is not None: + if db is not None: try: - deleted = await asyncio.to_thread(db.delete_session, session_id) - if deleted is False: - return False + rows = await asyncio.to_thread( + db._read_all, + "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: logger.error("[pheby] session db delete failed", exc_info=True) return False diff --git a/plugin/pheby/server.py b/plugin/pheby/server.py index aab7a4a..bd79712 100644 --- a/plugin/pheby/server.py +++ b/plugin/pheby/server.py @@ -426,8 +426,17 @@ class PhebyServer: limit = max(1, min(int(limit), proto.MAX_HISTORY_MESSAGES)) except (TypeError, ValueError): limit = proto.MAX_HISTORY_MESSAGES - history, found = await self.bridge.conversation_history( - conversation_id, limit) + cursor = message.get("before_message_id") + 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: await client.send_json(proto.error_event( proto.ERR_CONVERSATION_NOT_FOUND, @@ -438,6 +447,7 @@ class PhebyServer: "type": proto.S_CONVERSATION_HISTORY, "conversation_id": conversation_id, "messages": history, + "has_more": has_more, "attachments": self.store.list_for_conversation(conversation_id), **runtime, **({"request_id": request_id} if request_id else {}), diff --git a/tests/test_pheby.py b/tests/test_pheby.py index afe8c51..97911be 100644 --- a/tests/test_pheby.py +++ b/tests/test_pheby.py @@ -334,6 +334,16 @@ class TestConversations: ev = client.ws.events()[-1] 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 async def test_ids_survive_router_reload(self, tmp_path): server = make_server(tmp_path) @@ -443,8 +453,12 @@ class TestBridge: return "session-1" 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_compacted is True return [ {"_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}, @@ -453,13 +467,58 @@ class TestBridge: monkeypatch.setattr(hb, "_session_store", lambda: Store()) 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 older is False assert [message["message_id"] for message in history] == ["m10", "m11"] assert history[0]["ts"] == "2026-09-10T00:26:40.250000+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): # gateway session-store timestamps are naive LOCAL datetimes # (gateway.session_lifecycle._now). last_active must come back with a @@ -776,9 +835,13 @@ class TestBridge: class DB: deleted = None - def delete_session(self, session_id): - self.deleted = session_id - return True + def _read_all(self, sql, params): + assert "session_key" in sql and params == (session_key, "pheby") + 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() from pheby import hermes_bridge as hb @@ -786,7 +849,7 @@ class TestBridge: monkeypatch.setattr(hb, "_session_db", lambda: db) 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 store.saved is True assert await server.router.get_name(cid) is None