From 5f60194f9f26c45c7bafdf94b550bd6a62b9de7b Mon Sep 17 00:00:00 2001 From: Rod Boev Date: Thu, 4 Jun 2026 06:53:36 -0400 Subject: [PATCH 1/3] feat(gateway): paginate /api/sessions/{id}/messages and /api/jobs (#38370) --- gateway/platforms/api_server.py | 78 +++++++++------------------ tests/gateway/test_api_server_jobs.py | 51 ++++++++++++++++++ tests/gateway/test_session_api.py | 47 ++++++++++++++++ 3 files changed, 122 insertions(+), 54 deletions(-) diff --git a/gateway/platforms/api_server.py b/gateway/platforms/api_server.py index 5ba09d67492e..7e8ed77cb508 100644 --- a/gateway/platforms/api_server.py +++ b/gateway/platforms/api_server.py @@ -64,29 +64,6 @@ logger = logging.getLogger(__name__) - -def _hermes_version() -> str: - """Return the hermes-agent version string, or "dev" if it can't be resolved. - - Tries the installed package metadata first (authoritative for a pip/uv - install), then the in-tree ``hermes_cli.__version__`` (covers editable / - source checkouts where metadata may be stale or absent). Never raises — - a version probe must not be able to break the health endpoint. - """ - try: - from importlib.metadata import version - - return version("hermes-agent") - except Exception: - pass - try: - from hermes_cli import __version__ - - return __version__ - except Exception: - return "dev" - - # Default settings DEFAULT_HOST = "127.0.0.1" DEFAULT_PORT = 8642 @@ -455,19 +432,7 @@ def get(self, response_id: str) -> Optional[Dict[str, Any]]: (time.time(), response_id), ) self._conn.commit() - try: - return json.loads(row[0]) - except (json.JSONDecodeError, TypeError): - logger.warning( - "Corrupted JSON in response store for id=%s, evicting entry", - response_id, - ) - self._conn.execute( - "DELETE FROM responses WHERE response_id = ?", - (response_id,), - ) - self._conn.commit() - return None + return json.loads(row[0]) def put(self, response_id: str, data: Dict[str, Any]) -> None: """Store a response, evicting the oldest if at capacity.""" @@ -800,7 +765,6 @@ def _derive_chat_session_id( _cron_resume = None _cron_trigger = None - def _notify_cron_provider_jobs_changed() -> None: """Tell the active cron scheduler provider the job set changed after a REST mutation (no-op for the built-in). Best-effort — never breaks the handler.""" @@ -823,7 +787,6 @@ def _notify_cron_provider_jobs_changed() -> None: except Exception: # pragma: no cover - scanner is optional hardening _scan_cron_prompt = None - class APIServerAdapter(BasePlatformAdapter): """ OpenAI-compatible HTTP API server adapter. @@ -1369,9 +1332,7 @@ def _create_agent( async def _handle_health(self, request: "web.Request") -> "web.Response": """GET /health — simple health check.""" - return web.json_response( - {"status": "ok", "platform": "hermes-agent", "version": _hermes_version()} - ) + return web.json_response({"status": "ok", "platform": "hermes-agent"}) async def _handle_health_detailed(self, request: "web.Request") -> "web.Response": """GET /health/detailed — rich status for cross-container dashboard probing. @@ -1818,12 +1779,19 @@ async def _handle_session_messages(self, request: "web.Request") -> "web.Respons if err: return err db = self._ensure_session_db() - resolved_id = db.resolve_resume_session_id(session_id) - messages = db.get_messages(resolved_id) + limit = self._parse_nonnegative_int(request.query.get("limit"), default=200, maximum=1000) + offset = self._parse_nonnegative_int(request.query.get("offset"), default=0, maximum=1_000_000) + messages = db.get_messages(session_id) + total = len(messages) + page = messages[offset:offset + limit] return web.json_response({ "object": "list", - "session_id": resolved_id, - "data": [self._message_response(m) for m in messages], + "session_id": session_id, + "data": [self._message_response(m) for m in page], + "limit": limit, + "offset": offset, + "total": total, + "has_more": offset + len(page) < total, }) async def _handle_fork_session(self, request: "web.Request") -> "web.Response": @@ -3561,8 +3529,18 @@ async def _handle_list_jobs(self, request: "web.Request") -> "web.Response": return cron_err try: include_disabled = request.query.get("include_disabled", "").lower() in {"true", "1"} + limit = self._parse_nonnegative_int(request.query.get("limit"), default=200, maximum=1000) + offset = self._parse_nonnegative_int(request.query.get("offset"), default=0, maximum=1_000_000) jobs = _cron_list(include_disabled=include_disabled) - return web.json_response({"jobs": jobs}) + total = len(jobs) + page = jobs[offset:offset + limit] + return web.json_response({ + "jobs": page, + "limit": limit, + "offset": offset, + "total": total, + "has_more": offset + len(page) < total, + }) except Exception as e: return web.json_response({"error": _redact_api_error_text(e)}, status=500) @@ -3595,10 +3573,6 @@ async def _handle_create_job(self, request: "web.Request") -> "web.Response": return web.json_response( {"error": f"Prompt must be ≤ {self._MAX_PROMPT_LENGTH} characters"}, status=400, ) - if prompt and _scan_cron_prompt is not None: - scan_error = _scan_cron_prompt(prompt) - if scan_error: - return web.json_response({"error": scan_error}, status=400) if repeat is not None and (not isinstance(repeat, int) or repeat < 1): return web.json_response({"error": "Repeat must be a positive integer"}, status=400) @@ -3665,10 +3639,6 @@ async def _handle_update_job(self, request: "web.Request") -> "web.Response": return web.json_response( {"error": f"Prompt must be ≤ {self._MAX_PROMPT_LENGTH} characters"}, status=400, ) - if sanitized.get("prompt") and _scan_cron_prompt is not None: - scan_error = _scan_cron_prompt(sanitized["prompt"]) - if scan_error: - return web.json_response({"error": scan_error}, status=400) job = _cron_update(job_id, sanitized) if not job: return web.json_response({"error": "Job not found"}, status=404) diff --git a/tests/gateway/test_api_server_jobs.py b/tests/gateway/test_api_server_jobs.py index 082ab6cf1671..d364795a99ae 100644 --- a/tests/gateway/test_api_server_jobs.py +++ b/tests/gateway/test_api_server_jobs.py @@ -130,6 +130,57 @@ async def test_list_jobs_default_excludes_disabled(self, adapter): assert resp.status == 200 mock_list.assert_called_once_with(include_disabled=False) + @pytest.mark.asyncio + async def test_list_jobs_pagination_slices_and_reports_total(self, adapter): + """GET /api/jobs with limit/offset returns the correct page and metadata.""" + all_jobs = [{**SAMPLE_JOB, "id": f"aabbccddeef{i}", "name": f"job-{i}"} for i in range(5)] + app = _create_app(adapter) + async with TestClient(TestServer(app)) as cli: + with patch(f"{_MOD}._CRON_AVAILABLE", True), patch( + f"{_MOD}._cron_list", return_value=all_jobs + ): + # First page + resp = await cli.get("/api/jobs?limit=2&offset=0") + assert resp.status == 200 + data = await resp.json() + assert data["jobs"] == all_jobs[:2] + assert data["limit"] == 2 + assert data["offset"] == 0 + assert data["total"] == 5 + assert data["has_more"] is True + + # Second page + resp = await cli.get("/api/jobs?limit=2&offset=2") + assert resp.status == 200 + data = await resp.json() + assert data["jobs"] == all_jobs[2:4] + assert data["has_more"] is True + + # Last page (1 item) + resp = await cli.get("/api/jobs?limit=2&offset=4") + assert resp.status == 200 + data = await resp.json() + assert data["jobs"] == all_jobs[4:] + assert data["has_more"] is False + + @pytest.mark.asyncio + async def test_list_jobs_default_limit_returns_all(self, adapter): + """GET /api/jobs with no limit param returns all jobs with has_more=False.""" + all_jobs = [{**SAMPLE_JOB, "id": f"aabbccddeef{i}", "name": f"job-{i}"} for i in range(3)] + app = _create_app(adapter) + async with TestClient(TestServer(app)) as cli: + with patch(f"{_MOD}._CRON_AVAILABLE", True), patch( + f"{_MOD}._cron_list", return_value=all_jobs + ): + resp = await cli.get("/api/jobs") + assert resp.status == 200 + data = await resp.json() + assert data["jobs"] == all_jobs + assert data["limit"] == 200 + assert data["offset"] == 0 + assert data["total"] == 3 + assert data["has_more"] is False + # --------------------------------------------------------------------------- # 3-7. test_create_job and validation diff --git a/tests/gateway/test_session_api.py b/tests/gateway/test_session_api.py index 47f7b38eec44..bbddbf965ad7 100644 --- a/tests/gateway/test_session_api.py +++ b/tests/gateway/test_session_api.py @@ -409,6 +409,53 @@ async def fake_run(**kwargs): +@pytest.mark.asyncio +async def test_session_messages_pagination(adapter, session_db): + session_id = session_db.create_session("page-session", "api_server") + for i in range(5): + session_db.append_message(session_id, "user", f"msg {i}") + + app = _create_session_app(adapter) + async with TestClient(TestServer(app)) as cli: + # Default: all 5, limit=200, offset=0, has_more=False + resp = await cli.get(f"/api/sessions/{session_id}/messages") + assert resp.status == 200 + data = await resp.json() + assert data["object"] == "list" + assert len(data["data"]) == 5 + assert data["limit"] == 200 + assert data["offset"] == 0 + assert data["total"] == 5 + assert data["has_more"] is False + + # First page of 2, has_more=True + resp = await cli.get(f"/api/sessions/{session_id}/messages?limit=2&offset=0") + assert resp.status == 200 + data = await resp.json() + assert len(data["data"]) == 2 + assert data["limit"] == 2 + assert data["offset"] == 0 + assert data["total"] == 5 + assert data["has_more"] is True + + # Last page: offset=4, limit=2, 1 item, has_more=False + resp = await cli.get(f"/api/sessions/{session_id}/messages?limit=2&offset=4") + assert resp.status == 200 + data = await resp.json() + assert len(data["data"]) == 1 + assert data["offset"] == 4 + assert data["total"] == 5 + assert data["has_more"] is False + + # Malformed params fall back to defaults, no 400 + resp = await cli.get(f"/api/sessions/{session_id}/messages?limit=bad&offset=-1") + assert resp.status == 200 + data = await resp.json() + assert data["limit"] == 200 + assert data["offset"] == 0 + assert len(data["data"]) == 5 + + @pytest.mark.asyncio async def test_session_endpoints_require_auth_when_key_configured(auth_adapter): app = _create_session_app(auth_adapter) From 64fb509c14c9f83dd1b61154f6214e198859ecd6 Mon Sep 17 00:00:00 2001 From: Rod Boev Date: Fri, 12 Jun 2026 07:37:23 -0400 Subject: [PATCH 2/3] fix(gateway): keep api_server behavior while paginating (#38370) --- gateway/platforms/api_server.py | 47 ++++++++++++++++++++++++++++++--- 1 file changed, 43 insertions(+), 4 deletions(-) diff --git a/gateway/platforms/api_server.py b/gateway/platforms/api_server.py index 7e8ed77cb508..e2159e59d406 100644 --- a/gateway/platforms/api_server.py +++ b/gateway/platforms/api_server.py @@ -64,6 +64,22 @@ logger = logging.getLogger(__name__) + +def _hermes_version() -> str: + """Return the hermes-agent version string, or "dev" if it can't be resolved.""" + try: + from importlib.metadata import version + + return version("hermes-agent") + except Exception: + pass + try: + from hermes_cli import __version__ + + return __version__ + except Exception: + return "dev" + # Default settings DEFAULT_HOST = "127.0.0.1" DEFAULT_PORT = 8642 @@ -432,7 +448,19 @@ def get(self, response_id: str) -> Optional[Dict[str, Any]]: (time.time(), response_id), ) self._conn.commit() - return json.loads(row[0]) + try: + return json.loads(row[0]) + except (json.JSONDecodeError, TypeError): + logger.warning( + "Corrupted JSON in response store for id=%s, evicting entry", + response_id, + ) + self._conn.execute( + "DELETE FROM responses WHERE response_id = ?", + (response_id,), + ) + self._conn.commit() + return None def put(self, response_id: str, data: Dict[str, Any]) -> None: """Store a response, evicting the oldest if at capacity.""" @@ -1332,7 +1360,9 @@ def _create_agent( async def _handle_health(self, request: "web.Request") -> "web.Response": """GET /health — simple health check.""" - return web.json_response({"status": "ok", "platform": "hermes-agent"}) + return web.json_response( + {"status": "ok", "platform": "hermes-agent", "version": _hermes_version()} + ) async def _handle_health_detailed(self, request: "web.Request") -> "web.Response": """GET /health/detailed — rich status for cross-container dashboard probing. @@ -1779,14 +1809,15 @@ async def _handle_session_messages(self, request: "web.Request") -> "web.Respons if err: return err db = self._ensure_session_db() + resolved_id = db.resolve_resume_session_id(session_id) limit = self._parse_nonnegative_int(request.query.get("limit"), default=200, maximum=1000) offset = self._parse_nonnegative_int(request.query.get("offset"), default=0, maximum=1_000_000) - messages = db.get_messages(session_id) + messages = db.get_messages(resolved_id) total = len(messages) page = messages[offset:offset + limit] return web.json_response({ "object": "list", - "session_id": session_id, + "session_id": resolved_id, "data": [self._message_response(m) for m in page], "limit": limit, "offset": offset, @@ -3573,6 +3604,10 @@ async def _handle_create_job(self, request: "web.Request") -> "web.Response": return web.json_response( {"error": f"Prompt must be ≤ {self._MAX_PROMPT_LENGTH} characters"}, status=400, ) + if prompt and _scan_cron_prompt is not None: + scan_error = _scan_cron_prompt(prompt) + if scan_error: + return web.json_response({"error": scan_error}, status=400) if repeat is not None and (not isinstance(repeat, int) or repeat < 1): return web.json_response({"error": "Repeat must be a positive integer"}, status=400) @@ -3639,6 +3674,10 @@ async def _handle_update_job(self, request: "web.Request") -> "web.Response": return web.json_response( {"error": f"Prompt must be ≤ {self._MAX_PROMPT_LENGTH} characters"}, status=400, ) + if sanitized.get("prompt") and _scan_cron_prompt is not None: + scan_error = _scan_cron_prompt(sanitized["prompt"]) + if scan_error: + return web.json_response({"error": scan_error}, status=400) job = _cron_update(job_id, sanitized) if not job: return web.json_response({"error": "Job not found"}, status=404) From b4adf0dcb5392a8fcfc43a7a0fa6e629a278a72a Mon Sep 17 00:00:00 2001 From: Rod Boev Date: Mon, 13 Jul 2026 21:15:49 -0400 Subject: [PATCH 3/3] docs(gateway): document paginated API server responses --- website/docs/user-guide/features/api-server.md | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/website/docs/user-guide/features/api-server.md b/website/docs/user-guide/features/api-server.md index 4f1db5ab0c26..1e15d5718ab7 100644 --- a/website/docs/user-guide/features/api-server.md +++ b/website/docs/user-guide/features/api-server.md @@ -282,7 +282,7 @@ The server exposes a lightweight jobs CRUD surface for managing scheduled / back ### GET /api/jobs -List all scheduled jobs. +List scheduled jobs. Supports `limit`, `offset`, and `include_disabled=1|true`. The response includes the current page in `jobs`, plus `limit`, `offset`, `total`, and `has_more` so clients can keep paginating without re-counting locally. ### POST /api/jobs @@ -323,13 +323,15 @@ External UIs can manage Hermes sessions over REST without standing up the dashbo | `GET` | `/api/sessions/{id}` | Read session metadata | | `PATCH` | `/api/sessions/{id}` | Update title or `end_reason` | | `DELETE` | `/api/sessions/{id}` | Delete a session | -| `GET` | `/api/sessions/{id}/messages` | Message history for a session | +| `GET` | `/api/sessions/{id}/messages` | Message history for a session, paginated with `limit` and `offset` | | `POST` | `/api/sessions/{id}/fork` | Branch the session via `SessionDB` lineage (matches CLI `/branch` semantics) | | `POST` | `/api/sessions/{id}/chat` | Run one synchronous agent turn | | `POST` | `/api/sessions/{id}/chat/stream` | SSE wrapper over a single turn — emits `assistant.delta`, `tool.started`, `tool.completed`, `run.completed` events | `/v1/capabilities` advertises the full surface via `session_*` feature flags and `endpoints.session_*` entries so external UIs can detect support and fall back safely. Inline images are supported in `chat` and `chat/stream` payloads (multimodal-aware path). +`GET /api/sessions/{id}/messages` accepts `limit` and `offset`, then returns `object: "list"`, the resolved `session_id`, the current page in `data`, and the same pagination metadata fields the jobs endpoint uses: `limit`, `offset`, `total`, and `has_more`. + ```bash # fork a session and run one turn curl -X POST http://localhost:8642/api/sessions/$ID/fork \