fix(review): fail-closed compressor detachment + warm-cache first request (#93057 review)
Adversarial-review fixes for the #93057 snapshot-compaction PR: - Fail-closed detachment: only re-enable compression after bind_session_state successfully severs the engine's parent binding. A failed rebind keeps the historical compression_enabled=False behavior and warns, instead of running compaction against a compressor still bound to the parent's SessionDB (#38727 re-open). - Warm-cache parity: defer both compression gates (turn-prologue preflight + pre-API pressure check) until the fork's first provider response, so the first request replays the full snapshot as the intended cached read and compaction applies from the second request on — matching the documented budget mental model. - Tests: regression for the rebind-failure fail-closed path (red on pre-fix code) and the existing threshold-crossing test reworked to a two-request review asserting the warm first request + compacted second request. 116 tests green across all touched suites; ruff clean.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user