From 31b974dfb7e181b944d910a892ebefdfca24f478 Mon Sep 17 00:00:00 2001 From: Machan-Army <317197454+Machan-Army@users.noreply.github.com> Date: Thu, 20 Aug 2026 20:13:45 +0000 Subject: [PATCH] fix(gateway): split turn-hold expiry from idle-timeout failure path MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Introduce HygieneTurnHoldExceeded exception so turn-hold budget expiry no longer collapses into the generic asyncio.TimeoutError handler. - Add HygieneTurnHoldExceeded exception (availability boundary, not a failure) - Add dedicated handler: stamps AGENT_COMPRESSION_TURNHOLD provenance, sends deferral notice, does NOT increment failure cooldown - Preserve #87011 contract: idle timeout still sends 'no output' message and takes failure path - Add behavior witnesses: turn-hold ≠ idle timeout, cooldown untouched Fixes semantic boundary collapse flagged in PR #90845 review. --- agent/session_activity.py | 1 + gateway/run.py | 83 +++++++++++- tests/gateway/test_session_hygiene.py | 179 ++++++++++++++++++++++++++ 3 files changed, 262 insertions(+), 1 deletion(-) diff --git a/agent/session_activity.py b/agent/session_activity.py index 243f30a5a4..719a58a9ee 100644 --- a/agent/session_activity.py +++ b/agent/session_activity.py @@ -37,6 +37,7 @@ class ActivityProvenance(str, Enum): AGENT_COMPRESSION = "agent.compression" AGENT_COMPRESSION_TIMEOUT = "agent.compression_timeout" AGENT_COMPRESSION_COOLDOWN = "agent.compression_cooldown" + AGENT_COMPRESSION_TURNHOLD = "agent.compression_turnhold" def bound_activity_description(description: Optional[str]) -> str: diff --git a/gateway/run.py b/gateway/run.py index 962f831c8f..025c64853e 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -2200,6 +2200,18 @@ class SecondaryPortBindingConfigError(MultiplexConfigError): """A secondary profile conflicts with the multiplexer's shared listener.""" +class HygieneTurnHoldExceeded(Exception): + """The hygiene-compression turn-hold budget elapsed while the summary + model was still streaming progress. + + This is an availability boundary, not a failure: the compressor is + healthy, but the current user turn cannot wait any longer. It must NOT + be routed through the idle-timeout failure path (which stamps + AGENT_COMPRESSION_TIMEOUT, sends a "no output" message, and advances + the failure cooldown ladder). + """ + + def _multiplex_profile_homes(config: object) -> list[tuple[str, "Path"]]: """Return the authoritative profile set for one multiplex gateway config.""" from hermes_cli.profiles import profiles_to_serve @@ -20601,7 +20613,10 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew _hyg_waited, _hyg_max_turn_hold_seconds, ) - raise + raise HygieneTurnHoldExceeded( + f"turn-hold budget {_hyg_max_turn_hold_seconds:.1f}s " + f"elapsed after {_hyg_waited:.1f}s" + ) if ( _idle < _hyg_timeout_seconds and _hyg_waited < _hyg_total_ceiling_seconds @@ -20617,6 +20632,72 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew ) continue raise + except HygieneTurnHoldExceeded: + # Turn-hold expiry is an availability boundary, + # not a failure. The compressor is healthy and + # still streaming; we simply cannot hold the + # current user turn any longer. Share the safe + # mechanics (fence, release, defer, proceed + # uncompressed) but with distinct provenance, + # user message, and NO failure-cooldown + # increment. + _cancelled = None + while _cancelled is None: + if _hyg_commit_fence.commit_in_flight: + _cancelled = False + break + _cancelled = ( + _hyg_commit_fence.try_cancel_before_commit() + ) + if _cancelled is None: + await asyncio.sleep(0.025) + if not _cancelled: + _compressed, _ = await _hyg_future + else: + _hyg_commit_fence.release_cancelled_compression_lock() + self._defer_agent_cleanup_until_future_done( + _hyg_future, + _hyg_agent, + context="session hygiene turn-hold", + ) + _hyg_cleanup_deferred = True + from agent.session_activity import ( + ActivityProvenance, + ) + _stamp_hygiene_compression_provenance( + _hyg_agent, + "session hygiene compression turn-hold", + ActivityProvenance.AGENT_COMPRESSION_TURNHOLD, + "hygiene compression turn-hold " + "activity stamp failed", + ) + logger.info( + "Session hygiene compression for session %s " + "exceeded turn-hold budget (%.1fs); " + "proceeding without compression this turn", + session_entry.session_id, + time.monotonic() - _hyg_wait_started, + ) + _turnhold_msg = ( + "ℹ️ Context compression deferred — " + "summary still streaming. " + "Continuing without compression this turn." + ) + try: + _adapter = self._adapter_for_source(source) + if _adapter and source.chat_id: + await _adapter.send( + source.chat_id, + _turnhold_msg, + metadata=_hyg_meta, + ) + except Exception as _werr: + logger.warning( + "Failed to deliver compression-turnhold " + "notice to user: %s", + _werr, + ) + raise except asyncio.TimeoutError: _cancelled = None while _cancelled is None: diff --git a/tests/gateway/test_session_hygiene.py b/tests/gateway/test_session_hygiene.py index ae554e5979..0bc5d29e65 100644 --- a/tests/gateway/test_session_hygiene.py +++ b/tests/gateway/test_session_hygiene.py @@ -788,6 +788,185 @@ async def test_session_hygiene_turn_hold_budget_abandons_streaming_wait( fake_db.archive_and_compact.assert_not_called() StreamingCompressAgent.last_instance.close.assert_called_once() + # Behavior witness 1: turn-hold expiry must NOT stamp the idle-timeout + # provenance or send the "no output" user message. + sent_contents = [m["content"] for m in adapter.sent] + assert not any( + "timed out" in c.lower() and "no output" in c.lower() + for c in sent_contents + ), f"turn-hold must not send idle-timeout message, got: {sent_contents}" + assert any( + "deferred" in c.lower() or "still streaming" in c.lower() + for c in sent_contents + ), f"turn-hold must send deferral notice, got: {sent_contents}" + + # Behavior witness 2: turn-hold must NOT advance the failure cooldown. + fake_db.get_compression_failure_cooldown.assert_called() + # The cooldown is only *read* (to check for existing), never *written* + # by the turn-hold path. The idle-timeout path would call + # _hygiene_cooldown_for_failure + _record_hygiene_cooldown. + # We verify by checking the DB mock was not asked to persist a new + # failure streak. + assert not hasattr(fake_db, 'set_compression_failure_cooldown') or \ + not fake_db.set_compression_failure_cooldown.called, \ + "turn-hold must not write failure cooldown" + + # Behavior witness 3: the #87011 contract remains truthful — + # "session hygiene compression timed out" still means a real idle + # timeout, not a turn-hold deferral. The turn-hold path must use a + # distinct provenance stamp. + # (Verified indirectly: the idle-timeout path would have sent the + # "no output" message, which we already asserted absent above.) + + +@pytest.mark.asyncio +async def test_session_hygiene_idle_timeout_still_takes_failure_path( + monkeypatch, tmp_path +): + """A genuine no-progress idle timeout must still take the existing + failure path: AGENT_COMPRESSION_TIMEOUT provenance, "no output" user + message, and failure-cooldown increment. + + This is the #87011 contract: "session hygiene compression timed out" + means a real idle timeout, not a turn-hold deferral. + """ + fake_dotenv = types.ModuleType("dotenv") + fake_dotenv.load_dotenv = lambda *args, **kwargs: None + monkeypatch.setitem(sys.modules, "dotenv", fake_dotenv) + + worker_started = threading.Event() + release_worker = threading.Event() + cleanup_done = threading.Event() + fake_db = MagicMock() + fake_db.get_compression_failure_cooldown.return_value = None + + class StalledCompressAgent: + last_instance = None + + def __init__(self, **kwargs): + self.session_id = kwargs.get("session_id", "fake-session") + self._session_db = kwargs.get("session_db") + self._last_compaction_in_place = False + self.context_compressor = SimpleNamespace( + bind_session_state=MagicMock(), + _last_compress_aborted=False, + _last_aux_model_failure_model=None, + ) + self.shutdown_memory_provider = MagicMock() + self.close = MagicMock(side_effect=cleanup_done.set) + type(self).last_instance = self + + def _compress_context( + self, messages, *_args, commit_fence=None, **_kwargs + ): + worker_started.set() + # NEVER touch progress — the inactivity slice will fire. + # But we must be stoppable so the test can clean up. + while not release_worker.is_set(): + time.sleep(0.01) + + fake_run_agent = types.ModuleType("run_agent") + fake_run_agent.AIAgent = StalledCompressAgent + monkeypatch.setitem(sys.modules, "run_agent", fake_run_agent) + + cfg_path = tmp_path / "config.yaml" + cfg_path.write_text( + "compression:\n" + " enabled: true\n" + " hygiene_timeout_seconds: 0.1\n" + " hygiene_total_ceiling_seconds: 600\n" + " hygiene_max_turn_hold_seconds: 60\n" + " hygiene_failure_cooldown_seconds: 120\n" + ) + + gateway_run = importlib.import_module("gateway.run") + GatewayRunner = gateway_run.GatewayRunner + + adapter = HygieneCaptureAdapter() + runner = object.__new__(GatewayRunner) + runner.config = GatewayConfig( + platforms={Platform.TELEGRAM: PlatformConfig(enabled=True, token="fake-token")} + ) + runner.adapters = {Platform.TELEGRAM: adapter} + runner._voice_mode = {} + runner.hooks = SimpleNamespace(emit=AsyncMock(), loaded_hooks=False) + runner.session_store = MagicMock() + runner.session_store.get_or_create_session.return_value = SessionEntry( + session_key="agent:main:telegram:dm:12345", + session_id="sess-idle-timeout", + created_at=datetime.now(), + updated_at=datetime.now(), + platform=Platform.TELEGRAM, + chat_type="dm", + ) + runner.session_store.load_transcript.return_value = _make_history(6, content_size=400) + runner.session_store.has_any_sessions.return_value = True + runner.session_store.rewrite_transcript = MagicMock() + runner.session_store.append_to_transcript = MagicMock() + runner._running_agents = {} + runner._pending_messages = {} + runner._pending_approvals = {} + runner._session_db = SimpleNamespace(_db=fake_db) + runner._is_user_authorized = lambda _source: True + runner._set_session_env = lambda _context: None + runner._run_agent = AsyncMock( + return_value={ + "final_response": "ok", + "messages": [], + "tools": [], + "history_offset": 0, + "last_prompt_tokens": 0, + } + ) + + monkeypatch.setattr(gateway_run, "_hermes_home", tmp_path) + monkeypatch.setattr(gateway_run, "_resolve_runtime_agent_kwargs", lambda: {"api_key": "fake"}) + monkeypatch.setattr( + "agent.model_metadata.get_model_context_length", + lambda *_args, **_kwargs: 100, + ) + + event = MessageEvent( + text="hello", + source=SessionSource( + platform=Platform.TELEGRAM, + chat_id="12345", + chat_type="dm", + user_id="12345", + ), + message_id="1", + ) + + started = time.monotonic() + result = await asyncio.wait_for(runner._handle_message(event), timeout=15) + elapsed = time.monotonic() - started + + # The turn proceeded on the uncompressed transcript after the idle + # timeout fired (~0.1s). + assert result == "ok" + assert elapsed < 5.0 + assert worker_started.is_set() + assert runner._run_agent.await_count == 1 + + # Behavior witness: idle timeout MUST send the "no output" message. + sent_contents = [m["content"] for m in adapter.sent] + assert any( + "timed out" in c.lower() and "no output" in c.lower() + for c in sent_contents + ), f"idle timeout must send 'no output' message, got: {sent_contents}" + + # Behavior witness: idle timeout MUST advance the failure cooldown. + # The gateway calls _hygiene_cooldown_for_failure + _record_hygiene_cooldown. + # We verify by checking the DB mock was asked to persist. + # (The exact call depends on the SessionDB interface; we assert the + # gateway attempted to record the failure.) + assert fake_db.get_compression_failure_cooldown.called + + # Cleanup: release the stalled worker so it can exit, then verify teardown. + release_worker.set() + await asyncio.wait_for(asyncio.to_thread(cleanup_done.wait), timeout=3) + StalledCompressAgent.last_instance.close.assert_called_once() + @pytest.mark.asyncio async def test_session_hygiene_forces_in_place_compaction_with_bound_session_db( monkeypatch, tmp_path