fix(compression): arm the failure cooldown when codex compaction fails
Closes #75364. `_compress_context_via_codex_app_server` returns the transcript unchanged when the codex thread reports `interrupted` or `error`. The session is therefore still above threshold, and nothing records that the attempt failed — so the next turn retries immediately, and keeps retrying for as long as the condition persists. Every other compression path arms the shared failure cooldown, records an ineffective-compression strike, or both. This path records neither: * `_hygiene_compression_failure_cooldowns` is set only on `asyncio.TimeoutError`, or behind `_last_compress_aborted`, which is assigned exclusively in `context_compressor.py` on the Hermes summarizer path. * `compression_ineffective_count` lives in `ContextCompressor`, and this path returns before any compressor bookkeeping runs. `compress_context` already documents the rule this path was missing — "Every automatic entrypoint must honor compressor-owned cooldown and breaker state" — but the codex branch dispatches above that block and returns from inside it. `result.interrupted` needs no unusual configuration to occur: an ordinary user message arriving mid-compaction sets it (see `codex_app_server_session.py`, which produces the "compact turn interrupted" string). Observed in production on a Discord gateway session at ~315k tokens against a 258k window, where compaction was attempted on essentially every turn for ~70 minutes; the session's `compression_ineffective_count` was still 0 afterwards. This reuses the existing cooldown rather than adding a new mechanism: * arm `_record_compression_failure_cooldown` with the existing `_SUMMARY_FAILURE_COOLDOWN_SECONDS` when compaction returns interrupted/error; * honor an active cooldown on entry, matching the Hermes path. `force=True` bypasses both, so an explicit /compress is never braked by a failure it did not cause, and a successful compaction arms nothing.
This commit is contained in:
@@ -4699,6 +4699,46 @@ def compress_context(
|
||||
commit_fence.finish_commit()
|
||||
|
||||
|
||||
def _codex_compaction_cooldown_remaining(agent: Any) -> float:
|
||||
"""Seconds left on this session's compaction-failure cooldown (0 = clear)."""
|
||||
compressor = getattr(agent, "context_compressor", None)
|
||||
getter = getattr(compressor, "get_active_compression_failure_cooldown", None)
|
||||
if not callable(getter):
|
||||
return 0.0
|
||||
try:
|
||||
state = getter(refresh=True)
|
||||
except Exception:
|
||||
logger.debug("codex compaction cooldown lookup failed", exc_info=True)
|
||||
return 0.0
|
||||
if not state:
|
||||
return 0.0
|
||||
try:
|
||||
return max(0.0, float(state.get("remaining_seconds") or 0.0))
|
||||
except (TypeError, ValueError):
|
||||
return 0.0
|
||||
|
||||
|
||||
def _record_codex_compaction_failure(agent: Any, error: str) -> None:
|
||||
"""Arm the shared compression-failure cooldown after a failed compaction.
|
||||
|
||||
The codex path returns the transcript unchanged on failure, so the session
|
||||
is still above threshold and the next turn retries immediately. Every other
|
||||
compression path records a cooldown, an ineffective-compression strike, or
|
||||
both; this one recorded neither, so an interrupted compaction retried once
|
||||
per turn for as long as the condition persisted.
|
||||
"""
|
||||
from agent.context_compressor import _SUMMARY_FAILURE_COOLDOWN_SECONDS
|
||||
|
||||
compressor = getattr(agent, "context_compressor", None)
|
||||
recorder = getattr(compressor, "_record_compression_failure_cooldown", None)
|
||||
if not callable(recorder):
|
||||
return
|
||||
try:
|
||||
recorder(_SUMMARY_FAILURE_COOLDOWN_SECONDS, error)
|
||||
except Exception:
|
||||
logger.debug("codex compaction cooldown persist failed", exc_info=True)
|
||||
|
||||
|
||||
def _compress_context_via_codex_app_server(
|
||||
agent: Any,
|
||||
messages: list,
|
||||
@@ -4734,6 +4774,25 @@ def _compress_context_via_codex_app_server(
|
||||
existing_prompt = agent._build_system_prompt(system_message)
|
||||
return messages, existing_prompt
|
||||
|
||||
# Automatic entrypoints must honor the compressor-owned cooldown, the same
|
||||
# way the Hermes path below does. An active cooldown means a recent
|
||||
# compaction already failed; retrying every turn is what thrashes.
|
||||
if not force:
|
||||
_cooldown_remaining = _codex_compaction_cooldown_remaining(agent)
|
||||
if _cooldown_remaining > 0:
|
||||
logger.info(
|
||||
"codex app-server compaction skipped: failure cooldown active "
|
||||
"for %.0fs (session=%s messages=%d tokens=~%s)",
|
||||
_cooldown_remaining,
|
||||
getattr(agent, "session_id", None) or "none",
|
||||
len(messages),
|
||||
f"{approx_tokens:,}" if approx_tokens else "unknown",
|
||||
)
|
||||
existing_prompt = getattr(agent, "_cached_system_prompt", None)
|
||||
if not existing_prompt:
|
||||
existing_prompt = agent._build_system_prompt(system_message)
|
||||
return messages, existing_prompt
|
||||
|
||||
codex_session = getattr(agent, "_codex_session", None)
|
||||
if codex_session is None:
|
||||
logger.info(
|
||||
@@ -4787,6 +4846,12 @@ def _compress_context_via_codex_app_server(
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
# The transcript is returned unchanged, so the session is still over
|
||||
# threshold. Without a brake the next turn retries immediately.
|
||||
_record_codex_compaction_failure(
|
||||
agent,
|
||||
str(getattr(result, "error", None) or "compaction interrupted"),
|
||||
)
|
||||
existing_prompt = getattr(agent, "_cached_system_prompt", None)
|
||||
if not existing_prompt:
|
||||
existing_prompt = agent._build_system_prompt(system_message)
|
||||
|
||||
@@ -205,3 +205,119 @@ def test_codex_native_boundary_clears_stale_hermes_fallback_streak():
|
||||
assert _record_codex_app_server_compaction(agent, turn) is True
|
||||
assert compressor._fallback_compression_streak == 0
|
||||
assert compressor._verify_compaction_cleared_threshold is True
|
||||
|
||||
|
||||
class RecordingCooldownCompressor(SimpleNamespace):
|
||||
"""Compressor stub exposing the real cooldown API surface."""
|
||||
|
||||
def __init__(self, remaining=0.0):
|
||||
super().__init__(
|
||||
compression_count=0,
|
||||
last_compression_rough_tokens=0,
|
||||
last_prompt_tokens=123,
|
||||
last_completion_tokens=45,
|
||||
awaiting_real_usage_after_compression=False,
|
||||
)
|
||||
self.remaining = remaining
|
||||
self.recorded = []
|
||||
|
||||
def get_active_compression_failure_cooldown(self, *, refresh=False):
|
||||
if self.remaining <= 0:
|
||||
return None
|
||||
return {"remaining_seconds": self.remaining, "error": "prior failure"}
|
||||
|
||||
def _record_compression_failure_cooldown(self, seconds, error):
|
||||
self.recorded.append((seconds, error))
|
||||
self.remaining = float(seconds)
|
||||
|
||||
|
||||
def test_interrupted_codex_compaction_arms_the_failure_cooldown():
|
||||
"""Regression: the codex path returned unchanged with no brake, so the
|
||||
session stayed above threshold and the next turn retried immediately."""
|
||||
from agent.context_compressor import _SUMMARY_FAILURE_COOLDOWN_SECONDS
|
||||
|
||||
agent = DummyAgent(
|
||||
TurnResult(
|
||||
thread_id="thread-1",
|
||||
turn_id="compact-turn-1",
|
||||
interrupted=True,
|
||||
error="compact turn interrupted",
|
||||
),
|
||||
auto_compaction="hermes",
|
||||
)
|
||||
agent.context_compressor = RecordingCooldownCompressor()
|
||||
messages = [{"role": "user", "content": "hi"}]
|
||||
|
||||
returned, prompt = compress_context(
|
||||
agent, messages, "system", approx_tokens=100000, task_id="test"
|
||||
)
|
||||
|
||||
assert returned is messages
|
||||
assert prompt == "cached prompt"
|
||||
assert agent.context_compressor.recorded == [
|
||||
(_SUMMARY_FAILURE_COOLDOWN_SECONDS, "compact turn interrupted")
|
||||
]
|
||||
|
||||
|
||||
def test_codex_compaction_error_without_interrupt_also_arms_cooldown():
|
||||
agent = DummyAgent(
|
||||
TurnResult(thread_id="thread-1", turn_id="compact-turn-1", error="boom"),
|
||||
auto_compaction="hermes",
|
||||
)
|
||||
agent.context_compressor = RecordingCooldownCompressor()
|
||||
messages = [{"role": "user", "content": "hi"}]
|
||||
|
||||
compress_context(
|
||||
agent, messages, "system", approx_tokens=100000, task_id="test"
|
||||
)
|
||||
|
||||
assert len(agent.context_compressor.recorded) == 1
|
||||
assert agent.context_compressor.recorded[0][1] == "boom"
|
||||
|
||||
|
||||
def test_active_cooldown_blocks_automatic_codex_compaction():
|
||||
agent = DummyAgent(
|
||||
TurnResult(thread_id="thread-1", turn_id="compact-turn-1"),
|
||||
auto_compaction="hermes",
|
||||
)
|
||||
agent.context_compressor = RecordingCooldownCompressor(remaining=120.0)
|
||||
session = agent._codex_session
|
||||
messages = [{"role": "user", "content": "hi"}]
|
||||
|
||||
returned, prompt = compress_context(
|
||||
agent, messages, "system", approx_tokens=100000, task_id="test"
|
||||
)
|
||||
|
||||
assert returned is messages
|
||||
assert prompt == "cached prompt"
|
||||
assert session.calls == 0, "compaction ran despite an active cooldown"
|
||||
|
||||
|
||||
def test_force_bypasses_the_codex_compaction_cooldown():
|
||||
"""An explicit /compress is a user decision and must not be braked by a
|
||||
failure it did not cause."""
|
||||
agent = DummyAgent(TurnResult(thread_id="thread-1", turn_id="compact-turn-1"))
|
||||
agent.context_compressor = RecordingCooldownCompressor(remaining=120.0)
|
||||
session = agent._codex_session
|
||||
messages = [{"role": "user", "content": "hi"}]
|
||||
|
||||
compress_context(
|
||||
agent, messages, "system", approx_tokens=100000, task_id="test", force=True
|
||||
)
|
||||
|
||||
assert session.calls == 1
|
||||
|
||||
|
||||
def test_successful_codex_compaction_arms_no_cooldown():
|
||||
agent = DummyAgent(
|
||||
TurnResult(thread_id="thread-1", turn_id="compact-turn-1"),
|
||||
auto_compaction="hermes",
|
||||
)
|
||||
agent.context_compressor = RecordingCooldownCompressor()
|
||||
messages = [{"role": "user", "content": "hi"}]
|
||||
|
||||
compress_context(
|
||||
agent, messages, "system", approx_tokens=100000, task_id="test"
|
||||
)
|
||||
|
||||
assert agent.context_compressor.recorded == []
|
||||
|
||||
Reference in New Issue
Block a user