diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 5c9d210f97..f48c8b0b76 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -718,6 +718,17 @@ class CompressionCommitFence: self._progress_observed = False self._deadline: float | None = None self._retain_cancelled_lock_until_worker_done = False + # #97963: set by the worker (mark_commit_watermark_fenced) once its + # commit path is watermark-fenced — i.e. it captured the session's + # active-row watermark at compression start, so any row appended + # AFTER that point survives a late commit verbatim as concurrent + # tail (archive_and_compact / publish_compression_child clone rows + # above the watermark instead of archiving them). Hosts read this + # at the turn-hold boundary to decide whether a detached worker may + # KEEP its commit admission (safe: newer turns cannot be clobbered) + # or must be cancelled as before (unfenced commit; discard is the + # only safe outcome). Plain bool store — atomic in CPython. + self._commit_watermark_fenced = False if total_ceiling_seconds is not None: self.set_total_ceiling_seconds(total_ceiling_seconds) @@ -857,6 +868,24 @@ class CompressionCommitFence: """Prevent a timed-out live worker from overlapping a retry.""" self._retain_cancelled_lock_until_worker_done = True + def mark_commit_watermark_fenced(self) -> None: + """Record that this attempt's commit is bounded by a start watermark. + + Called by the compression worker right after it captures + ``get_active_message_watermark()`` under the durable compression + lock (#75316/#87484). A watermark-fenced commit archives ONLY rows + at or below the watermark; rows appended later — e.g. the user turn + the host released at the turn-hold boundary (#97963) — are cloned + as live concurrent tail. That is exactly the property a host needs + before letting a detached worker keep its commit admission. + """ + self._commit_watermark_fenced = True + + @property + def commit_watermark_fenced(self) -> bool: + """Lock-free read: the worker's commit is watermark-bounded.""" + return self._commit_watermark_fenced + def allow_cancelled_lock_release(self) -> None: """Undo :meth:`retain_compression_lock_until_worker_done`. @@ -3530,6 +3559,18 @@ def compress_context( _commit_watermark = _lock_db.get_active_message_watermark( _lock_sid ) + # #97963: a captured watermark makes the eventual + # commit safe against rows appended after this + # point (they survive as cloned concurrent tail on + # BOTH commit paths — archive_and_compact and + # publish_compression_child). Tell the fence so a + # host at the turn-hold boundary can keep this + # attempt's commit admission instead of burning it. + if commit_fence is not None: + try: + commit_fence.mark_commit_watermark_fenced() + except AttributeError: + pass # test doubles without the method except Exception as _wm_err: # Watermark capture is safety-additive: without it the # commit falls back to archive-everything (historical diff --git a/cli-config.yaml.example b/cli-config.yaml.example index c45203cfcb..de265fa619 100644 --- a/cli-config.yaml.example +++ b/cli-config.yaml.example @@ -699,6 +699,19 @@ compression: # summarization on a short idle thread. Example: 1800 = compact after 30 min idle. idle_compact_after_seconds: 0 + # Gateway session-hygiene turn-hold budget (default: 10). Max seconds an + # arriving user turn is held while a still-streaming hygiene summary + # finishes. Distinct from hygiene_timeout_seconds (compressor inactivity + # budget): this bounds user-visible latency so chat transports (Telegram + # ~30s) do not drop a silent connection. On expiry the turn proceeds + # uncompressed; the detached worker keeps its commit admission (when the + # commit is watermark-fenced) and the summary is adopted at the next safe + # boundary. Thinking-model summarizers often need longer than 10s to emit + # the first content token — raise to 300 (or >= your summarizer's real + # time-to-first-content) only if you want THIS turn to wait for the + # compression instead of adopting it one turn late. + hygiene_max_turn_hold_seconds: 10 + # Proactive tool-result prune (default: 0 = disabled). Opt-in token trigger # for a deterministic, no-LLM prune of OLD tool-result payloads, run # independently of `threshold` above. On large-window models (512K/1M) the diff --git a/gateway/run.py b/gateway/run.py index 11d43932b2..93ef09b73d 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -21472,6 +21472,175 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew # uncompressed) but with distinct provenance, # user message, and NO failure-cooldown # increment. + # + # #97963: decouple the TURN from the + # COMPRESSION. When the worker's commit is + # watermark-fenced (it captured the session's + # active-row watermark at compression start, + # so rows appended after that point — this + # released turn included — survive its late + # commit verbatim as cloned concurrent tail), + # the already-running attempt KEEPS its commit + # admission: the user's turn proceeds on the + # uncompressed transcript NOW, and the summary + # is adopted when the detached worker reaches + # its own watermark-fenced commit transaction + # (archive_and_compact / the rotation publish + # path — the next safe boundary). Before this, + # the fence was ALWAYS cancelled here, burning + # the full summary attempt — for a thinking + # summary model whose reasoning prefix alone + # exceeds the 10s hold, that made hygiene + # auto-compression fail 100% of the time while + # paying the summary model per turn. The turn + # itself is still released at the same budget: + # only the fate of the detached worker's + # RESULT changes. If the commit is NOT + # watermark-fenced (no session_db, watermark + # capture failed, legacy lock API), a late + # commit could clobber newer turns, so cancel + # exactly as before — never worse than the + # status quo. + _hyg_keep_admission = bool( + getattr( + _hyg_commit_fence, + "commit_watermark_fenced", + False, + ) + ) and not _hyg_commit_fence.is_cancelled + if _hyg_keep_admission: + self._defer_agent_cleanup_until_future_done( + _hyg_future, + _hyg_agent, + context="session hygiene turn-hold", + ) + _hyg_cleanup_deferred = True + # NO retry-after here (#97963 (b)): the + # attempt is still running toward a real + # commit, and arming the flat 60s + # retry-after would ALSO block the + # agent-side preflight compressor from a + # fresh chance ("Skipping preflight + # compression: same-session cooldown + # active"). Re-attempt spacing is covered + # by the durable compression lock instead: + # the next turn's hygiene pre-check skips + # while this worker's lease is held + # (_session_has_compression_in_flight). + # The flat retry-after is recorded by the + # done-callback below ONLY if the worker + # ends without committing anything. + _hyg_deferred_sid = session_entry.session_id + _hyg_deferred_key = session_key + _hyg_deferred_agent = _hyg_agent + + def _hyg_adopt_or_space_retry( + _fut, + _gw=self, + _sid=_hyg_deferred_sid, + _skey=_hyg_deferred_key, + _agent=_hyg_deferred_agent, + ): + try: + _exc = _fut.exception() + except ( + asyncio.CancelledError, + Exception, + ): + _exc = None + _committed = False + else: + _committed = _exc is None and ( + bool( + getattr( + _agent, + "_last_compaction_in_place", + False, + ) + ) + or getattr( + _agent, "session_id", _sid + ) + != _sid + ) + if _committed: + logger.info( + "Session hygiene compression for " + "session %s finished after the " + "turn-hold was released — summary " + "adopted at the watermark-fenced " + "commit boundary (#97963)", + _sid, + ) + try: + _reset_hygiene_failure_streak( + _gw, _skey + ) + except Exception as _rs_err: + logger.debug( + "hygiene streak reset after " + "deferred adoption failed: %s", + _rs_err, + ) + else: + # Nothing to adopt (summary failed, + # fence refused the commit, or the + # attempt was superseded). Restore + # the pre-#97963 spacing so + # sustained traffic does not spawn + # and abandon a fresh compressor + # every turn. Flat and + # non-escalating: the streak must + # not advance for a deferral. + _record_hygiene_cooldown( + _gw, _sid, + _HYGIENE_TURNHOLD_RETRY_SECONDS, + "hygiene compression deferred: " + "turn-hold budget expired and the " + "detached attempt did not commit", + ) + + _hyg_future.add_done_callback( + _hyg_adopt_or_space_retry + ) + 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 — " + "the watermark-fenced worker keeps its " + "commit admission and the summary will be " + "adopted when it finishes", + session_entry.session_id, + time.monotonic() - _hyg_wait_started, + ) + _turnhold_msg = t( + "gateway.compress.turnhold_deferred" + ) + 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 _cancelled = None while _cancelled is None: if _hyg_commit_fence.commit_in_flight: @@ -22038,6 +22207,14 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew _hyg_agent, context="session hygiene" ) + except HygieneTurnHoldExceeded: + # Availability boundary, not a failure — already logged + # at INFO by the turn-hold handler. Must not hit the + # generic "auto-compress failed" warning below: that + # log is how thinking-model deployments read as + # permanently broken (#97963; surfaced by @686f6c61 + # in PR #99657). + pass except Exception as e: logger.warning( "Session hygiene auto-compress failed: %s", e diff --git a/hermes_cli/config_defaults.py b/hermes_cli/config_defaults.py index 2e64ec06e6..c6d72c5ac1 100644 --- a/hermes_cli/config_defaults.py +++ b/hermes_cli/config_defaults.py @@ -957,6 +957,10 @@ DEFAULT_CONFIG = { # waiting. Kept well under chat-transport idle timeouts # (Telegram ~30s). On expiry the turn proceeds # uncompressed — an availability boundary, not a failure. + # The detached worker keeps its commit admission when its + # commit is watermark-fenced, so the finished summary is + # adopted at the next safe boundary instead of being + # discarded (#97963 — thinking summary models). "context_timeout_seconds": 120, # inactivity budget for in-agent compress_context # (conversation loop, /compress, preflight, etc.). # Same progress-aware semantics as hygiene_timeout_seconds: diff --git a/tests/gateway/test_session_hygiene_turnhold_adoption.py b/tests/gateway/test_session_hygiene_turnhold_adoption.py new file mode 100644 index 0000000000..1b6c209cbd --- /dev/null +++ b/tests/gateway/test_session_hygiene_turnhold_adoption.py @@ -0,0 +1,431 @@ +"""Regression tests for #97963 — hygiene turn-hold must not burn a +watermark-fenced compression attempt. + +The 10s ``hygiene_max_turn_hold_seconds`` budget (#92318) releases the +arriving user turn while a thinking summary model is still streaming its +reasoning prefix. Before the fix, that release ALWAYS cancelled the commit +fence, so 100% of the summary attempt (including the full thinking prefix) +was discarded on every turn — auto-compression permanently failed for any +deployment whose summary model thinks longer than the hold. + +The fix decouples the turn from the compression: when the worker's commit is +watermark-fenced (rows appended after compression start survive its commit +verbatim as concurrent tail), the detached worker KEEPS its commit admission +and the summary is adopted at its own watermark-fenced commit boundary. The +turn is still released at the same budget — the invariant pinned by +``test_session_hygiene_turn_hold_budget_abandons_streaming_wait`` (#90845) +is untouched (that test's worker is NOT watermark-fenced and still takes the +cancel path). +""" + +import asyncio +import importlib +import sys +import threading +import time +import types +from datetime import datetime +from types import SimpleNamespace +from unittest.mock import AsyncMock, MagicMock + +import pytest + +from gateway.config import GatewayConfig, Platform, PlatformConfig +from gateway.platforms.base import BasePlatformAdapter, MessageEvent, SendResult +from gateway.session import SessionEntry, SessionSource + + +def _make_history(n_messages: int, content_size: int = 100) -> list: + history = [] + content = "x" * content_size + for i in range(n_messages): + role = "user" if i % 2 == 0 else "assistant" + history.append({"role": role, "content": content, "timestamp": f"t{i}"}) + return history + + +class _CaptureAdapter(BasePlatformAdapter): + def __init__(self): + super().__init__( + PlatformConfig(enabled=True, token="fake-token"), Platform.TELEGRAM + ) + self.sent = [] + + async def connect(self, *, is_reconnect: bool = False) -> bool: + return True + + async def disconnect(self) -> None: + return None + + async def send(self, chat_id, content, reply_to=None, metadata=None): + self.sent.append({"chat_id": chat_id, "content": content}) + return SendResult(success=True, message_id="x") + + async def get_chat_info(self, chat_id: str): + return {"id": chat_id} + + +def _write_turnhold_config(tmp_path): + cfg_path = tmp_path / "config.yaml" + cfg_path.write_text( + "compression:\n" + " enabled: true\n" + " hygiene_timeout_seconds: 60\n" + " hygiene_total_ceiling_seconds: 600\n" + " hygiene_max_turn_hold_seconds: 0.3\n" + " hygiene_failure_cooldown_seconds: 120\n" + ) + + +def _build_runner(gateway_run, adapter, fake_db): + runner = object.__new__(gateway_run.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-97963", + 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, + } + ) + return runner + + +def _make_event(): + return MessageEvent( + text="hello", + source=SessionSource( + platform=Platform.TELEGRAM, + chat_id="12345", + chat_type="dm", + user_id="12345", + ), + message_id="1", + ) + + +def _install_fakes(monkeypatch, gateway_run, tmp_path, agent_cls): + fake_dotenv = types.ModuleType("dotenv") + fake_dotenv.load_dotenv = lambda *args, **kwargs: None + monkeypatch.setitem(sys.modules, "dotenv", fake_dotenv) + fake_run_agent = types.ModuleType("run_agent") + fake_run_agent.AIAgent = agent_cls + monkeypatch.setitem(sys.modules, "run_agent", fake_run_agent) + 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, + ) + + +async def _drain_deferred(runner, timeout=10.0): + tasks = getattr(runner, "_deferred_agent_cleanup_tasks", None) or set() + if tasks: + await asyncio.wait_for( + asyncio.gather(*list(tasks), return_exceptions=True), timeout + ) + + +@pytest.mark.asyncio +async def test_turn_hold_keeps_admission_and_adopts_watermark_fenced_summary( + monkeypatch, tmp_path +): + """A watermark-fenced worker keeps its commit admission at turn-hold + expiry; its late summary is ADOPTED (committed), not discarded — while + the turn itself is still released at the budget (#90845 invariant). + """ + worker_started = threading.Event() + release_worker = threading.Event() + committed = threading.Event() + cleanup_done = threading.Event() + fake_db = MagicMock() + fake_db.get_compression_failure_cooldown.return_value = None + + class FencedStreamingAgent: + last_instance = None + + def __init__(self, **kwargs): + self.session_id = kwargs.get("session_id", "sess-97963") + 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 + ): + # Real compress_context marks the fence right after capturing + # the active-row watermark under the durable compression lock. + if commit_fence is not None: + commit_fence.mark_commit_watermark_fenced() + worker_started.set() + # Thinking-model shape: continuous progress, no commit yet — + # only the turn-hold budget can release the waiting turn. + # Bounded spin: a failing assertion before release_worker.set() + # must not leave this executor thread alive forever (pytest + # would hang at interpreter exit joining executor threads). + _spin_started = time.monotonic() + while not release_worker.is_set(): + if time.monotonic() - _spin_started > 20: + return (messages, None) + if commit_fence is not None: + commit_fence.touch_progress() + time.sleep(0.01) + if commit_fence is not None and not commit_fence.begin_commit(): + return (messages, None) + try: + self._session_db.archive_and_compact( + self.session_id, + [{"role": "assistant", "content": "summary"}], + watermark=6, + ) + self._last_compaction_in_place = True + committed.set() + return ([{"role": "assistant", "content": "summary"}], None) + finally: + if commit_fence is not None: + commit_fence.finish_commit() + + gateway_run = importlib.import_module("gateway.run") + _write_turnhold_config(tmp_path) + _install_fakes(monkeypatch, gateway_run, tmp_path, FencedStreamingAgent) + + adapter = _CaptureAdapter() + runner = _build_runner(gateway_run, adapter, fake_db) + + started = time.monotonic() + result = await asyncio.wait_for(runner._handle_message(_make_event()), timeout=15) + elapsed = time.monotonic() - started + + # #90845/#92318 invariant intact: the turn is released at the budget. + assert result == "ok" + assert elapsed < 5.0, f"turn held for {elapsed:.1f}s despite the turn-hold budget" + assert worker_started.is_set() + assert runner._run_agent.await_count == 1 + + # (b) NO retry-after was armed while the attempt is still running — + # arming it would block the agent-side preflight from adopting the + # finished summary ("same-session cooldown active", #97963). + assert not fake_db.record_compression_failure_cooldown.called, ( + "keep-admission path must not arm the retry-after while the " + "detached attempt is still running" + ) + + # The detached worker finishes late; its commit is ADMITTED (adoption), + # not refused — the summary attempt is no longer burned. + release_worker.set() + await asyncio.wait_for(asyncio.to_thread(committed.wait, 5), timeout=6) + assert committed.is_set(), ( + "watermark-fenced worker must keep its commit admission after " + "turn-hold expiry (fence was cancelled — attempt burned)" + ) + fake_db.archive_and_compact.assert_called_once() + # The commit went through the watermark-fenced path (concurrent tail + # rows above the watermark survive the compaction). + assert fake_db.archive_and_compact.call_args.kwargs.get("watermark") == 6 + + await _drain_deferred(runner) + await asyncio.wait_for(asyncio.to_thread(cleanup_done.wait, 5), timeout=6) + FencedStreamingAgent.last_instance.close.assert_called_once() + + # Successful adoption resets the hygiene failure streak and still never + # advances it (the deferral is not a failure). + assert not fake_db.increment_hygiene_failure_streak.called + assert fake_db.reset_hygiene_failure_streak.called + # Deferral notice still reaches the user. + sent = [m["content"] for m in adapter.sent] + assert any( + "deferred" in c.lower() or "still streaming" in c.lower() for c in sent + ), f"turn-hold must send deferral notice, got: {sent}" + + +@pytest.mark.asyncio +async def test_turn_hold_kept_admission_arms_flat_retry_only_when_nothing_commits( + monkeypatch, tmp_path +): + """If the kept-admission worker ends WITHOUT committing (summary failed + / attempt superseded), the flat non-escalating retry-after is restored so + sustained traffic does not spawn-and-abandon a compressor every turn — + but only AFTER the attempt truly ended, and without touching the streak. + """ + worker_started = threading.Event() + release_worker = threading.Event() + fake_db = MagicMock() + fake_db.get_compression_failure_cooldown.return_value = None + + class FencedNoCommitAgent: + def __init__(self, **kwargs): + self.session_id = kwargs.get("session_id", "sess-97963") + 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() + + def _compress_context( + self, messages, *_args, commit_fence=None, **_kwargs + ): + if commit_fence is not None: + commit_fence.mark_commit_watermark_fenced() + worker_started.set() + _spin_started = time.monotonic() + while not release_worker.is_set(): + if time.monotonic() - _spin_started > 20: + return (messages, None) + if commit_fence is not None: + commit_fence.touch_progress() + time.sleep(0.01) + # Summary failed — return unchanged, no commit. + return (messages, None) + + gateway_run = importlib.import_module("gateway.run") + _write_turnhold_config(tmp_path) + _install_fakes(monkeypatch, gateway_run, tmp_path, FencedNoCommitAgent) + + adapter = _CaptureAdapter() + runner = _build_runner(gateway_run, adapter, fake_db) + + result = await asyncio.wait_for(runner._handle_message(_make_event()), timeout=15) + assert result == "ok" + assert worker_started.is_set() + # While the attempt still runs: no cooldown, so preflight adoption + # stays possible. + assert not fake_db.record_compression_failure_cooldown.called + + release_worker.set() + await _drain_deferred(runner) + # Let the done-callback fire. + for _ in range(100): + if fake_db.record_compression_failure_cooldown.called: + break + await asyncio.sleep(0.05) + + # Nothing committed → flat retry-after restored (spacing), streak intact. + assert fake_db.record_compression_failure_cooldown.called, ( + "a kept-admission attempt that ends without committing must restore " + "the flat turn-hold retry-after spacing" + ) + args = fake_db.record_compression_failure_cooldown.call_args[0] + retry = args[1] - time.time() + assert retry <= 120, ( + f"retry-after must stay flat (~60s), got {retry:.0f}s" + ) + assert "turn-hold" in (args[2] or "") + assert not fake_db.increment_hygiene_failure_streak.called, ( + "turn-hold deferral must never advance the failure streak" + ) + + +@pytest.mark.asyncio +async def test_turn_hold_without_watermark_fence_still_cancels( + monkeypatch, tmp_path +): + """A worker whose commit is NOT watermark-fenced (no session_db / + watermark capture failed) must still be cancelled at turn-hold expiry — + a late unfenced commit could clobber newer turns. Never worse than the + status quo. (Complements the pinned #90845 test, which exercises the + same path through the public surface.) + """ + worker_started = threading.Event() + release_worker = threading.Event() + fake_db = MagicMock() + fake_db.get_compression_failure_cooldown.return_value = None + + class UnfencedStreamingAgent: + def __init__(self, **kwargs): + self.session_id = kwargs.get("session_id", "sess-97963") + 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() + + def _compress_context( + self, messages, *_args, commit_fence=None, **_kwargs + ): + # Deliberately NO mark_commit_watermark_fenced(). + worker_started.set() + _spin_started = time.monotonic() + while not release_worker.is_set(): + if time.monotonic() - _spin_started > 20: + return (messages, None) + if commit_fence is not None: + commit_fence.touch_progress() + time.sleep(0.01) + if commit_fence is not None and not commit_fence.begin_commit(): + return (messages, None) + try: + self._session_db.archive_and_compact( + self.session_id, + [{"role": "assistant", "content": "too late"}], + ) + return ([{"role": "assistant", "content": "too late"}], None) + finally: + if commit_fence is not None: + commit_fence.finish_commit() + + gateway_run = importlib.import_module("gateway.run") + _write_turnhold_config(tmp_path) + _install_fakes(monkeypatch, gateway_run, tmp_path, UnfencedStreamingAgent) + + adapter = _CaptureAdapter() + runner = _build_runner(gateway_run, adapter, fake_db) + + result = await asyncio.wait_for(runner._handle_message(_make_event()), timeout=15) + assert result == "ok" + assert worker_started.is_set() + + release_worker.set() + await _drain_deferred(runner) + await asyncio.sleep(0.2) + # The unfenced late commit was refused — discard as before the fix. + fake_db.archive_and_compact.assert_not_called() + # Legacy path still records the flat retry-after immediately. + assert fake_db.record_compression_failure_cooldown.called + assert not fake_db.increment_hygiene_failure_streak.called diff --git a/website/docs/user-guide/configuration.md b/website/docs/user-guide/configuration.md index d3e71e1a30..0c270e5660 100644 --- a/website/docs/user-guide/configuration.md +++ b/website/docs/user-guide/configuration.md @@ -909,7 +909,7 @@ Older configs with `compression.summary_model`, `compression.summary_provider`, `hygiene_total_ceiling_seconds` (default `600`) bounds the total wait even while tokens are still moving, so a degenerate trickle stream can't hold a turn hostage indefinitely. It is clamped to at least `hygiene_timeout_seconds`. -`hygiene_max_turn_hold_seconds` (default `10`) is the gateway's **turn-hold budget** — the maximum wall-clock the incoming message is held waiting on hygiene compression before the gateway stops waiting and proceeds on the uncompressed transcript. It exists because `hygiene_total_ceiling_seconds` alone can leave the wire silent for far longer than a chat transport's idle-timeout: a summary model that keeps streaming tokens keeps resetting the inactivity slice, so without a turn-hold budget the wait can stretch toward the ceiling while zero bytes reach the user — Telegram (and similar transports) then drop the connection and the turn appears frozen. Capping the turn's wait at this budget (well under the typical ~30s transport idle-timeout) guarantees the message is answered promptly; the compression worker keeps running detached and its commit is fenced (`CompressionCommitFence`), so when it eventually finishes it cannot overwrite the turns appended after the wait was abandoned. Raise it if your summary model routinely needs longer and your transport tolerates it; lower it for snappier recovery on very slow backends. +`hygiene_max_turn_hold_seconds` (default `10`) is the gateway's **turn-hold budget** — the maximum wall-clock the incoming message is held waiting on hygiene compression before the gateway stops waiting and proceeds on the uncompressed transcript. It exists because `hygiene_total_ceiling_seconds` alone can leave the wire silent for far longer than a chat transport's idle-timeout: a summary model that keeps streaming tokens keeps resetting the inactivity slice, so without a turn-hold budget the wait can stretch toward the ceiling while zero bytes reach the user — Telegram (and similar transports) then drop the connection and the turn appears frozen. Capping the turn's wait at this budget (well under the typical ~30s transport idle-timeout) guarantees the message is answered promptly. **The compression is not lost when the budget expires**: the worker keeps running detached and — when its commit is watermark-fenced (the normal case with a session DB) — it keeps its commit admission, so the finished summary is adopted at the next safe boundary and turns appended after the wait was abandoned survive verbatim as concurrent tail. This matters especially for **thinking/reasoning summary models** (DeepSeek, QwQ, etc.) whose reasoning phase alone can exceed the budget: their summaries land one turn late instead of never. If the commit cannot be safely fenced, the late result is discarded (`CompressionCommitFence`) and it cannot overwrite newer turns. Raise the budget if you'd rather have compression apply within the same turn and your transport tolerates the wait; lower it for snappier recovery on very slow backends. `hygiene_failure_cooldown_seconds` controls that per-session cooldown after a hygiene compression timeout or abort. During the cooldown, the gateway skips repeated hygiene attempts for the same oversized session so every incoming message does not block on the same broken auxiliary backend. `/compress`, `/reset`, or a healthy later turn can still recover the session.