From a995973d45f386d34221f8a2d4208da3c1898180 Mon Sep 17 00:00:00 2001 From: HexLab98 Date: Fri, 28 Aug 2026 14:51:46 +0700 Subject: [PATCH 1/2] fix(gateway): stop hygiene retry livelock after commit-fence cancel (#96953) A /stop or /restart abort left hygiene with no cooldown, so the next turn re-armed auto-compression and waited up to 600s behind a fence that would refuse the commit again. Record a cooldown on fence-cancel and unwind, stop extending that wait once the fence is cancelled, and skip a new hygiene agent while a compression lock is already held. --- gateway/run.py | 304 +++++++++++++++++++++++++++++++++++-------------- 1 file changed, 220 insertions(+), 84 deletions(-) diff --git a/gateway/run.py b/gateway/run.py index 3ce89c3b739e1..62fba594b41c8 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -442,6 +442,24 @@ async def run_codex_hygiene_compaction( # route already persisted its own failure cooldown in that case. return "failed:no-boundary" +def hygiene_wait_should_extend( + *, + idle: float, + timeout: float, + waited: float, + ceiling: float, + fence_cancelled: bool = False, +) -> bool: + """Whether the hygiene host should keep waiting for a slow summary. + + A cancelled commit fence cannot produce a commit (#96953): extending the + wait up to the 600s ceiling only queues inbound messages behind a doomed + attempt. Stop extending immediately so the turn can continue. + """ + if fence_cancelled: + return False + return idle < timeout and waited < ceiling + def _record_hygiene_cooldown( gateway, @@ -10566,7 +10584,10 @@ async def _session_has_compression_in_flight(self, session_key: str) -> bool: holder = await asyncio.to_thread( raw_db.get_compression_lock_holder, str(session_id) ) - return bool(holder) + # Production returns Optional[str]. Reject non-strings so a + # MagicMock auto-attr (or any unexpected truthy) cannot look + # like a held lock and skip hygiene (#96953). + return isinstance(holder, str) and bool(holder) except (AttributeError, TypeError): return False except Exception: @@ -20713,6 +20734,22 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g ) _needs_compress = False + if _needs_compress and await self._session_has_compression_in_flight( + session_key + ): + # A prior hygiene/agent compression still holds the + # durable lock (typically a shielded worker left behind + # by /stop or /restart). Starting another attempt waits + # up to the 600s ceiling behind a commit that the fence + # will refuse, while inbound messages demote to queue + # (#96953). + logger.info( + "Session hygiene: skipping compression for %s; " + "another compression is already in flight", + session_entry.session_id, + ) + _needs_compress = False + if _needs_compress: logger.info( "Session hygiene: %s messages, ~%s tokens (%s) — auto-compressing " @@ -20898,6 +20935,8 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g # the turn forever. _hyg_wait_started = time.monotonic() while True: + if _hyg_commit_fence.is_cancelled: + raise asyncio.TimeoutError # #76354 S3: charge the idle budget # from the LAST PROGRESS event, not # from the start of this wait slice — @@ -20943,6 +20982,15 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g _slice, max(_turn_hold_remaining, 0.005), ) + # Re-check the fence on a short poll so a + # /stop or /restart cancel is not stuck + # behind a full idle window (#96953). + _idle_left = max( + _hyg_timeout_seconds + - _hyg_commit_fence.seconds_since_progress(), + 0.005, + ) + _slice = min(_slice, 0.25) try: _compressed, _ = await asyncio.wait_for( asyncio.shield(_hyg_future), @@ -20950,6 +20998,8 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g ) break except asyncio.TimeoutError: + if _hyg_commit_fence.is_cancelled: + raise _hyg_waited = time.monotonic() - _hyg_wait_started _idle = _hyg_commit_fence.seconds_since_progress() # Bounded turn-hold (#TKT-0029): @@ -20984,19 +21034,23 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g 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 + if hygiene_wait_should_extend( + idle=_idle, + timeout=_hyg_timeout_seconds, + waited=_hyg_waited, + ceiling=_hyg_total_ceiling_seconds, + fence_cancelled=_hyg_commit_fence.is_cancelled, ): - logger.info( - "Session hygiene compression for " - "session %s still streaming after " - "%.0fs (last progress %.1fs ago) — " - "extending wait (ceiling %.0fs)", - session_entry.session_id, - _hyg_waited, _idle, - _hyg_total_ceiling_seconds, - ) + if _slice >= _idle_left - 1e-9: + logger.info( + "Session hygiene compression for " + "session %s still streaming after " + "%.0fs (last progress %.1fs ago) — " + "extending wait (ceiling %.0fs)", + session_entry.session_id, + _hyg_waited, _idle, + _hyg_total_ceiling_seconds, + ) continue raise except HygieneTurnHoldExceeded: @@ -21104,6 +21158,13 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g # established release behavior for a # provider call that may never return. _hyg_commit_fence.retain_compression_lock_until_worker_done() + # Capture fence state BEFORE try_cancel + # — that call itself sets is_cancelled, + # which would mis-label a genuine idle + # timeout as a fence cancel (#96953). + _hyg_fence_cancelled = ( + _hyg_commit_fence.is_cancelled + ) _cancelled = None while _cancelled is None: # #76354 F1: a hung commit retains the @@ -21143,6 +21204,16 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g context="session hygiene timeout", ) _hyg_cleanup_deferred = True + _hyg_timeout_error = ( + "session hygiene compression " + "cancelled at commit fence" + if _hyg_fence_cancelled + else ( + "session hygiene compression " + "timed out with no output from " + "the summary model" + ) + ) if _hyg_failure_cooldown_seconds >= 0: _hyg_cooldown = await asyncio.to_thread( _hygiene_cooldown_for_failure, @@ -21151,12 +21222,16 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g _hyg_failure_cooldown_seconds, ) _timeout_reason = ( - "session hygiene compression total " - "ceiling exhausted" - if _hyg_total_exhausted - else "session hygiene compression " - "timed out with no output from the " - "summary model" + _hyg_timeout_error + if _hyg_fence_cancelled + else ( + "session hygiene compression total " + "ceiling exhausted" + if _hyg_total_exhausted + else "session hygiene compression " + "timed out with no output from the " + "summary model" + ) ) _record_hygiene_cooldown( self, session_entry.session_id, @@ -21168,59 +21243,73 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g ) _stamp_hygiene_compression_provenance( _hyg_agent, - "session hygiene compression timed out", + ( + "session hygiene compression " + "cancelled at commit fence" + if _hyg_fence_cancelled + else "session hygiene compression timed out" + ), ActivityProvenance.AGENT_COMPRESSION_TIMEOUT, "hygiene compression timeout " "activity stamp failed", ) - _hyg_elapsed = ( - time.monotonic() - _hyg_wait_started - ) - if _hyg_total_exhausted: + if _hyg_fence_cancelled: logger.warning( - "Session hygiene compression for session %s " - "reached its total ceiling after %.1fs " - "(progress observed=%s); continuing without " + "Session hygiene compression for " + "session %s was cancelled at the " + "commit fence; continuing without " "compression", session_entry.session_id, - _hyg_elapsed, - _hyg_commit_fence.progress_observed, ) else: - logger.warning( - "Session hygiene compression for session %s " - "made no progress for %.1fs (total wait " - "%.1fs, ceiling %.1fs); continuing without " - "compression", - session_entry.session_id, - _hyg_commit_fence.seconds_since_progress(), - _hyg_elapsed, - _hyg_total_ceiling_seconds, - ) - _timeout_msg = ( - _hygiene_compression_timeout_message( - total_exhausted=_hyg_total_exhausted, - elapsed=_hyg_elapsed, - idle_timeout=_hyg_timeout_seconds, - progress_observed=( - _hyg_commit_fence.progress_observed - ), + _hyg_elapsed = ( + time.monotonic() - _hyg_wait_started ) - ) - try: - _adapter = self._adapter_for_source(source) - if _adapter and source.chat_id: - await _adapter.send( - source.chat_id, - _timeout_msg, - metadata=_hyg_meta, + if _hyg_total_exhausted: + logger.warning( + "Session hygiene compression for session %s " + "reached its total ceiling after %.1fs " + "(progress observed=%s); continuing without " + "compression", + session_entry.session_id, + _hyg_elapsed, + _hyg_commit_fence.progress_observed, + ) + else: + logger.warning( + "Session hygiene compression for session %s " + "made no progress for %.1fs (total wait " + "%.1fs, ceiling %.1fs); continuing without " + "compression", + session_entry.session_id, + _hyg_commit_fence.seconds_since_progress(), + _hyg_elapsed, + _hyg_total_ceiling_seconds, + ) + _timeout_msg = ( + _hygiene_compression_timeout_message( + total_exhausted=_hyg_total_exhausted, + elapsed=_hyg_elapsed, + idle_timeout=_hyg_timeout_seconds, + progress_observed=( + _hyg_commit_fence.progress_observed + ), ) - except Exception as _werr: - logger.warning( - "Failed to deliver compression-timeout " - "warning to user: %s", - _werr, ) + try: + _adapter = self._adapter_for_source(source) + if _adapter and source.chat_id: + await _adapter.send( + source.chat_id, + _timeout_msg, + metadata=_hyg_meta, + ) + except Exception as _werr: + logger.warning( + "Failed to deliver compression-timeout " + "warning to user: %s", + _werr, + ) raise except BaseException: # #76354 F2: non-timeout unwind while the @@ -21239,6 +21328,30 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g context="session hygiene unwind", ) _hyg_cleanup_deferred = True + # #96953: restart drain / task cancel used + # to re-raise with no cooldown, so the + # next turn immediately re-armed hygiene + # and waited up to 600s behind a fence + # that would refuse the commit again. + if _hyg_failure_cooldown_seconds >= 0: + try: + _hyg_cooldown = _hygiene_cooldown_for_failure( + self, + session_key, + _hyg_failure_cooldown_seconds, + ) + _record_hygiene_cooldown( + self, session_entry.session_id, + _hyg_cooldown, + "session hygiene compression " + "cancelled at commit fence", + ) + except Exception as _cd_err: + logger.debug( + "hygiene unwind cooldown " + "record failed: %s", + _cd_err, + ) raise # _compress_context ends the old session and creates @@ -21387,6 +21500,23 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g _hyg_aborted = _comp is not None and getattr( _comp, "_last_compress_aborted", False ) + # Fence-cancelled _compress_context returns + # the original transcript with + # _last_compress_aborted still False + # (failure_class=commit_fence_cancelled, + # chunk_count=0). Treat that no-op as an + # abort so hygiene records a cooldown + # instead of retrying into the 600s wait + # (#96953). A successful rotate/in-place + # commit is not an abort even if a later + # invalidation flipped the fence. + _hyg_fence_cancelled = bool( + _hyg_commit_fence.is_cancelled + and not _hyg_rotated + and not _hyg_in_place + ) + if _hyg_fence_cancelled: + _hyg_aborted = True if not _hyg_aborted: # Recovery decision lives in the # extracted, unit-tested predicate — the @@ -21421,8 +21551,13 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g _record_hygiene_cooldown( self, session_entry.session_id, _hyg_cooldown, - getattr( - _comp, "_last_summary_error", None + ( + "session hygiene compression " + "cancelled at commit fence" + if _hyg_fence_cancelled + else getattr( + _comp, "_last_summary_error", None + ) ), ) from agent.session_activity import ( @@ -21435,29 +21570,30 @@ async def _handle_message_with_agent(self, event, source, _quick_key: str, run_g "hygiene compression abort " "activity stamp failed", ) - _err = getattr(_comp, "_last_summary_error", None) or "unknown error" - # Force-redact: provider exception text - # may contain credentials; this message - # reaches gateway users directly. - from agent.redact import redact_sensitive_text - _err = redact_sensitive_text(_err, force=True) - _warn_msg = ( - "⚠️ Context compression aborted " - f"({_err}). No messages were dropped — " - "conversation is unchanged. Run /compress " - "to retry, /reset for a clean session, or " - "check your auxiliary.compression model " - "configuration." - ) - try: - _adapter = self._adapter_for_source(source) - if _adapter and source.chat_id: - await _adapter.send(source.chat_id, _warn_msg, metadata=_hyg_meta) - except Exception as _werr: - logger.warning( - "Failed to deliver compression-failure warning to user: %s", - _werr, + if not _hyg_fence_cancelled: + _err = getattr(_comp, "_last_summary_error", None) or "unknown error" + # Force-redact: provider exception text + # may contain credentials; this message + # reaches gateway users directly. + from agent.redact import redact_sensitive_text + _err = redact_sensitive_text(_err, force=True) + _warn_msg = ( + "⚠️ Context compression aborted " + f"({_err}). No messages were dropped — " + "conversation is unchanged. Run /compress " + "to retry, /reset for a clean session, or " + "check your auxiliary.compression model " + "configuration." ) + try: + _adapter = self._adapter_for_source(source) + if _adapter and source.chat_id: + await _adapter.send(source.chat_id, _warn_msg, metadata=_hyg_meta) + except Exception as _werr: + logger.warning( + "Failed to deliver compression-failure warning to user: %s", + _werr, + ) # Separately: if the user's CONFIGURED aux # model failed and we recovered by falling # back to the main model, tell them — a From 806e13d92309111f106da58151259d6f2120b7e5 Mon Sep 17 00:00:00 2001 From: HexLab98 Date: Fri, 28 Aug 2026 14:51:55 +0700 Subject: [PATCH 2/2] test(gateway): cover hygiene fence-cancel cooldown and in-flight skip (#96953) Prove a fence-cancelled helper (no abort flag) persists cooldown so the next turn does not re-arm compression, the host does not wait out the 600s ceiling after cancel, a held lock skips the sibling agent, and unwind cancellation records the same brake. (cherry picked from commit d2e178cd96fbcc8b2baf68488a8f46c70bff31a2) --- .../test_compression_in_flight_check.py | 14 + .../test_hygiene_failure_cooldown_ladder.py | 26 ++ tests/gateway/test_session_hygiene.py | 258 ++++++++++++++++++ 3 files changed, 298 insertions(+) diff --git a/tests/gateway/test_compression_in_flight_check.py b/tests/gateway/test_compression_in_flight_check.py index d994777f8a3e5..b1d51230e97a8 100644 --- a/tests/gateway/test_compression_in_flight_check.py +++ b/tests/gateway/test_compression_in_flight_check.py @@ -48,6 +48,20 @@ async def test_returns_false_when_no_session_store(): assert await runner._session_has_compression_in_flight("k") is False +@pytest.mark.asyncio +async def test_returns_false_when_holder_is_not_a_string(): + """Lock holders are session-id strings. A MagicMock auto-attr must not + look like an in-flight compression and skip hygiene (#96953).""" + runner = _make_runner(holder_value=MagicMock()) + assert await runner._session_has_compression_in_flight("k") is False + runner = _make_runner(holder_value=True) + assert await runner._session_has_compression_in_flight("k") is False + runner = _make_runner(holder_value="") + assert await runner._session_has_compression_in_flight("k") is False + runner = _make_runner(holder_value="agent-1") + assert await runner._session_has_compression_in_flight("k") is True + + @pytest.mark.asyncio async def test_db_call_runs_off_event_loop(): """Regression core: get_compression_lock_holder MUST execute in non-event-loop thread.""" diff --git a/tests/gateway/test_hygiene_failure_cooldown_ladder.py b/tests/gateway/test_hygiene_failure_cooldown_ladder.py index d8811dca3a276..0e9b81ded867f 100644 --- a/tests/gateway/test_hygiene_failure_cooldown_ladder.py +++ b/tests/gateway/test_hygiene_failure_cooldown_ladder.py @@ -26,6 +26,7 @@ _record_hygiene_cooldown, _reset_hygiene_failure_streak, hygiene_compaction_recovered, + hygiene_wait_should_extend, ) from gateway.run import GatewayRunner from gateway.session_state import PersistentState, SessionState @@ -373,6 +374,31 @@ def test_sub_five_percent_wobble_is_not_recovery(self): ) is False +class TestHygieneWaitShouldExtend: + """Host must not keep waiting after the commit fence is already cancelled.""" + + def test_extends_while_idle_and_under_ceiling(self): + assert hygiene_wait_should_extend( + idle=1.0, timeout=30.0, waited=10.0, ceiling=600.0, + ) is True + + def test_stops_when_idle_budget_exhausted(self): + assert hygiene_wait_should_extend( + idle=30.0, timeout=30.0, waited=10.0, ceiling=600.0, + ) is False + + def test_stops_at_ceiling(self): + assert hygiene_wait_should_extend( + idle=1.0, timeout=30.0, waited=600.0, ceiling=600.0, + ) is False + + def test_fence_cancel_stops_even_with_fresh_progress(self): + assert hygiene_wait_should_extend( + idle=0.0, timeout=30.0, waited=0.1, ceiling=600.0, + fence_cancelled=True, + ) is False + + # --------------------------------------------------------------------------- # Integration with the persist helper # --------------------------------------------------------------------------- diff --git a/tests/gateway/test_session_hygiene.py b/tests/gateway/test_session_hygiene.py index f09b1e60c99cf..1bb266b05b540 100644 --- a/tests/gateway/test_session_hygiene.py +++ b/tests/gateway/test_session_hygiene.py @@ -1556,3 +1556,261 @@ def _compress_context(self, messages, *_args, **_kwargs): assert escalated["remaining_seconds"] == pytest.approx(900, abs=5) finally: db.close() + + +# --------------------------------------------------------------------------- +# Commit-fence cancel must not livelock hygiene (#96953) +# --------------------------------------------------------------------------- + +@pytest.mark.asyncio +async def test_hygiene_fence_cancel_records_cooldown_without_abort_flag( + monkeypatch, tmp_path +): + """A fence-cancelled hygiene worker returns the original transcript with + ``_last_compress_aborted`` still False (failure_class=commit_fence_cancelled). + + That used to skip the abort-cooldown block, so the next turn immediately + re-armed hygiene and waited up to the 600s ceiling behind a doomed attempt. + """ + from hermes_state import SessionDB + + gateway_run = importlib.import_module("gateway.run") + session_id = "sess-fence-cancel" + + class FenceCancelCompressAgent: + instances = 0 + + def __init__(self, **kwargs): + type(self).instances += 1 + self.session_id = kwargs.get("session_id", session_id) + 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_summary_error=None, + _last_aux_model_failure_model=None, + ) + self.shutdown_memory_provider = MagicMock() + self.close = MagicMock() + + def _compress_context(self, messages, *_args, commit_fence=None, **_kwargs): + if commit_fence is not None: + assert commit_fence.try_cancel_before_commit() is True + return (messages, None) + + db = SessionDB(db_path=tmp_path / "state.db") + try: + db.create_session(session_id, "telegram") + runner1, adapter1, event1 = _make_cooldown_runner( + monkeypatch, tmp_path, FenceCancelCompressAgent, db, session_id + ) + assert await runner1._handle_message(event1) == "ok" + assert FenceCancelCompressAgent.instances == 1 + state = db.get_compression_failure_cooldown(session_id) + assert state is not None and state["remaining_seconds"] > 0, ( + "fence-cancelled hygiene compression did not persist a cooldown; " + f"got {state!r}" + ) + assert not any( + "Context compression aborted" in s["content"] for s in adapter1.sent + ), "fence-cancel during /stop or /restart must not toast an abort" + + class ShouldNotRunAgent: + instances = 0 + + def __init__(self, **kwargs): + type(self).instances += 1 + 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() + + def _compress_context(self, messages, *_args, **_kwargs): + return (messages, None) + + runner2, _adapter2, event2 = _make_cooldown_runner( + monkeypatch, tmp_path, ShouldNotRunAgent, db, session_id + ) + assert await runner2._handle_message(event2) == "ok" + assert ShouldNotRunAgent.instances == 0, ( + "REGRESSION (#96953): hygiene re-armed after a commit-fence " + "cancel instead of honoring the failure cooldown" + ) + assert runner2._run_agent.await_count == 1 + finally: + db.close() + + +@pytest.mark.asyncio +async def test_hygiene_does_not_wait_ceiling_after_fence_cancel( + monkeypatch, tmp_path +): + """Once the commit fence is cancelled, the host must stop extending the + wait — even if the shielded worker is still alive and touching progress. + """ + from hermes_state import SessionDB + + worker_started = threading.Event() + release_worker = threading.Event() + cleanup_done = threading.Event() + session_id = "sess-fence-wait" + + class HungAfterFenceCancelAgent: + last_instance = None + + def __init__(self, **kwargs): + self.session_id = kwargs.get("session_id", session_id) + 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 + ): + if commit_fence is not None: + commit_fence.try_cancel_before_commit() + worker_started.set() + # Keep the worker alive (and keep reporting "progress") so a + # host that still extends to the 600s ceiling would stall here. + deadline = time.monotonic() + 2.0 + while time.monotonic() < deadline: + if commit_fence is not None: + commit_fence.touch_progress() + if release_worker.is_set(): + break + time.sleep(0.02) + return (messages, None) + + db = SessionDB(db_path=tmp_path / "state.db") + try: + db.create_session(session_id, "telegram") + runner, adapter, event = _make_cooldown_runner( + monkeypatch, tmp_path, HungAfterFenceCancelAgent, db, session_id + ) + started = time.monotonic() + result = await runner._handle_message(event) + elapsed = time.monotonic() - started + + assert result == "ok" + assert worker_started.wait(timeout=2) + assert elapsed < 2.0, ( + f"hygiene host waited {elapsed:.1f}s after fence cancel — " + "must not extend toward the 600s ceiling (#96953)" + ) + assert runner._run_agent.await_count == 1 + state = db.get_compression_failure_cooldown(session_id) + assert state is not None and state["remaining_seconds"] > 0 + assert not any( + "Context compression timed out" in s["content"] for s in adapter.sent + ), "fence-cancel is not a summary-model timeout; no timeout toast" + release_worker.set() + await asyncio.wait_for(asyncio.to_thread(cleanup_done.wait), timeout=2) + finally: + db.close() + + +@pytest.mark.asyncio +async def test_hygiene_skips_when_compression_already_in_flight( + monkeypatch, tmp_path +): + """Do not spawn a sibling hygiene compressor while a lock is already held.""" + from hermes_state import SessionDB + + session_id = "sess-in-flight" + + class ShouldNotRunAgent: + instances = 0 + + def __init__(self, **kwargs): + type(self).instances += 1 + 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() + + def _compress_context(self, messages, *_args, **_kwargs): + return (messages, None) + + db = SessionDB(db_path=tmp_path / "state.db") + try: + db.create_session(session_id, "telegram") + runner, _adapter, event = _make_cooldown_runner( + monkeypatch, tmp_path, ShouldNotRunAgent, db, session_id + ) + runner._session_has_compression_in_flight = AsyncMock(return_value=True) + assert await runner._handle_message(event) == "ok" + assert ShouldNotRunAgent.instances == 0 + assert runner._run_agent.await_count == 1 + finally: + db.close() + + +@pytest.mark.asyncio +async def test_hygiene_unwind_records_cooldown(monkeypatch, tmp_path): + """Restart-drain cancellation must persist a cooldown before re-raising. + + ``except BaseException`` used to revoke the fence and re-raise with no + cooldown, so the next turn after /restart re-triggered hygiene immediately. + """ + from hermes_state import SessionDB + + worker_started = threading.Event() + release_worker = threading.Event() + cleanup_done = threading.Event() + session_id = "sess-unwind" + + class SlowCompressAgent: + last_instance = None + + def __init__(self, **kwargs): + self.session_id = kwargs.get("session_id", session_id) + 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, **_kwargs): + worker_started.set() + release_worker.wait(timeout=5) + return (messages, None) + + db = SessionDB(db_path=tmp_path / "state.db") + try: + db.create_session(session_id, "telegram") + runner, _adapter, event = _make_cooldown_runner( + monkeypatch, tmp_path, SlowCompressAgent, db, session_id + ) + task = asyncio.create_task(runner._handle_message(event)) + assert await asyncio.to_thread(worker_started.wait, 2) + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + state = db.get_compression_failure_cooldown(session_id) + assert state is not None and state["remaining_seconds"] > 0, ( + "hygiene unwind did not persist a cooldown; got " + f"{state!r}" + ) + release_worker.set() + await asyncio.wait_for(asyncio.to_thread(cleanup_done.wait), timeout=2) + finally: + db.close()