fix(gateway): bound the hygiene-compression turn-hold so a streaming summary cannot freeze the turn

Session hygiene auto-compression runs inline on the incoming-message path
and awaits the summary worker with a progress-aware inactivity budget
(hygiene_timeout_seconds) that extends up to hygiene_total_ceiling_seconds
(default 600s). A summary model that keeps streaming tokens keeps resetting
the inactivity slice, so the wait can stretch toward the ceiling while zero
bytes reach the user — chat transports (Telegram ~30s idle-timeout) drop the
connection and the turn appears frozen, even though the gateway is healthy.

Add hygiene_max_turn_hold_seconds (default 10), a turn-hold budget that caps
the wall-clock the incoming message waits on hygiene compression. The wait
slice is additionally capped at the remaining budget so the budget is
re-evaluated even when the worker keeps the inactivity slice large. On
exceeding the budget the gateway abandons the inline wait and proceeds on
the uncompressed transcript via the existing timeout path, which revokes the
worker's commit admission (CompressionCommitFence) and defers cleanup — so a
stale compression finishing later can never overwrite the turns appended
after the wait was abandoned.

Well under the typical transport idle-timeout, this guarantees the message
is answered promptly while the detached compression completes in the
background. Configurable via compression.hygiene_max_turn_hold_seconds.

Adds a regression test: a worker that streams progress continuously (so the
inactivity slice never fires) must be abandoned once it exceeds the
turn-hold budget, the turn proceeds uncompressed, and the stale commit is
fenced (no session mutation, role alternation intact).
This commit is contained in:
Machan-Army
2026-08-17 14:00:45 +00:00
committed by kshitij
parent 23bae43cfa
commit 543c85acbf
3 changed files with 228 additions and 0 deletions

View File

@@ -20165,6 +20165,14 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
_hyg_hard_msg_limit = 5000
_hyg_timeout_seconds = 30.0
_hyg_total_ceiling_seconds = 600.0
# Max wall-clock the user's TURN is held waiting on hygiene
# compression before the gateway stops waiting and proceeds on the
# uncompressed transcript (#TKT-0029). The compressor keeps running
# detached; its commit is fenced (revoke_commit_admission) so a
# stale result can never clobber turns appended after the wait was
# abandoned. Capped well below typical transport idle-timeouts
# (Telegram ~30s) so the wire never goes silent long enough to sever.
_hyg_max_turn_hold_seconds = 10.0
_hyg_failure_cooldown_seconds = 300.0
_hyg_config_context_length = None
_hyg_provider = None
@@ -20232,6 +20240,14 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
_hyg_total_ceiling_seconds = max(
_hyg_total_ceiling_seconds, _hyg_timeout_seconds,
)
_raw_turn_hold = _comp_cfg.get("hygiene_max_turn_hold_seconds")
if _raw_turn_hold is not None:
try:
_parsed = float(_raw_turn_hold)
if _parsed > 0:
_hyg_max_turn_hold_seconds = _parsed
except (TypeError, ValueError):
pass
_raw_cooldown = _comp_cfg.get("hygiene_failure_cooldown_seconds")
if _raw_cooldown is not None:
try:
@@ -20523,6 +20539,31 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
- _hyg_commit_fence.seconds_since_progress(),
0.005,
)
# Bounded turn-hold (#TKT-0029): cap
# this slice at the remaining
# turn-hold budget so the wait is
# re-evaluated against
# _hyg_max_turn_hold_seconds at
# least that often — otherwise a
# continuously-streaming worker
# (which keeps the inactivity slice
# large) would hold the turn until
# the total ceiling before the
# budget check ever runs.
_turn_hold_remaining = (
_hyg_max_turn_hold_seconds
- (time.monotonic() - _hyg_wait_started)
)
if _turn_hold_remaining <= 0:
# Budget already exhausted —
# force an immediate timeout so
# the abandonment path below runs.
_slice = 0.005
else:
_slice = min(
_slice,
max(_turn_hold_remaining, 0.005),
)
try:
_compressed, _ = await asyncio.wait_for(
asyncio.shield(_hyg_future),
@@ -20532,6 +20573,35 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
except asyncio.TimeoutError:
_hyg_waited = time.monotonic() - _hyg_wait_started
_idle = _hyg_commit_fence.seconds_since_progress()
# Bounded turn-hold (#TKT-0029):
# never hold the user's TURN
# longer than
# _hyg_max_turn_hold_seconds,
# even if the summary model is
# still streaming. Past the
# budget we stop waiting and
# fall through to the timeout
# path below, which revokes
# commit admission and proceeds
# on the uncompressed
# transcript — the wire never
# stays silent long enough to
# trip a transport idle-timeout.
if (
_hyg_waited
>= _hyg_max_turn_hold_seconds
):
logger.info(
"Session hygiene compression for "
"session %s exceeded the turn-hold "
"budget (%.1fs >= %.1fs) — "
"abandoning inline wait, proceeding "
"without compression this turn",
session_entry.session_id,
_hyg_waited,
_hyg_max_turn_hold_seconds,
)
raise
if (
_idle < _hyg_timeout_seconds
and _hyg_waited < _hyg_total_ceiling_seconds

View File

@@ -633,6 +633,161 @@ async def test_session_hygiene_timeout_continues_to_agent_and_sets_cooldown(monk
SlowCompressAgent.last_instance.close.assert_called_once()
@pytest.mark.asyncio
async def test_session_hygiene_turn_hold_budget_abandons_streaming_wait(
monkeypatch, tmp_path
):
"""A compression that still streams progress must not hold the turn hostage.
Regression test for the bounded turn-hold (#TKT-0029). The worker keeps
ticking the commit fence (touch_progress), so the per-slice inactivity
timeout NEVER fires — without a turn-hold budget the gateway would extend
the wait up to the total ceiling (default 600s) while zero bytes hit the
wire, severing the transport. The turn must instead be abandoned once it
exceeds ``hygiene_max_turn_hold_seconds``, proceed on the uncompressed
transcript, and fence the stale commit.
"""
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 StreamingCompressAgent:
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()
# Stream progress continuously so the inactivity slice never
# times out; only the turn-hold budget can abandon this wait.
while not release_worker.is_set():
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"}],
)
self._last_compaction_in_place = True
return ([{"role": "assistant", "content": "too late"}], None)
finally:
if commit_fence is not None:
commit_fence.finish_commit()
fake_run_agent = types.ModuleType("run_agent")
fake_run_agent.AIAgent = StreamingCompressAgent
monkeypatch.setitem(sys.modules, "run_agent", fake_run_agent)
cfg_path = tmp_path / "config.yaml"
cfg_path.write_text(
"compression:\n"
" enabled: true\n"
# Inactivity budget is huge, so the slice timeout can never fire on
# its own; the turn-hold budget is the ONLY thing that abandons.
" 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"
)
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-turnhold",
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 well under the 600s
# ceiling — the turn-hold budget (~0.3s) abandoned the streaming wait.
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
# The stale commit must be fenced: the late worker never mutates the session.
fake_db.archive_and_compact.assert_not_called()
release_worker.set()
await asyncio.wait_for(asyncio.to_thread(cleanup_done.wait), timeout=3)
fake_db.archive_and_compact.assert_not_called()
StreamingCompressAgent.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

View File

@@ -881,6 +881,7 @@ compression:
hygiene_hard_message_limit: 5000 # Gateway safety valve — see below
hygiene_timeout_seconds: 30 # Max seconds of NO summary-model output before hygiene compression is cut off
hygiene_total_ceiling_seconds: 600 # Absolute cap on the hygiene wait even while tokens are still streaming
hygiene_max_turn_hold_seconds: 10 # Max wall-clock the incoming turn waits on hygiene compression before proceeding uncompressed — see below
hygiene_failure_cooldown_seconds: 300 # First rung of the per-session hygiene-failure backoff (x1/x3/x9, capped at 1h)
context_timeout_seconds: 120 # Inactivity budget for in-agent compress_context (loop /compress / preflight) — see below
context_total_ceiling_seconds: 600 # Absolute cap on the *pre-commit* in-agent compress_context wait even while tokens are still streaming (an already-started SessionDB commit is never abandoned; overruns are logged + surfaced)
@@ -908,6 +909,8 @@ 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_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.
The value is the **first rung** of an escalating ladder, not a fixed interval: consecutive failures for the same session wait `1x`, `3x`, then `9x` this value, capped at one hour. A session whose summary model is permanently broken therefore backs off instead of retrying forever on a fixed interval, and a run that actually shrinks the transcript resets it to the first rung. Escalation is per-session and process-local — a gateway restart resets it to the first rung while the cooldown deadline itself survives.