fix(gateway): split turn-hold expiry from idle-timeout failure path
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.
This commit is contained in:
@@ -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:
|
||||
|
||||
+82
-1
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user