From 1c21e96ed022761409b9fd17f9ca66efb276a547 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 22 Jul 2026 04:56:23 -0700 Subject: [PATCH] refactor(api): dedupe compressed-transcript persist sites, drop config opt-out MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Rework on top of the salvaged #58133 commits: - Remove the compression.persist_in_response_store config key — this is a bug fix (stored transcripts must reflect what the agent will actually replay), not behavior that should be opt-out-able. - Drop the per-request load_config() imports the handler-level persist blocks added. - Dedupe the two handler-level persist blocks: the compressed-transcript substitution already lives in _build_response_conversation_history (via result["_compressed"]), so the handlers only need to propagate the effective (possibly rotation-changed) session_id. The streaming path does this via a new session_id_snapshot arg on _persist_response_snapshot; the non-streaming path picks up result["session_id"] directly. - Rotation propagation no longer gates on history-from-store: the first request in a chain can also rotate, and its stored session_id must be the child session or the next previous_response_id request resumes the pre-rotation session and re-compresses every turn. --- gateway/platforms/api_server.py | 82 +++++++-------------------------- hermes_cli/config.py | 10 ---- 2 files changed, 16 insertions(+), 76 deletions(-) diff --git a/gateway/platforms/api_server.py b/gateway/platforms/api_server.py index d04803f5ebb79..4b3b4d91228c2 100644 --- a/gateway/platforms/api_server.py +++ b/gateway/platforms/api_server.py @@ -3213,7 +3213,6 @@ class APIServerAdapter(BasePlatformAdapter): store: bool, session_id: str, gateway_session_key: Optional[str] = None, - history_from_store: bool = False, ) -> "web.StreamResponse": """Write an SSE stream for POST /v1/responses (OpenAI Responses API). @@ -3311,6 +3310,7 @@ class APIServerAdapter(BasePlatformAdapter): response_env: Dict[str, Any], *, conversation_history_snapshot: Optional[List[Dict[str, Any]]] = None, + session_id_snapshot: Optional[str] = None, ) -> None: if not store: return @@ -3321,7 +3321,7 @@ class APIServerAdapter(BasePlatformAdapter): "response": response_env, "conversation_history": conversation_history_snapshot, "instructions": instructions, - "session_id": session_id, + "session_id": session_id_snapshot or session_id, }) if conversation: self._response_store.set_conversation(conversation, response_id) @@ -3724,38 +3724,15 @@ class APIServerAdapter(BasePlatformAdapter): result, final_response_text, ) - # Persist compressed messages when compression occurred and - # history was loaded from response_store (same logic as the - # non-streaming path). - if history_from_store and isinstance(result, dict): - _result_sid = result.get("session_id") - _did_compress = bool(result.get("_compressed")) - _rotated = bool(_result_sid and _result_sid != session_id) - if _did_compress or _rotated: - try: - from hermes_cli.config import load_config - _comp_cfg = load_config().get("compression", {}) - if _comp_cfg.get("persist_in_response_store", True): - _agent_messages = result.get("messages") - if isinstance(_agent_messages, list) and _agent_messages: - _mode = "in-place" if _did_compress and not _rotated else "rotation" - logger.info( - "Compression persisted in response_store (streaming, %s): " - "%d messages (was %d before compression)", - _mode, - len(_agent_messages), - len(conversation_history) + 1, - ) - full_history = list(_agent_messages) - except Exception as e: - logger.warning( - "Failed to persist compressed response_store snapshot " - "(streaming): %s", - e, - ) + # Compression-aware transcript substitution happens inside + # _build_response_conversation_history (result["_compressed"]); + # here we only propagate a compression-rotated session_id so + # previous_response_id chaining resumes the child session. + _result_sid = result.get("session_id") if isinstance(result, dict) else None _persist_response_snapshot( completed_env, conversation_history_snapshot=full_history, + session_id_snapshot=_result_sid if isinstance(_result_sid, str) and _result_sid else None, ) terminal_snapshot_persisted = True await _write_event("response.completed", { @@ -3908,14 +3885,12 @@ class APIServerAdapter(BasePlatformAdapter): logger.debug("Both conversation_history and previous_response_id provided; using conversation_history") stored_session_id = None - _history_from_store = False if not conversation_history and previous_response_id: stored = self._response_store.get(previous_response_id) if stored is None: return web.json_response(_openai_error(f"Previous response not found: {previous_response_id}"), status=404) conversation_history = list(stored.get("conversation_history", [])) stored_session_id = stored.get("session_id") - _history_from_store = True # If no instructions provided, carry forward from previous if instructions is None: instructions = stored.get("instructions") @@ -4018,7 +3993,6 @@ class APIServerAdapter(BasePlatformAdapter): store=store, session_id=session_id, gateway_session_key=gateway_session_key, - history_from_store=_history_from_store, ) async def _compute_response(): @@ -4071,39 +4045,15 @@ class APIServerAdapter(BasePlatformAdapter): final_response, ) - # If compression occurred during the agent run and the history was - # loaded from the response_store (not supplied explicitly by the - # client), persist the compressed messages instead of the original - # uncompressed chain. This prevents repeated re-compression on - # subsequent requests in the same conversation. + # Persist the effective session ID surfaced by _run_agent so that + # compression-triggered session rotations propagate to the stored + # response and the X-Hermes-Session-Id header. Without this, + # previous_response_id chaining keeps resuming the pre-rotation + # session and re-triggers compression on every subsequent request. _effective_session_id = session_id - if _history_from_store and isinstance(result, dict): - _result_sid = result.get("session_id") - _did_compress = bool(result.get("_compressed")) - _rotated = bool(_result_sid and _result_sid != session_id) - if _did_compress or _rotated: - try: - from hermes_cli.config import load_config - _comp_cfg = load_config().get("compression", {}) - if _comp_cfg.get("persist_in_response_store", True): - _agent_messages = result.get("messages") - if isinstance(_agent_messages, list) and _agent_messages: - _mode = "in-place" if _did_compress and not _rotated else "rotation" - logger.info( - "Compression persisted in response_store (%s): " - "%d messages (was %d before compression)", - _mode, - len(_agent_messages), - len(conversation_history) + 1, - ) - full_history = list(_agent_messages) - if _rotated and _result_sid: - _effective_session_id = _result_sid - except Exception as e: - logger.warning( - "Failed to persist compressed response_store snapshot: %s", - e, - ) + _result_sid = result.get("session_id") if isinstance(result, dict) else None + if isinstance(_result_sid, str) and _result_sid: + _effective_session_id = _result_sid # Build output items from the current turn only. AIAgent returns a # full transcript in result["messages"], while older/mocked paths may diff --git a/hermes_cli/config.py b/hermes_cli/config.py index 20fd2f7776be1..7a527b75f53c0 100644 --- a/hermes_cli/config.py +++ b/hermes_cli/config.py @@ -1529,16 +1529,6 @@ DEFAULT_CONFIG = { # session_search and recoverable, not deleted. # Default False during rollout; will flip on # after live validation. - "persist_in_response_store": True, # When True, if compression occurs during a - # /v1/responses request that loads history from the - # response_store (via previous_response_id or - # conversation name), the compressed messages are - # persisted as the stored conversation_history - # snapshot. This prevents repeated re-compression - # on subsequent requests in the same chain. - # Only applies when the client does NOT supply an - # explicit conversation_history array. Set to False - # to preserve legacy behavior (store uncompressed). }, # Kanban subsystem (orchestrator workers + dispatcher-driven child tasks).