perf(sessions): index prompts without hydrating full transcripts

This commit is contained in:
brooklyn!
2026-09-16 04:21:49 -07:00
parent a63d6af683
commit 7e69873ca6
4 changed files with 538 additions and 0 deletions

View File

@@ -564,6 +564,65 @@ async def get_session_messages(
"returned": len(projected_messages)}}
def _timeline_session_id(db, session_id: str, owner: str) -> str:
# Durable jump addresses are exact ids, never title/prefix guesses. A NULL
# legacy owner belongs to this profile's store, just like /messages pages.
def owned(sid):
row = db._read_one("SELECT profile_name FROM sessions WHERE id = ?", (sid,))
return row is not None and row["profile_name"] in (None, owner)
if not owned(session_id):
raise HTTPException(status_code=404, detail=_NOT_FOUND)
sid = db.resolve_resume_session_id(session_id)
if not owned(sid):
raise HTTPException(status_code=404, detail=_NOT_FOUND)
return sid
@manage_router.get("/api/sessions/{session_id}/timeline")
async def get_session_timeline(
session_id: str, profile: Optional[str] = None,
limit: int = Query(500, ge=1, le=500), after_row_id: int = Query(0, ge=0),
):
"""Prompt metadata only, including compacted display history (never rewind rows).
``next_cursor`` is a stable logical first-row id; pass it as ``after_row_id``.
Entry ``row_id`` addresses the current representative for /messages/around.
"""
from hermes_state_timeline import get_session_timeline as read_timeline
owner = _serving_profile(profile)
def _read(db):
sid = _timeline_session_id(db, session_id, owner)
return {"session_id": sid, "profile": owner,
**read_timeline(db, sid, limit=limit, after_row_id=after_row_id)}
return await asyncio.to_thread(_with_db, profile, _read, read_only=True)
@manage_router.get("/api/sessions/{session_id}/messages/around")
async def get_session_messages_around(
session_id: str, row_id: int = Query(..., ge=1), profile: Optional[str] = None,
limit: int = Query(120, ge=1, le=120),
):
"""Bounded display page starting at a timeline prompt; no intervening payloads."""
from hermes_state_timeline import get_session_messages_around as read_around
owner = _serving_profile(profile)
def _read(db):
sid = _timeline_session_id(db, session_id, owner)
page = read_around(db, sid, row_id, limit=limit)
if page is None:
raise HTTPException(status_code=404, detail="Prompt not found")
return {"session_id": sid, "profile": owner, **page}
result = await asyncio.to_thread(_with_db, profile, _read, read_only=True)
result["messages"] = _project_for_display(result["messages"])
return result
@manage_router.delete("/api/sessions/{session_id}")
async def delete_session_endpoint(session_id: str, profile: Optional[str] = None):
def _delete(db):

166
hermes_state_timeline.py Normal file
View File

@@ -0,0 +1,166 @@
"""Read-only prompt index and bounded transcript jumps; no transcript-wide payload hydration."""
from __future__ import annotations
import re
from contextlib import contextmanager
from agent.compaction_display import project_compaction_message_for_display
from agent.context_compressor import user_originated_turn_view
_SYNTHETIC_PROMPT = re.compile(
r"^\s*(?:\[IMPORTANT: Background process |\[ASYNC (?:DELEGATION )?(?:BATCH )?COMPLETE\b|"
r"A background fan-out of \d+ subagent\(s\) you dispatched earlier has finished\.|"
r"A background subagent you dispatched earlier has finished\.)",
re.IGNORECASE,
)
def _prompt_preview(db, content, display_kind, summary):
message = project_compaction_message_for_display({
"role": "user", "content": db._decode_content(content),
"display_kind": display_kind, "_compressed_summary": bool(summary),
})
if message is None or user_originated_turn_view(message) is None:
return ""
content = message.get("content")
if isinstance(content, list):
content = " ".join(
part if isinstance(part, str) else part.get("text", "")
for part in content if isinstance(part, (str, dict)))
if not isinstance(content, str):
return ""
text = " ".join(content.split())
if not text or _SYNTHETIC_PROMPT.match(text):
return ""
return text if len(text) <= 120 else text[:119].rstrip() + "…"
@contextmanager
def _snapshot(db):
# The count and page must see the same compaction/rewind generation.
with db._read_ctx() as conn:
conn.execute("BEGIN")
try:
yield conn
finally:
if conn.in_transaction:
conn.execute("ROLLBACK")
def _display_rows_sql(conn, session_id, *, users_only=False):
"""Return only representative ids and their stable first-row order, never bodies.
Legacy stores cannot backfill on a GET. SQL groups their payload identities in
SQLite; only user content crosses the Python boundary for carrier normalization.
Current stores use the durable display index, including protected-tail copies.
"""
role = " AND role = 'user'" if users_only else ""
indexed = conn.execute(
"SELECT 1 FROM messages WHERE session_id = ? AND (active = 1 OR compacted = 1) "
f"{role} AND (display_order IS NULL OR display_identity IS NULL) LIMIT 1",
(session_id,),
).fetchone() is None
if indexed:
return f"""WITH display_rows AS (
SELECT (SELECT candidate.id FROM messages candidate
WHERE candidate.session_id = :sid
AND candidate.display_order = m.display_order
AND (candidate.active = 1 OR candidate.compacted = 1)
ORDER BY candidate.active DESC, candidate.id DESC LIMIT 1) AS row_id,
m.display_order AS sort_id
FROM messages m WHERE session_id = :sid AND (active = 1 OR compacted = 1){role}
GROUP BY m.display_order
)"""
return f"""WITH ranked AS (
SELECT id, MIN(id) OVER identity AS sort_id,
ROW_NUMBER() OVER (identity ORDER BY active DESC, id DESC) AS preference
FROM messages WHERE session_id = :sid AND (active = 1 OR compacted = 1){role}
WINDOW identity AS (PARTITION BY role,
CASE WHEN role = 'user' THEN timeline_identity_content(content, display_kind) ELSE content END,
timestamp, tool_call_id, tool_calls, tool_name)
), display_rows AS (SELECT id AS row_id, sort_id FROM ranked WHERE preference = 1)"""
def _register_functions(db, conn):
from agent.context_compressor import split_user_originated_turn
def identity_content(content, display_kind):
handoff, live = split_user_originated_turn({
"role": "user", "content": db._decode_content(content), "display_kind": display_kind})
return db._encode_content(live.get("content")) if handoff is not None and live is not None else content
conn.create_function("timeline_identity_content", 2, identity_content, deterministic=True)
conn.create_function("timeline_preview", 3,
lambda content, kind, summary: _prompt_preview(db, content, kind, summary),
deterministic=True)
def get_session_messages_around(db, session_id, row_id, *, limit=120):
"""Read at most *limit* display rows starting at an exact human prompt.
Existence probes/counts contain ids only. Full payloads are fetched only for
the selected bounded page, even when the anchor is deep in a transcript.
"""
with _snapshot(db) as conn:
_register_functions(db, conn)
anchor = conn.execute(
"SELECT content, display_kind, _compressed_summary FROM messages "
"WHERE session_id = ? AND id = ? AND role = 'user' AND (active = 1 OR compacted = 1)",
(session_id, row_id),
).fetchone()
if anchor is None or not _prompt_preview(db, *anchor):
return None
sql = _display_rows_sql(conn, session_id)
params = {"sid": session_id, "row_id": row_id, "limit": limit}
selected = conn.execute(sql + """
SELECT sort_id FROM display_rows WHERE row_id = :row_id
""", params).fetchone()
if selected is None:
return None
params["start"] = selected["sort_id"]
counts = conn.execute(sql + """
SELECT COUNT(*) AS total, COALESCE(SUM(sort_id < :start), 0) AS offset FROM display_rows
""", params).fetchone()
rows = conn.execute(sql + """
SELECT m.* FROM (SELECT row_id, sort_id FROM display_rows
WHERE sort_id >= :start ORDER BY sort_id LIMIT :limit) AS page
JOIN messages m ON m.id = page.row_id ORDER BY page.sort_id
""", params).fetchall()
messages = [db._row_to_message_dict(row, warn_context="timeline jump", summary_flag=True) for row in rows]
return {"messages": messages, "pagination": {
"row_id": row_id, "limit": limit, "returned": len(messages), "order": "oldest",
"offset": counts["offset"], "total": counts["total"],
"has_older": counts["offset"] > 0,
"has_newer": counts["offset"] + len(messages) < counts["total"],
}}
def get_session_timeline(db, session_id, *, limit=500, after_row_id=0):
"""Chronological prompts. Cursor is the first physical row id of a logical turn."""
with _snapshot(db) as conn:
_register_functions(db, conn)
sql = _display_rows_sql(conn, session_id, users_only=True) + """,
prompts AS MATERIALIZED (
SELECT row_id, sort_id, m.timestamp,
timeline_preview(m.content, m.display_kind, m._compressed_summary) AS preview
FROM display_rows JOIN messages m ON m.id = row_id
), eligible AS MATERIALIZED (SELECT * FROM prompts WHERE preview <> '')
"""
params = {"sid": session_id, "after": after_row_id, "limit": limit + 1}
rows = conn.execute(sql + """
SELECT row_id, sort_id, timestamp, preview, (SELECT COUNT(*) FROM eligible) AS total
FROM eligible WHERE sort_id > :after ORDER BY sort_id LIMIT :limit
""", params).fetchall()
total = rows[0]["total"] if rows else conn.execute(
sql + "SELECT COUNT(*) FROM eligible", params).fetchone()[0]
has_more = len(rows) > limit
page = rows[:limit]
return {
"entries": [{"row_id": row["row_id"], "preview": row["preview"], "timestamp": row["timestamp"]}
for row in page],
"pagination": {"limit": limit, "after_row_id": after_row_id, "returned": len(page),
"total": total, "has_more": has_more,
"next_cursor": page[-1]["sort_id"] if has_more else None},
}

View File

@@ -0,0 +1,99 @@
"""Reproducible HTTP payload benchmark using generated data, never a live store.
Run with the development Python from the repository root. Prints a JSON receipt;
fixture setup is excluded from timings. No provider, credential, or app access.
"""
import json
import os
from pathlib import Path
import statistics
import sys
import tempfile
import time
sys.path.insert(0, str(Path(__file__).resolve().parents[2]))
def main():
with tempfile.TemporaryDirectory(prefix="hermes-timeline-bench-") as directory:
os.environ["HERMES_HOME"] = directory
from fastapi import FastAPI
from fastapi.testclient import TestClient
from hermes_state import SessionDB
from hermes_cli.web_routers.sessions import manage_router
sid = "generated-tool-heavy"
prompt_count = 600
rows_per_turn = 10
with SessionDB(db_path=Path(directory) / "state.db") as db:
db.create_session(session_id=sid, source="desktop")
for start in range(0, prompt_count, 100):
batch = []
for turn in range(start, start + 100):
batch.append({"role": "user", "content": f"Investigate generated task {turn}", "timestamp": turn + 1})
for tool in range(4):
call_id = f"call-{turn}-{tool}"
batch.append({"role": "assistant", "content": "", "tool_calls": [{
"id": call_id, "type": "function", "function": {
"name": "terminal", "arguments": json.dumps({"command": "x" * 4096})}}]})
batch.append({"role": "tool", "tool_call_id": call_id,
"content": "generated tool output\n" * 800})
batch.append({"role": "assistant", "content": f"Finished task {turn}"})
db.append_messages_batch(sid, batch)
app = FastAPI()
app.include_router(manage_router)
with TestClient(app) as client:
def full_messages():
total_bytes = count = requests = 0
for offset in range(0, prompt_count * rows_per_turn + 1, 500):
response = client.get(f"/api/sessions/{sid}/messages", params={
"limit": 500, "offset": offset, "order": "oldest", "include_compacted": True})
response.raise_for_status()
page = response.json()["messages"]
total_bytes += len(response.content)
count += len(page)
requests += 1
if len(page) < 500:
break
assert count == prompt_count * rows_per_turn
return {"bytes": total_bytes, "rows": count, "requests": requests}
def timeline():
total_bytes = count = requests = 0
cursor = 0
seen = set()
while True:
response = client.get(f"/api/sessions/{sid}/timeline", params={"limit": 500, "after_row_id": cursor})
response.raise_for_status()
page = response.json()
total_bytes += len(response.content)
count += len(page["entries"])
requests += 1
seen.update(entry["row_id"] for entry in page["entries"])
if not page["pagination"]["has_more"]:
break
cursor = page["pagination"]["next_cursor"]
assert count == len(seen) == prompt_count
return {"bytes": total_bytes, "rows": count, "requests": requests}
samples = {"full_messages": [], "timeline": []}
for _ in range(3):
for name, procedure in (("full_messages", full_messages), ("timeline", timeline)):
started = time.perf_counter()
result = procedure()
result["elapsed_ms"] = (time.perf_counter() - started) * 1000
samples[name].append(result)
receipt = {
"fixture": {"prompts": prompt_count, "message_rows": prompt_count * rows_per_turn,
"tool_results": prompt_count * 4},
**{name: {**runs[-1], "elapsed_ms": statistics.median(r["elapsed_ms"] for r in runs),
"samples_ms": [r["elapsed_ms"] for r in runs]} for name, runs in samples.items()},
}
receipt["payload_reduction_percent"] = 100 * (1 - receipt["timeline"]["bytes"] / receipt["full_messages"]["bytes"])
receipt["speedup"] = receipt["full_messages"]["elapsed_ms"] / receipt["timeline"]["elapsed_ms"]
print(json.dumps(receipt, indent=2))
if __name__ == "__main__":
main()

View File

@@ -0,0 +1,214 @@
"""Real read-only timeline requests against temporary profile stores."""
from pathlib import Path
import pytest
from fastapi import FastAPI
from fastapi.testclient import TestClient
from hermes_state import SessionDB
@pytest.fixture
def timeline_store(tmp_path, monkeypatch):
monkeypatch.setattr(Path, "home", lambda: tmp_path)
home = tmp_path / ".hermes"
home.mkdir()
monkeypatch.setenv("HERMES_HOME", str(home))
monkeypatch.setattr("hermes_state.DEFAULT_DB_PATH", home / "state.db")
db = SessionDB(db_path=home / "state.db")
db.create_session(session_id="timeline-root", source="desktop")
from hermes_cli.web_routers.sessions import manage_router
app = FastAPI()
app.include_router(manage_router)
with TestClient(app) as client:
yield db, client, home
db.close()
def test_timeline_pages_project_human_prompts_without_tool_payloads(timeline_store):
db, client, _ = timeline_store
db.append_messages_batch("timeline-root", [
{"role": "user", "content": " First\n prompt ", "timestamp": 200},
{"role": "assistant", "content": "answer", "tool_calls": [
{"id": "tool-1", "function": {"name": "terminal", "arguments": "x" * 10000}}]},
{"role": "tool", "content": "secret tool payload" * 1000, "tool_call_id": "tool-1"},
{"role": "user", "content": "é" * 150, "timestamp": 100},
])
expected_ids = [row["id"] for row in db.get_messages("timeline-root") if row["role"] == "user"]
response = client.get("/api/sessions/timeline-root/timeline?limit=1")
assert response.status_code == 200
first = response.json()
assert first["session_id"] == "timeline-root"
assert first["profile"] == "default"
assert first["entries"] == [{"row_id": expected_ids[0], "preview": "First prompt", "timestamp": 200}]
assert first["pagination"]["total"] == 2
assert first["pagination"]["returned"] == 1
assert first["pagination"]["has_more"] is True
second = client.get("/api/sessions/timeline-root/timeline", params={
"after_row_id": first["pagination"]["next_cursor"], "limit": 1}).json()
assert second["entries"][0]["row_id"] == expected_ids[1]
assert len(second["entries"][0]["preview"]) == 120
assert second["pagination"]["has_more"] is False
assert second["pagination"]["next_cursor"] is None
assert b"secret tool payload" not in response.content
@pytest.mark.parametrize("legacy", [False, True])
def test_compacted_timeline_and_jump_share_display_order_and_visibility(timeline_store, legacy):
from agent.context_compressor import COMPRESSION_CONTINUATION_USER_CONTENT, SUMMARY_PREFIX, _SUMMARY_END_MARKER
db, client, _ = timeline_store
sid = "timeline-root"
db.append_messages_batch(sid, [
{"role": "user", "content": "old ask", "timestamp": 10},
{"role": "assistant", "content": "old answer", "timestamp": 11},
{"role": "user", "content": "retained ask", "timestamp": 12},
{"role": "assistant", "content": "retained answer", "timestamp": 13},
])
retained = db.get_messages(sid)[2:]
db.archive_and_compact(sid, [
{"role": "user", "content": SUMMARY_PREFIX + "summary", "_compressed_summary": True},
*retained,
])
db.append_messages_batch(sid, [
{"role": "user", "content": "hidden ask", "display_kind": "hidden"},
{"role": "user", "content": COMPRESSION_CONTINUATION_USER_CONTENT},
{"role": "user", "content": "[IMPORTANT: Background process 12 completed]"},
{"role": "user", "content": "[ASYNC DELEGATION BATCH COMPLETE] completed"},
{"role": "user", "content": "A background subagent you dispatched earlier has finished."},
{"role": "user", "content": SUMMARY_PREFIX + "handoff\n" + _SUMMARY_END_MARKER + "\nlive ask",
"_compressed_summary": True, "display_kind": "hidden"},
{"role": "assistant", "content": "live answer"},
{"role": "user", "content": "rewound ask"},
])
rewound_id = db.get_messages(sid)[-1]["id"]
db._write_sql("UPDATE messages SET active = 0, compacted = 0 WHERE id = ?", (rewound_id,))
if legacy:
db._write_sql("UPDATE messages SET display_identity = NULL, display_order = NULL WHERE session_id = ?", (sid,))
timeline = client.get(f"/api/sessions/{sid}/timeline").json()
assert [entry["preview"] for entry in timeline["entries"]] == ["old ask", "retained ask", "live ask"]
anchor = timeline["entries"][1]["row_id"]
response = client.get(f"/api/sessions/{sid}/messages/around?row_id={anchor}&limit=2")
assert response.status_code == 200
jump = response.json()
assert [row["content"] for row in jump["messages"]] == ["retained ask", "retained answer"]
assert jump["pagination"]["has_older"] is True
assert jump["pagination"]["has_newer"] is True
live_id = timeline["entries"][-1]["row_id"]
last = client.get(f"/api/sessions/{sid}/messages/around?row_id={live_id}").json()
assert last["messages"][0]["display_content"] == "live ask"
assert last["messages"][1]["content"] == "live answer"
assert last["pagination"]["has_newer"] is False
assert client.get(f"/api/sessions/{sid}/messages/around?row_id={rewound_id}").status_code == 404
def test_cursor_survives_compaction_between_pages(timeline_store):
db, client, _ = timeline_store
sid = "timeline-root"
db.append_messages_batch(sid, [
{"role": "user", "content": f"ask {i}", "timestamp": 100 - i} for i in range(4)
])
first = client.get(f"/api/sessions/{sid}/timeline?limit=2").json()
old_row_id = first["entries"][-1]["row_id"]
retained = db.get_messages(sid)[1:]
db.archive_and_compact(sid, retained)
second = client.get(f"/api/sessions/{sid}/timeline", params={
"limit": 2, "after_row_id": first["pagination"]["next_cursor"]}).json()
assert [r["preview"] for r in first["entries"] + second["entries"]] == [f"ask {i}" for i in range(4)]
assert second["pagination"]["total"] == 4
assert second["pagination"]["has_more"] is False
current = client.get(f"/api/sessions/{sid}/timeline").json()
assert current["entries"][1]["row_id"] > old_row_id
empty = client.get(f"/api/sessions/{sid}/timeline?after_row_id=99999").json()
assert empty["entries"] == []
assert empty["pagination"]["total"] == 4
assert empty["pagination"]["has_more"] is False
def test_exact_owner_lineage_validation_and_bounded_jump(timeline_store):
db, client, home = timeline_store
sid = "timeline-root"
db.end_session(sid, end_reason="compression")
db.create_session(session_id="timeline-tip", source="desktop", parent_session_id=sid)
db.append_messages_batch("timeline-tip", [
{"role": "user", "content": "tip ask"},
*[{"role": "assistant", "content": f"step {i}"} for i in range(130)],
{"role": "user", "content": "last ask"},
{"role": "assistant", "content": "last answer"},
])
db.create_session(session_id="delegate", source="tool", parent_session_id="timeline-tip")
db.append_message("delegate", role="user", content="delegate ask")
db.create_session(session_id="foreign", source="desktop")
foreign = db.append_message("foreign", role="user", content="foreign ask")
expected_sid = client.get(f"/api/sessions/{sid}/messages").json()["session_id"]
first = client.get(f"/api/sessions/{sid}/timeline?limit=1").json()
assert first["session_id"] == expected_sid == "timeline-tip"
row_id = first["entries"][0]["row_id"]
jump = client.get(f"/api/sessions/{sid}/messages/around?row_id={row_id}").json()
assert jump["messages"][0]["id"] == row_id
assert len(jump["messages"]) == jump["pagination"]["returned"] == 120
assert jump["pagination"]["has_older"] is False
assert jump["pagination"]["has_newer"] is True
assert client.get(f"/api/sessions/{sid}/messages/around?row_id={foreign}").status_code == 404
assert client.get(f"/api/sessions/{sid}/messages/around?row_id={row_id + 1}").status_code == 404
assert client.get(f"/api/sessions/{sid}/messages/around?row_id={row_id}&limit=121").status_code == 422
assert client.get(f"/api/sessions/{sid}/timeline?limit=501").status_code == 422
assert client.get("/api/sessions/timeline-ro/timeline").status_code == 404
assert client.get("/api/sessions/absent/timeline").status_code == 404
work = home / "profiles" / "work"
work.mkdir(parents=True)
with SessionDB(db_path=work / "state.db") as other:
other.create_session(session_id=sid, source="desktop")
other.append_message(sid, role="user", content="work ask")
for profile, expected in (("default", "tip ask"), ("work", "work ask"), ("default", "tip ask")):
page = client.get(f"/api/sessions/{sid}/timeline?profile={profile}").json()
assert page["profile"] == profile
assert page["entries"][0]["preview"] == expected
db._write_sql("UPDATE sessions SET profile_name = 'wrong-owner' WHERE id = 'timeline-tip'")
assert client.get(f"/api/sessions/{sid}/timeline").status_code == 404
assert client.get(f"/api/sessions/{sid}/messages/around?row_id={row_id}").status_code == 404
def test_timeline_sql_never_reads_tool_columns_or_writes(timeline_store, monkeypatch):
import sqlite3
from contextlib import contextmanager
from hermes_state_timeline import get_session_timeline
db, _, home = timeline_store
db.append_messages_batch("timeline-root", [
{"role": "user", "content": [{"type": "text", "text": "multimodal ask"},
{"type": "image_url", "image_url": {"url": "data:unused"}}]},
{"role": "tool", "content": "unread tool result", "tool_calls": [{"unused": "payload"}]},
])
reader = SessionDB(db_path=home / "state.db", read_only=True)
original = reader._read_ctx
accesses = []
def authorize(action, table, column, *_):
if action == sqlite3.SQLITE_READ:
accesses.append((table, column))
if table == "messages" and column in {"tool_calls", "reasoning", "api_content", "codex_reasoning_items"}:
return sqlite3.SQLITE_DENY
if action in {sqlite3.SQLITE_UPDATE, sqlite3.SQLITE_INSERT, sqlite3.SQLITE_DELETE}:
return sqlite3.SQLITE_DENY
return sqlite3.SQLITE_OK
@contextmanager
def guarded():
with original() as conn:
conn.set_authorizer(authorize)
try:
yield conn
finally:
conn.set_authorizer(None)
monkeypatch.setattr(reader, "_read_ctx", guarded)
try:
page = get_session_timeline(reader, "timeline-root")
assert page["entries"][0]["preview"] == "multimodal ask"
assert page["pagination"]["total"] == 1
assert accesses
finally:
reader.close()