From 7a21bfe68abd241dd3fb934a2b96ef182753bcf1 Mon Sep 17 00:00:00 2001 From: Uttkarsh Tiwari <38376852+uttkarsh-26@users.noreply.github.com> Date: Sat, 25 Jul 2026 22:33:48 +0530 Subject: [PATCH] fix(compression): suppress duplicate completion notices --- agent/conversation_compression.py | 10 ++- gateway/run.py | 3 + .../agent/test_compression_concurrent_fork.py | 64 +++++++++++++++++++ .../test_compression_progress_notices.py | 21 +++--- .../test_codex_app_server_compaction.py | 25 ++++++++ 5 files changed, 109 insertions(+), 14 deletions(-) diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 2e5c0cbde2..f2a77bbe24 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -193,6 +193,7 @@ CONTEXT_OVERFLOW_BLOCKED_WARNING_TEMPLATE = ( # same constants the emission sites use) through the gateway noise filter. ROUTINE_COMPRESSION_STATUS_SAMPLES = ( COMPACTION_STATUS, + COMPACTION_DONE_STATUS, PRE_API_COMPRESSION_STATUS_TEMPLATE.format(tokens=123456), PREFLIGHT_COMPRESSION_STATUS_TEMPLATE.format(tokens=120000, threshold=100000), IDLE_COMPACTION_STATUS_TEMPLATE.format(idle_seconds=3600, tokens=120000), @@ -2493,6 +2494,7 @@ def compress_context( if _compaction_status: agent._emit_status(_compaction_status) _compaction_done_emitted = False + _compaction_succeeded = False def _complete_compaction_lifecycle() -> None: nonlocal _compaction_done_emitted @@ -2502,7 +2504,7 @@ def compress_context( # A suppressed start (quiet context engine) opened no visible # compaction phase — emit no terminal edge either. Failure warnings # go through agent._emit_warning and are never suppressed here. - if _compaction_status_emitted: + if _compaction_status_emitted and _compaction_succeeded: _emit_compaction_done(agent) # ── Compression lock ──────────────────────────────────────────────── @@ -4110,6 +4112,7 @@ def compress_context( else None ), ) + _compaction_succeeded = True return compressed, new_system_prompt finally: # Release the lock on the OLD session_id only AFTER rotation completed @@ -4185,13 +4188,15 @@ def _compress_context_via_codex_app_server( pass _compaction_done_emitted = False + _compaction_succeeded = False def _complete_compaction_lifecycle() -> None: nonlocal _compaction_done_emitted if _compaction_done_emitted: return _compaction_done_emitted = True - _emit_compaction_done(agent) + if _compaction_succeeded: + _emit_compaction_done(agent) _activity_heartbeat: Optional[_CompressionActivityHeartbeat] = None try: @@ -4265,6 +4270,7 @@ def _compress_context_via_codex_app_server( existing_prompt = getattr(agent, "_cached_system_prompt", None) if not existing_prompt: existing_prompt = agent._build_system_prompt(system_message) + _compaction_succeeded = True _complete_compaction_lifecycle() return messages, existing_prompt diff --git a/gateway/run.py b/gateway/run.py index 690c8811a0..8a0b5473e9 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -49,6 +49,7 @@ from typing import Awaitable, Callable, Dict, Optional, Any, List, Tuple, Union, from agent.async_utils import consume_detached_task_result, safe_schedule_threadsafe from agent.conversation_compression import ( + COMPACTION_DONE_STATUS, COMPACTION_STATUS, COMPRESSION_RETRY_CONTEXT_REDUCED_STATUS_TEMPLATE, COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE, @@ -147,6 +148,7 @@ _TELEGRAM_NOISY_STATUS_RE = re.compile( r"|max\s+retries\s+\(\d+\).*(?:trying\s+fallback|exhausted|invalid\s+responses)" r"|stream\s+(?:drop|drop\s+mid\s+tool-call).+retry\s+\d" r"|stale\s+connections\s+from\s+a\s+previous\s+provider\s+issue" + rf"|{re.escape(COMPACTION_DONE_STATUS)}" r")", re.IGNORECASE | re.DOTALL, ) @@ -338,6 +340,7 @@ _COMPRESSION_PROGRESS_STATUS_RE = re.compile( _status_template_to_regex(_template) for _template in ( COMPACTION_STATUS, + COMPACTION_DONE_STATUS, PRE_API_COMPRESSION_STATUS_TEMPLATE, PREFLIGHT_COMPRESSION_STATUS_TEMPLATE, IDLE_COMPACTION_STATUS_TEMPLATE, diff --git a/tests/agent/test_compression_concurrent_fork.py b/tests/agent/test_compression_concurrent_fork.py index eea66417d3..3145e702b9 100644 --- a/tests/agent/test_compression_concurrent_fork.py +++ b/tests/agent/test_compression_concurrent_fork.py @@ -167,6 +167,70 @@ def test_compression_activity_heartbeat_touches_agent_during_long_compress(tmp_p assert db.get_compression_lock_holder(session_id) is None +def test_lock_contender_preserves_terminal_compaction_lifecycle(tmp_path: Path) -> None: + """A lock loser still closes the structured compaction lifecycle. + + The gateway independently filters this routine notice for chat surfaces + unless ``compression.progress_notices`` is enabled. The low-level event + must remain available so the desktop can retire its compaction phase. + """ + from agent.conversation_compression import COMPACTION_DONE_STATUS + + db = SessionDB(db_path=tmp_path / "state.db") + session_id = "LOCK_CONTENDER_STATUS_TEST" + db.create_session(session_id, source="discord") + assert db.try_acquire_compression_lock(session_id, "winner", ttl_seconds=60) + + agent = _build_agent_with_db(db, session_id) + status_events: list[tuple[str, str]] = [] + setattr( + agent, + "status_callback", + lambda event, message: status_events.append((event, message)), + ) + messages = [{"role": "user", "content": f"m{i}"} for i in range(20)] + + returned, _system_prompt = agent._compress_context( + messages, + "sys", + approx_tokens=120_000, + ) + + assert returned is messages + assert getattr(agent, "_compression_skipped_due_to_lock", None) == "winner" + assert status_events.count(("compacted", COMPACTION_DONE_STATUS)) == 1 + + +def test_failed_session_split_does_not_announce_compaction_complete(tmp_path: Path) -> None: + """A failed durable split must not emit a successful completion edge.""" + from agent.conversation_compression import COMPACTION_DONE_STATUS + + db = SessionDB(db_path=tmp_path / "state.db") + session_id = "FAILED_SPLIT_STATUS_TEST" + db.create_session(session_id, source="discord") + agent = _build_agent_with_db(db, session_id) + setattr(agent, "compression_in_place", False) + db.publish_compression_child = MagicMock(side_effect=RuntimeError("split boom")) + status_events: list[tuple[str, str]] = [] + setattr( + agent, + "status_callback", + lambda event, message: status_events.append((event, message)), + ) + messages = [{"role": "user", "content": f"m{i}"} for i in range(20)] + + agent._compress_context( + messages, + "sys", + approx_tokens=120_000, + force=True, + ) + + db.publish_compression_child.assert_called_once() + assert ("compacted", COMPACTION_DONE_STATUS) not in status_events + assert db.get_compression_lock_holder(session_id) is None + + def test_compression_activity_heartbeat_stops_on_compress_exception(tmp_path: Path) -> None: """Exception paths must stop the heartbeat and release the compression lock.""" db = SessionDB(db_path=tmp_path / "state.db") diff --git a/tests/gateway/test_compression_progress_notices.py b/tests/gateway/test_compression_progress_notices.py index 0c26534774..12d7b93186 100644 --- a/tests/gateway/test_compression_progress_notices.py +++ b/tests/gateway/test_compression_progress_notices.py @@ -78,23 +78,20 @@ def test_enabled_still_suppresses_non_compression_noise( @pytest.mark.parametrize("enabled", [True, False], ids=["enabled", "default"]) @pytest.mark.parametrize("platform", CHAT_PLATFORMS) -def test_compaction_completion_notice_reaches_chat(monkeypatch, platform, enabled): - """The #69546 'compacted' lifecycle edge is deliverable on chat surfaces. - - COMPACTION_DONE_STATUS already flows through the status callback on - compaction completion and is not matched by the noise regex — the opt-in - gate must not change that in either mode, so users who enable - progress_notices see the completion stat notice paired with the start. - """ +def test_compaction_completion_notice_respects_progress_notices_gate( + monkeypatch, platform, enabled +): + """The completion edge follows the same opt-in gate as the start edge.""" monkeypatch.setattr( gateway_run, "_load_gateway_config", lambda: {"compression": {"progress_notices": enabled}}, ) - assert ( - _prepare_gateway_status_message(platform, "compacted", COMPACTION_DONE_STATUS) - == COMPACTION_DONE_STATUS - ) + result = _prepare_gateway_status_message(platform, "compacted", COMPACTION_DONE_STATUS) + if enabled: + assert result == COMPACTION_DONE_STATUS + else: + assert result is None def test_enabled_gate_does_not_leak_to_raw_platforms(progress_notices_enabled): diff --git a/tests/run_agent/test_codex_app_server_compaction.py b/tests/run_agent/test_codex_app_server_compaction.py index e72b9eea7d..2895a56312 100644 --- a/tests/run_agent/test_codex_app_server_compaction.py +++ b/tests/run_agent/test_codex_app_server_compaction.py @@ -150,6 +150,31 @@ def test_codex_app_server_compaction_heartbeat_refreshes_activity_while_waiting( +def test_codex_app_server_compression_failure_preserves_bookkeeping(): + agent = DummyAgent(TurnResult(error="compact failed")) + messages = [{"role": "user", "content": "hi"}] + + returned, prompt = compress_context( + agent, + messages, + "system", + approx_tokens=100000, + force=True, + ) + + assert returned is messages + assert prompt == "cached prompt" + assert agent._codex_session.calls == 1 + assert agent.context_compressor.compression_count == 0 + assert agent.context_compressor.last_prompt_tokens == 123 + assert agent.warnings + assert agent.touch_calls[0] == "context compression started" + assert agent.touch_calls[-1] == "context compression failed" + assert agent.status_events == [ + ("lifecycle", COMPACTION_STATUS), + ("warn", "⚠ Codex app-server compaction failed: compact failed"), + ] +