Skip to content
Open
18 changes: 18 additions & 0 deletions gateway/platforms/api_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -4505,7 +4505,11 @@ async def _handle_session_messages(self, request: "web.Request") -> "web.Respons
resolved_id = await asyncio.to_thread(db.resolve_resume_session_id, session_id)
raw_limit = request.query.get("limit")
raw_offset = request.query.get("offset", "0")
raw_before_id = request.query.get("before_id")
order = request.query.get("order")
include_compacted = _coerce_request_bool(
request.query.get("include_compacted"), default=False
)
if order not in (None, "oldest", "latest"):
return web.json_response(
_openai_error(
Expand All @@ -4517,9 +4521,11 @@ async def _handle_session_messages(self, request: "web.Request") -> "web.Respons
try:
offset = int(raw_offset)
requested_limit = None if raw_limit is None else int(raw_limit)
before_id = None if raw_before_id is None else int(raw_before_id)
except (TypeError, ValueError):
offset = -1
requested_limit = -1
before_id = -1
if offset < 0 or (requested_limit is not None and requested_limit < 0):
return web.json_response(
_openai_error(
Expand All @@ -4529,6 +4535,15 @@ async def _handle_session_messages(self, request: "web.Request") -> "web.Respons
status=400,
)

if before_id is not None and (before_id < 1 or order != "latest" or offset != 0):
return web.json_response(
_openai_error(
"before_id requires a positive integer, order=latest, and no offset",
code="invalid_pagination",
),
status=400,
)

default_page = requested_limit is None
latest_page = order == "latest" or (order is None and default_page)
limit = 500 if default_page else min(requested_limit, 500)
Expand All @@ -4538,6 +4553,8 @@ async def _handle_session_messages(self, request: "web.Request") -> "web.Respons
limit=limit,
offset=offset,
latest=latest_page,
include_compacted=include_compacted,
before_id=before_id,
)
return web.json_response({
"object": "list",
Expand All @@ -4548,6 +4565,7 @@ async def _handle_session_messages(self, request: "web.Request") -> "web.Respons
"offset": offset,
"order": order or ("latest" if default_page else "oldest"),
"returned": len(messages),
"next_before_id": messages[0]["id"] if messages else None,
},
})

Expand Down
8 changes: 8 additions & 0 deletions hermes_cli/web_routers/sessions.py
Original file line number Diff line number Diff line change
Expand Up @@ -612,12 +612,18 @@ async def get_session_messages(
offset: int = Query(0, ge=0),
order: Optional[str] = Query(None),
include_compacted: bool = Query(False),
before_id: Optional[int] = Query(None, ge=1),
):
if order not in (None, "oldest", "latest"):
raise HTTPException(
status_code=400,
detail="order must be one of: oldest, latest",
)
if before_id is not None and (order != "latest" or offset != 0):
raise HTTPException(
status_code=400,
detail="before_id requires order=latest and is incompatible with offset",
)

def _read():
db = _open_session_db_for_profile(profile, read_only=True)
Expand All @@ -640,6 +646,7 @@ def _read():
offset=offset,
latest=latest_page,
include_compacted=include_compacted,
before_id=before_id,
)
finally:
db.close()
Expand Down Expand Up @@ -676,6 +683,7 @@ def _read():
"offset": offset,
"order": order or ("latest" if limit is None else "oldest"),
"returned": len(projected_messages),
"next_before_id": messages[0]["id"] if messages else None,
},
}

Expand Down
23 changes: 20 additions & 3 deletions hermes_state.py
Original file line number Diff line number Diff line change
Expand Up @@ -11967,6 +11967,7 @@ def get_messages(
offset: int = 0,
latest: bool = False,
after_id: Optional[int] = None,
before_id: Optional[int] = None,
) -> List[Dict[str, Any]]:
"""Load messages for a session in insertion order.

Expand Down Expand Up @@ -11998,11 +11999,16 @@ def get_messages(
``after_id`` enables keyset pagination (``id > after_id``): O(1)
page seeks on huge transcripts where OFFSET degrades to O(n) per
page. Ascending order only (incompatible with ``latest``/``offset``).
``before_id`` is the descending counterpart (``id < before_id``),
used to page backwards from a prior latest page without a moving-tail
OFFSET. Pages remain chronological and new appends cannot shift them.
"""
if after_id is not None and (latest or offset):
raise ValueError("after_id is incompatible with latest/offset paging")
if after_id is not None and include_compacted:
raise ValueError("after_id is incompatible with include_compacted (deduped display reads use offset paging)")
if before_id is not None and (not latest or offset or after_id is not None):
raise ValueError("before_id requires latest paging and is incompatible with offset/after_id")
if include_inactive:
# Audit / debug reads: every row, including soft-deleted.
active_clause = ""
Expand All @@ -12013,7 +12019,9 @@ def get_messages(
active_clause = " AND (active = 1 OR compacted = 1)"
else:
active_clause = " AND active = 1"
keyset_clause = " AND id > ?" if after_id is not None else ""
keyset_clause = " AND id > ?" if after_id is not None else (
" AND id < ?" if before_id is not None else ""
)
sql = (
"SELECT * FROM messages WHERE session_id = ?"
f"{active_clause}{keyset_clause} ORDER BY id {'DESC' if latest else 'ASC'}"
Expand All @@ -12030,11 +12038,18 @@ def get_messages(
# full display set (a session's rows are bounded; the UI-level
# 500-row cap lives in the endpoint, not here), dedupe in Python,
# then apply paging.
# Cursor eligibility must be fixed before choosing the preferred
# copy. A compaction between pages can otherwise move an unseen
# logical message above before_id and make its older copy vanish.
display_cursor_clause = " AND id < ?" if before_id is not None else ""
display_params = [session_id]
if before_id is not None:
display_params.append(before_id)
with self._read_ctx() as conn:
cursor = conn.execute(
"SELECT * FROM messages WHERE session_id = ?" + active_clause
+ " ORDER BY id ASC",
[session_id],
+ display_cursor_clause + " ORDER BY id ASC",
display_params,
)
all_rows = cursor.fetchall()
seen: dict = {}
Expand Down Expand Up @@ -12080,6 +12095,8 @@ def get_messages(
if latest:
rows = rows[::-1]
else:
if before_id is not None:
params.append(before_id)
if limit is not None or offset:
# SQLite's OFFSET requires LIMIT; -1 means "no limit".
sql += " LIMIT ? OFFSET ?"
Expand Down
56 changes: 56 additions & 0 deletions tests/gateway/test_session_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,7 @@ async def test_session_messages_default_to_latest_bounded_page(adapter, session_
"offset": 0,
"order": "latest",
"returned": 500,
"next_before_id": payload["data"][0]["id"],
}
assert payload["data"][0]["content"] == "msg 1"
assert payload["data"][-1]["content"] == "msg 500"
Expand All @@ -116,6 +117,61 @@ async def test_session_messages_default_to_latest_bounded_page(adapter, session_
]


@pytest.mark.asyncio
async def test_session_messages_pages_compacted_history_from_latest(adapter, session_db):
session_db.create_session("long-chat", "api_server")
session_db.append_message("long-chat", "user", "old")
session_db.append_message("long-chat", "assistant", "old reply")
session_db.archive_and_compact(
"long-chat", [{"role": "system", "content": "summary"}]
)
session_db.append_message("long-chat", "user", "recent")
session_db.append_message("long-chat", "assistant", "recent reply")
rewound_id = session_db.append_message("long-chat", "user", "rewound")
session_db._conn.execute(
"UPDATE messages SET active = 0, compacted = 0 WHERE id = ?", (rewound_id,)
)
session_db._conn.commit()

app = _create_session_app(adapter)
async with TestClient(TestServer(app)) as cli:
response = await cli.get(
"/api/sessions/long-chat/messages"
"?include_compacted=true&order=latest&limit=2"
)
assert response.status == 200
payload = await response.json()
# Appending between requests used to move the offset-relative tail,
# duplicating one recent row and permanently skipping an older row.
session_db.append_message("long-chat", "assistant", "appended later")
earlier_response = await cli.get(
"/api/sessions/long-chat/messages"
"?include_compacted=true&order=latest&limit=3"
f"&before_id={payload['pagination']['next_before_id']}"
)
assert earlier_response.status == 200
earlier = await earlier_response.json()

assert [message["content"] for message in payload["data"]] == [
"recent",
"recent reply",
]
assert payload["pagination"] == {
"limit": 2,
"offset": 0,
"order": "latest",
"returned": 2,
"next_before_id": payload["data"][0]["id"],
}
assert [message["content"] for message in earlier["data"]] == [
"old",
"old reply",
"summary",
]
assert not ({message["id"] for message in payload["data"]} & {message["id"] for message in earlier["data"]})
assert "rewound" not in [message["content"] for message in payload["data"] + earlier["data"]]


@pytest.mark.asyncio
async def test_run_agent_binds_api_session_context_for_tool_env(adapter, monkeypatch):
"""API-server request sessions should reach tools and terminal subprocess env."""
Expand Down
1 change: 1 addition & 0 deletions tests/hermes_cli/test_web_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -2278,6 +2278,7 @@ def test_get_session_messages_omitted_limit_defaults_to_500(self):
"offset": 0,
"order": "latest",
"returned": 500,
"next_before_id": payload["messages"][0]["id"],
}
assert len(payload["messages"]) == 500
assert payload["messages"][0]["content"] == "msg 1"
Expand Down
58 changes: 58 additions & 0 deletions tests/hermes_state/test_get_messages_include_compacted.py
Original file line number Diff line number Diff line change
Expand Up @@ -206,6 +206,64 @@ def test_dedupe_applies_before_paging(self, db):
assert [m["id"] for m in page] == all_ids[2:]
assert len(page) == 2

def test_backward_cursor_survives_compaction_copying_unseen_messages(self, db):
"""A copy above the cursor must not hide its unseen pre-cursor row."""
sid = "cursor-compaction"
db.create_session(sid, source="cli")
db.append_messages_batch(
sid,
[
{"role": "user", "content": "q1"},
{"role": "assistant", "content": "a1"},
{"role": "user", "content": "q2"},
{"role": "assistant", "content": "a2"},
{"role": "user", "content": "q3"},
{"role": "assistant", "content": "a3"},
],
)
original = db.get_messages(sid)

first = db.get_messages(
sid, include_compacted=True, latest=True, limit=2
)
cursor = first[0]["id"]
assert [m["content"] for m in first] == ["q3", "a3"]

# Exercise the production transition: archive the old generation and
# publish a summary plus fresh active copies of two messages the client
# has not fetched yet. Their timestamps are preserved exactly, as in a
# protected-tail copy, so display projection recognizes them as the
# same logical messages.
copied = [
{
"role": message["role"],
"content": message["content"],
"timestamp": message["timestamp"],
}
for message in original[2:4]
]
db.archive_and_compact(
sid,
[
{
"role": "user",
"content": "summary of q1/a1",
"_compressed_summary": True,
},
*copied,
],
)

second = db.get_messages(
sid,
include_compacted=True,
latest=True,
before_id=cursor,
limit=2,
)
assert [m["content"] for m in second] == ["q2", "a2"]
assert len({m["content"] for m in first + second}) == 4

def test_distinct_tool_calls_with_same_content_are_not_merged(self, db):
"""Two real tool messages that happen to share role/content/timestamp
must stay separate: the dedupe key includes the tool fields, so only
Expand Down
11 changes: 11 additions & 0 deletions tests/test_hermes_state.py
Original file line number Diff line number Diff line change
Expand Up @@ -4539,6 +4539,17 @@ def test_latest_pages_count_back_from_newest_but_remain_chronological(self, db):
assert [m["content"] for m in page2] == ["msg-2", "msg-3", "msg-4", "msg-5"]
assert [m["content"] for m in page3] == ["msg-0", "msg-1"]

def test_before_id_latest_pages_are_stable_across_appends(self, db):
self._seed(db)
page1 = db.get_messages("s1", limit=4, latest=True)
cursor = page1[0]["id"]
db.append_message("s1", "user", "appended-after-page-1")
page2 = db.get_messages("s1", limit=4, latest=True, before_id=cursor)
page3 = db.get_messages("s1", limit=4, latest=True, before_id=page2[0]["id"])
contents = [m["content"] for m in page3 + page2 + page1]
assert contents == [f"msg-{i}" for i in range(10)]
assert len({m["id"] for m in page3 + page2 + page1}) == 10

def test_after_id_keyset_pages_forward_in_insertion_order(self, db):
self._seed(db)
page1 = db.get_messages("s1", limit=4, after_id=0)
Expand Down
Loading
Loading