diff --git a/agent/background_review.py b/agent/background_review.py index fffa902892..cc7641a99c 100644 --- a/agent/background_review.py +++ b/agent/background_review.py @@ -206,13 +206,15 @@ _REVIEW_MAX_ITERATIONS = 16 # Default aggregate INPUT-token budget for one review fork (#93057). The # fork's first request replays the full snapshot — a warm prompt-cache read -# that is cheap and intended (cache parity). After that, detached in-memory -# compaction bounds each request to roughly the compression threshold, but -# nothing capped the SUM across the review's tool loop: one production -# review made 8 requests replaying 1,487,951 input tokens total (four of -# them at 350k-384k). This budget caps the aggregate; the review tool loop -# stops before the provider call that would cross it (see -# ``_review_input_budget_exhausted`` in agent/conversation_loop.py). +# that is cheap and intended (cache parity), which is why both compression +# gates are deferred until the first provider response arrives +# (_review_fork_first_request_pending in agent/turn_context.py). After that, +# detached in-memory compaction bounds each request to roughly the +# compression threshold, but nothing capped the SUM across the review's tool +# loop: one production review made 8 requests replaying 1,487,951 input +# tokens total (four of them at 350k-384k). This budget caps the aggregate; +# the review tool loop stops before the provider call that would cross it +# (see ``_review_input_budget_exhausted`` in agent/conversation_loop.py). # 2x the historical 300k foreground trigger keeps legitimate reviews # comfortable while capping the pathological case. Override with # ``auxiliary.background_review.max_input_tokens``; 0 or a negative value @@ -1334,12 +1336,16 @@ def _run_review_in_thread( # with session_db=None / session_id="" makes every # compressor persist guard a no-op. # • Force in-place mode (never rotation) even if the parent's - # config selected rotation, and re-enable compression so the - # trigger gates in conversation_loop.py can fire. + # config selected rotation, and re-enable compression ONLY + # after the rebind succeeds (fail-closed — see below). While + # enabled, both compression gates stay deferred until the + # fork's first provider response so request #1 replays the + # full snapshot as a warm cache read. _review_compressor = getattr(review_agent, "context_compressor", None) _bind_review_compressor = getattr( _review_compressor, "bind_session_state", None ) + _review_compression_detached = False if callable(_bind_review_compressor): try: # Plugin/third-party context engines may not accept these @@ -1348,14 +1354,42 @@ def _run_review_in_thread( # and must never abort the review (same tolerance as the # init-time binding in agent_init.py). _bind_review_compressor(session_db=None, session_id="") + _review_compression_detached = True except Exception: - logger.debug( + # FAIL-CLOSED (adversarial review, #93057): if the rebind + # could not sever the engine's session binding, the + # compressor may still point at the parent's + # SessionDB/session_id. Enabling compression in that + # state would let durable cooldown/streak/ineffective- + # count writes land on the parent's row and re-open the + # #38727 sibling race. Keep the historical + # compression_enabled=False behavior instead and warn; + # the review still runs, bounded by the iteration cap + # and the aggregate input budget below. + logger.warning( "background-review compressor detachment failed; " - "keeping the engine's existing session binding", + "keeping compression DISABLED on this review fork " + "(fail-closed, issue #93057 / #38727)", exc_info=True, ) + # Force in-place mode (never rotation) even if the parent's + # config selected rotation. Re-enable compression ONLY after the + # compressor's session binding was successfully severed; an + # engine without a bind hook keeps the historical disabled + # behavior as well. review_agent.compression_in_place = True - review_agent.compression_enabled = True + review_agent.compression_enabled = _review_compression_detached + if _review_compression_detached: + # Warm-cache parity: the fork's FIRST provider request + # replays the parent's full snapshot as a warm prompt-cache + # read, so compaction must not rewrite the snapshot before + # that first request goes out. Defer both compression gates + # until the first provider response arrives (see + # _review_fork_first_request_pending in agent/turn_context.py + # and the pre-API gate in agent/conversation_loop.py); from + # the second request on, the fork's transcript is its own and + # compaction bounds it. + review_agent._review_defer_compaction_before_first_response = True # Aggregate input budget: compaction bounds any single request; # this bounds the WHOLE review. Iterations are already capped by # _REVIEW_MAX_ITERATIONS. Checked in agent/conversation_loop.py diff --git a/agent/conversation_loop.py b/agent/conversation_loop.py index 4f717cd835..3cc248c805 100644 --- a/agent/conversation_loop.py +++ b/agent/conversation_loop.py @@ -42,6 +42,7 @@ from agent.error_classifier import FailoverReason, classify_api_error from agent.message_metadata import append_message from agent.turn_context import ( _compression_warrants_another_preflight_pass, + _review_fork_first_request_pending, build_turn_context, compose_user_api_content, reanchor_current_turn_user_idx, @@ -2674,6 +2675,7 @@ def run_conversation( )() if ( agent.compression_enabled + and not _review_fork_first_request_pending(agent) and len(messages) > 1 and compression_attempts < max_compression_attempts and not _preflight_compression_blocked diff --git a/agent/turn_context.py b/agent/turn_context.py index 9df563d4ec..e20097a56b 100644 --- a/agent/turn_context.py +++ b/agent/turn_context.py @@ -320,6 +320,24 @@ def compression_made_progress( _compression_made_progress = compression_made_progress +def _review_fork_first_request_pending(agent: Any) -> bool: + """Whether a detached review fork has yet to send its first provider request. + + The background-review fork (issue #93057) replays the parent's FULL + snapshot on its first provider request as a warm prompt-cache read + (same-model cache parity). Compaction must not rewrite the snapshot + before that first request goes out — a compacted transcript would miss + the parent's cached prefix and turn a cheap cached replay into a cold + over-threshold write. Once the first provider response has arrived the + fork's tool loop is its own context, and both compression gates resume. + Dormant for every agent without the attribute. + """ + return bool( + getattr(agent, "_review_defer_compaction_before_first_response", False) + and not getattr(agent, "_turn_received_provider_response", False) + ) + + def _compression_warrants_another_preflight_pass( orig_tokens: int, new_tokens: int, threshold_tokens: int ) -> bool: @@ -879,11 +897,15 @@ def build_turn_context( _preflight_compression_blocked = False agent._turn_received_provider_response = False agent._turn_preflight_display_snapshot = None - if agent.compression_enabled and _should_run_preflight_estimate( - messages, - agent.context_compressor.protect_first_n, - agent.context_compressor.protect_last_n, - agent.context_compressor.threshold_tokens, + if ( + agent.compression_enabled + and not _review_fork_first_request_pending(agent) + and _should_run_preflight_estimate( + messages, + agent.context_compressor.protect_first_n, + agent.context_compressor.protect_last_n, + agent.context_compressor.threshold_tokens, + ) ): _preflight_tokens = estimate_request_tokens_rough( messages, diff --git a/hermes_cli/config_defaults.py b/hermes_cli/config_defaults.py index 61e8edc24d..8bfe2a4790 100644 --- a/hermes_cli/config_defaults.py +++ b/hermes_cli/config_defaults.py @@ -1267,11 +1267,14 @@ DEFAULT_CONFIG = { "extra_body": {}, "reasoning_effort": "", # per-task thinking level: none|minimal|low|medium|high|xhigh|max|ultra (empty = provider default) # Aggregate INPUT-token budget for one review fork (issue #93057). - # The fork compacts an oversized snapshot in memory before further - # provider calls; this caps the SUM of input tokens replayed - # across the whole review tool loop (iterations are separately - # capped at 16). The loop stops before the provider call that - # would cross the budget. 0 or a negative value = unlimited. + # The fork's FIRST request replays the full snapshot as a warm + # prompt-cache read (compaction is deferred until the first + # provider response arrives); after that it compacts an oversized + # snapshot in memory before further provider calls. This caps the + # SUM of input tokens replayed across the whole review tool loop + # (iterations are separately capped at 16). The loop stops before + # the provider call that would cross the budget. 0 or a negative + # value = unlimited. "max_input_tokens": 600000, }, "moa_reference": { diff --git a/tests/agent/test_compression_concurrent_fork.py b/tests/agent/test_compression_concurrent_fork.py index d9fc1c82de..eea66417d3 100644 --- a/tests/agent/test_compression_concurrent_fork.py +++ b/tests/agent/test_compression_concurrent_fork.py @@ -31,6 +31,7 @@ from __future__ import annotations import copy import inspect import json +import logging import os import sqlite3 import threading @@ -1190,7 +1191,8 @@ def test_real_lock_api_internal_errors_fail_closed_skips_compression( def test_review_fork_compacts_oversized_snapshot_in_memory(tmp_path: Path) -> None: - """An oversized review snapshot compacts in memory without mutating the parent. + """An oversized review snapshot replays warm on the first request, then + compacts in memory before further requests — without mutating the parent. Regression for #93057: the fork historically pinned ``compression_enabled = False`` because it shares the parent's session_id (issue #38727). That @@ -1200,11 +1202,14 @@ def test_review_fork_compacts_oversized_snapshot_in_memory(tmp_path: Path) -> No the parent's SessionDB/session_id and enables in-memory-only compaction. This test drives the REAL ``_run_review_in_thread`` + ``run_conversation`` - with a threshold-crossing snapshot and asserts: - • compression actually fired (a real threshold crossing, not just - construction-time flag state); - • the outbound provider request carries the compaction summary and none - of the middle snapshot turns; + with a threshold-crossing snapshot across two provider requests and + asserts: + • the FIRST request replays the full snapshot untouched (warm + prompt-cache parity) — no compaction summary, middle turns present; + • compression actually fired before the SECOND request (a real + threshold crossing, not just setup-time binding state), and that + request carries the compaction summary and none of the middle + snapshot turns; • the fork keeps the parent's session_id (prompt-cache parity) but its agent-level AND compressor-level session bindings are detached; • the parent's durable transcript, session row, and child-session graph @@ -1237,6 +1242,31 @@ def test_review_fork_compacts_oversized_snapshot_in_memory(tmp_path: Path) -> No captured = {} + def _tool_response(prompt_tokens: int) -> SimpleNamespace: + message = SimpleNamespace( + content=None, + reasoning_content=None, + reasoning=None, + tool_calls=[ + SimpleNamespace( + id="call_1", + type="function", + function=SimpleNamespace( + name="web_search", arguments='{"query": "x"}' + ), + ) + ], + ) + return SimpleNamespace( + choices=[SimpleNamespace(message=message, finish_reason="tool_calls")], + model="test/model", + usage=SimpleNamespace( + prompt_tokens=prompt_tokens, + completion_tokens=1, + total_tokens=prompt_tokens + 1, + ), + ) + def _final_response(): return SimpleNamespace( choices=[ @@ -1271,6 +1301,9 @@ def test_review_fork_compacts_oversized_snapshot_in_memory(tmp_path: Path) -> No captured["input_budget"] = getattr( self, "_review_input_token_budget", "missing" ) + captured["defer_first_request"] = getattr( + self, "_review_defer_compaction_before_first_response", "missing" + ) captured["compressor_session_db"] = getattr( self.context_compressor, "_session_db", "missing" ) @@ -1288,8 +1321,16 @@ def test_review_fork_compacts_oversized_snapshot_in_memory(tmp_path: Path) -> No {"role": "assistant", "content": "summary acknowledged"}, ] ) + # Compress on the first pressure check after the first response, then + # stand down so the compacted request proceeds instead of looping. + _should_compress_calls = {"count": 0} + + def _should_compress(_tokens): + _should_compress_calls["count"] += 1 + return _should_compress_calls["count"] == 1 + self.context_compressor.should_compress = MagicMock( - side_effect=lambda _tokens: True + side_effect=_should_compress ) self.context_compressor.should_compress_info = MagicMock( return_value=(True, "over threshold") @@ -1306,17 +1347,33 @@ def test_review_fork_compacts_oversized_snapshot_in_memory(tmp_path: Path) -> No self.context_compressor.select_context = MagicMock(return_value=None) self._compression_feasibility_checked = True self.client = MagicMock() - self.client.chat.completions.create.side_effect = [_final_response()] + self.client.chat.completions.create.side_effect = [ + _tool_response(100), + _final_response(), + ] self._disable_streaming = True self._use_prompt_caching = False + def _fake_execute_tool_calls(assistant_message, messages, *_args): + tool_call = assistant_message.tool_calls[0] + messages.append( + { + "role": "tool", + "name": tool_call.function.name, + "tool_call_id": tool_call.id, + "content": "ok", + } + ) + + self._execute_tool_calls = _fake_execute_tool_calls + result = real_run_conversation(self, *args, **kwargs) captured["compression_calls"] = self.context_compressor.compress.call_count - captured["create_calls"] = self.client.chat.completions.create.call_count - last_call = self.client.chat.completions.create.call_args - captured["outbound"] = ( - last_call.kwargs.get("messages") if last_call else None - ) + create = self.client.chat.completions.create + captured["create_calls"] = create.call_count + captured["outbound"] = [ + call.kwargs.get("messages") for call in create.call_args_list + ] return result try: @@ -1329,15 +1386,30 @@ def test_review_fork_compacts_oversized_snapshot_in_memory(tmp_path: Path) -> No "fork, which removed the only bound on the replayed snapshot " "(issue #93057)." ) - assert captured["create_calls"] >= 1 - outbound_contents = [ - str(m.get("content", "")) for m in captured["outbound"] - ] + assert captured["create_calls"] == 2, ( + f"expected a 2-request review (tool call + final), " + f"got {captured['create_calls']}" + ) + first_outbound, second_outbound = captured["outbound"] + first_contents = [str(m.get("content", "")) for m in first_outbound] + second_contents = [str(m.get("content", "")) for m in second_outbound] + # Warm-cache parity: the first request replays the full snapshot + # untouched — middle turns present, no compaction summary yet. + assert any("review turn 12" in text for text in first_contents), ( + "the review fork's FIRST request must replay the full snapshot " + "(warm prompt-cache read) — compaction must not rewrite it before " + "the first provider call" + ) + assert not any( + "[CONTEXT COMPACTION]" in text for text in first_contents + ), f"first request was compacted prematurely: {first_contents!r}" + # The SECOND request carries the compaction summary and none of the + # middle snapshot turns. assert any( "[CONTEXT COMPACTION] review summary" in text - for text in outbound_contents - ), f"outbound request did not contain the compaction summary: {outbound_contents!r}" - assert not any("review turn 12" in text for text in outbound_contents), ( + for text in second_contents + ), f"outbound request did not contain the compaction summary: {second_contents!r}" + assert not any("review turn 12" in text for text in second_contents), ( "outbound request still replays the middle of the snapshot — " "the review replayed an unbounded transcript despite compaction" ) @@ -1357,6 +1429,7 @@ def test_review_fork_compacts_oversized_snapshot_in_memory(tmp_path: Path) -> No assert captured["compressor_session_id"] == "" assert captured["compression_enabled"] is True assert captured["compression_in_place"] is True + assert captured["defer_first_request"] is True assert isinstance(captured["input_budget"], int) and captured["input_budget"] > 0 # Parent session must be byte-for-byte unchanged after the review @@ -1375,6 +1448,97 @@ def test_review_fork_compacts_oversized_snapshot_in_memory(tmp_path: Path) -> No db.close() +def test_review_fork_fails_closed_when_compressor_rebind_raises( + tmp_path: Path, caplog +) -> None: + """A failed compressor detachment must keep the fork's compression OFF. + + Regression for the #93057 adversarial review: if ``bind_session_state`` + cannot sever the engine's binding to the parent's SessionDB/session_id, + enabling compression would run it against the parent's live session + binding — durable cooldown/streak/ineffective-count writes on the + parent's row and the sibling-session race behind #38727 re-opened. The + fork must fail CLOSED: keep the historical ``compression_enabled = False`` + behavior and warn. The review still runs (the iteration cap and the + aggregate input budget still bound it). + """ + import agent.background_review as br + from agent.context_compressor import ContextCompressor + + parent_sid = "REVIEW_FORK_REBIND_FAIL_CLOSED_93057" + + db = SessionDB(db_path=tmp_path / "state.db") + db.create_session(parent_sid, source="discord") + parent = _build_agent_with_db(db, parent_sid) + parent._cached_system_prompt = "stable parent prompt" + + snapshot = [ + { + "role": "user" if i % 2 == 0 else "assistant", + "content": f"review turn {i}", + } + for i in range(8) + ] + + captured = {} + + def _capture_fork_flags(self, *args, **kwargs): + captured["compression_enabled"] = self.compression_enabled + captured["compression_in_place"] = self.compression_in_place + captured["input_budget"] = getattr( + self, "_review_input_token_budget", "missing" + ) + return { + "completed": True, + "final_response": "review complete", + "api_call_count": 0, + } + + # The worker does a local ``from run_agent import AIAgent``; patching the + # class method covers that import path. + from run_agent import AIAgent + + _real_bind = ContextCompressor.bind_session_state + + def _failing_bind(self, session_db=None, session_id=""): + # Only the detachment rebind may fail; any other binding passes + # through so the fork's construction path stays intact. + if session_db is None: + raise RuntimeError("detachment boom") + return _real_bind(self, session_db, session_id) + + try: + with ( + patch.object(AIAgent, "run_conversation", _capture_fork_flags), + patch.object( + ContextCompressor, "bind_session_state", _failing_bind + ), + ): + with caplog.at_level(logging.WARNING, logger="agent.background_review"): + br._run_review_in_thread(parent, snapshot, "review this conversation") + + assert captured["compression_enabled"] is False, ( + "FIX REGRESSION: a failed compressor rebind must leave " + "compression_enabled False on the review fork (fail-closed). " + "Enabling compression with the engine still bound to the " + "parent's session re-opens the #38727 race (issue #93057)." + ) + assert any( + "detachment failed" in record.message for record in caplog.records + ), ( + "the failed rebind must log a warning so operators can see the " + "fork fell back to the pre-fix behavior" + ) + assert ( + isinstance(captured["input_budget"], int) and captured["input_budget"] > 0 + ), ( + "the aggregate input budget must still be armed on the fail-closed " + "path — it bounds the review even when compaction stays off" + ) + finally: + db.close() + + # ── Lease-refresher bounded-failure tolerance (salvage follow-up, #54465) ──── # A single falsy refresh (transient DB blip) must NOT permanently kill the # lease — only a *persistent* failure (genuine lost-ownership) should stop the