diff --git a/gateway/run_turn.py b/gateway/run_turn.py index 2b93610155..4998cbef2c 100644 --- a/gateway/run_turn.py +++ b/gateway/run_turn.py @@ -79,8 +79,7 @@ class GatewayTurnMixin: skey or "", model, override_model, override_runtime.get("provider"), ) return override_model, override_runtime - # No api_key on the override: env-based resolution below, then the override's - # model/provider applied on top. + # No api_key on the override: env-based resolution below, override model/provider on top. logger.debug( "Session model override (no api_key, fallback): session=%s config_model=%s override_model=%s", skey or "", model, override_model, @@ -121,8 +120,8 @@ class GatewayTurnMixin: if override and skey: model, runtime_kwargs = self._apply_session_model_override(skey, model, runtime_kwargs) - # No model.default but a provider resolved (e.g. `hermes auth add openai-codex` without - # `hermes model`): fall back to the provider's first catalog model so the API call has one. + # Provider resolved but no model.default (`hermes auth add` without `hermes model`): use the + # provider's first catalog model. if not model and runtime_kwargs.get("provider"): with suppress(Exception): from hermes_cli.models import get_default_model_for_provider @@ -132,9 +131,8 @@ class GatewayTurnMixin: "No model configured — defaulting to %s for provider %s", model, runtime_kwargs["provider"], ) - # Final safety net: an empty model (e.g. transient config-cache miss on a post-interrupt - # recovery turn) makes every API call fail HTTP 400 and the session goes silent — reuse the - # last model resolved for this session, else the most recent process-wide. + # Final safety net: an empty model (transient config-cache miss) makes every API call 400 and + # the session goes silent — reuse the last model resolved for this session, else process-wide. if not model: _lr_state = self._peek_session_state(skey) if skey else None _lr_star = self._peek_session_state("*") @@ -159,8 +157,7 @@ class GatewayTurnMixin: def _resolve_turn_agent_config(self, user_message: str, model: str, runtime_kwargs: dict) -> dict: """Effective model/runtime config for one turn. With `/fast` priority on, fast-mode - ``request_overrides`` are deep-merged OVER the per-provider ones so both reach the model. - """ + ``request_overrides`` are deep-merged OVER the per-provider ones so both reach the model.""" from gateway.run import _deep_merge_request_overrides from hermes_cli.models import resolve_fast_mode_overrides @@ -183,8 +180,7 @@ class GatewayTurnMixin: ), } if getattr(self, "_service_tier", None) != "priority": - # None (normal) or auto/cold — the bounded window is applied per request by - # agent.fast_mode, not pinned into request_overrides. + # None / auto / cold: the bounded window is applied per request by agent.fast_mode. route["request_overrides"] = base_request_overrides return route try: @@ -254,7 +250,7 @@ class GatewayTurnMixin: topic-binding heal). Returns ``(source, session_entry, session_key)`` or ``None`` to drop the event.""" # Topic-mode DMs: rewrite a stale/foreign thread_id to the user's last-active topic so a - # cross-topic Reply or stripped plain reply doesn't fragment the conversation. + # cross-topic Reply doesn't fragment the conversation. recovered = await asyncio.to_thread(self._recover_telegram_topic_thread_id, source) if recovered is not None: logger.info( @@ -287,9 +283,8 @@ class GatewayTurnMixin: ) return else: - # Internal wakes must observe reset policy without counting as user activity, or - # periodic Kanban/process notifications keep the routing key alive across every - # daily/idle boundary. + # Internal wakes observe reset policy without counting as user activity, or periodic + # notifications keep the routing key alive across every daily/idle boundary. session_entry = await self.async_session_store.get_or_create_session( source, touch_activity=not bool(getattr(event, "internal", False)), ) @@ -334,9 +329,8 @@ class GatewayTurnMixin: if canonical_session_id and canonical_session_id != bound_session_id: bound_session_id = canonical_session_id if bound_session_id and bound_session_id != session_entry.session_id: - # Route through SessionStore so the session_key → session_id mapping is persisted and - # the previous lane session ended cleanly (mutating session_entry in place split-brained - # the JSON index from downstream code). + # Route through SessionStore so the key → id mapping persists and the previous lane + # session ends cleanly (in-place mutation split-brained the JSON index). switched = await self.async_session_store.switch_session(session_key, bound_session_id) if switched is not None: session_entry = switched @@ -350,13 +344,11 @@ class GatewayTurnMixin: async def _hmwa_open_session(self, session_entry, session_key, source): """Consume auto-reset / fresh-reset flags and emit ``session:start`` for new sessions. Returns ``(_was_auto_reset, _is_new_session)``.""" - # Consume was_auto_reset immediately so it cannot re-fire on later messages and wipe - # model/reasoning overrides set between turns. + # Consume was_auto_reset immediately so it cannot re-fire and wipe overrides set between turns. _was_auto_reset = getattr(session_entry, "was_auto_reset", False) if _was_auto_reset: - # Full conversation boundary: one funnel call clears every conversation-scoped - # per-session dict; evict the cached agent (keyed on the stable session_key) so - # context_compressor._previous_summary cannot leak prior history into new summaries. + # Conversation boundary: the funnel clears every conversation-scoped dict; evict the cached + # agent so context_compressor._previous_summary cannot leak into new summaries. self._clear_conversation_scope(session_key, reason="auto_reset") self._evict_cached_agent(session_key) session_entry.was_auto_reset = False @@ -380,8 +372,7 @@ class GatewayTurnMixin: from gateway.run import _AUTO_RESET_CONTEXT_NOTES, _auto_reset_reason_text reset_reason = getattr(session_entry, 'auto_reset_reason', None) or 'idle' context_note = _AUTO_RESET_CONTEXT_NOTES.get(reset_reason, _AUTO_RESET_CONTEXT_NOTES["idle"]) - # Long-lived Slack/Discord channels: point the agent at the specific prior same-channel - # session so session_search recalls that context (deterministic, no extra API/DB calls). + # Long-lived channels: point the agent at the prior same-channel session for session_search. try: continuity_note = build_channel_continuity_note(session_entry, source) except Exception: @@ -464,8 +455,7 @@ class GatewayTurnMixin: timeout=_float_env("HERMES_TURN_LEASE_TIMEOUT", DEFAULT_LEASE_WAIT), ) except TurnLeaseTimeoutError: - # The broad session-context cleanup finally starts later; restore the tokens here or - # this early exit leaks task-local identity. + # The cleanup finally starts later; restore the tokens here or this exit leaks identity. self._clear_session_env(_session_env_tokens) raise if _lease_token is not None: @@ -498,8 +488,7 @@ class GatewayTurnMixin: hs.provider = _model_cfg.get("provider") or None hs.base_url = _model_cfg.get("base_url") or None - # Only the enabled flag is shared with the agent's compression config; hygiene's - # threshold is deliberately separate (runs higher). + # Only the enabled flag is shared with the agent's compression config (hygiene runs higher). _comp_cfg = data.get("compression", {}) if not isinstance(_comp_cfg, dict): return @@ -565,8 +554,7 @@ class GatewayTurnMixin: except Exception: hs.config_context_length = None - # custom_providers per-model context_length fallback (as in run_agent.py); must run - # after runtime resolution so base_url is set. + # custom_providers per-model context_length fallback (as in run_agent.py); needs base_url. if hs.config_context_length is None and hs.base_url: try: try: @@ -610,9 +598,8 @@ class GatewayTurnMixin: else: _approx_tokens, _token_source = estimate_messages_tokens_rough(history), "estimated" - # Hard safety valve: force compression at an extreme message count regardless of token - # estimates, breaking the spiral where API disconnects prevent token data → no compression - # → more disconnects. Default 5000 sits clear of legitimate 1M+ context sessions. + # Hard safety valve: force compression at an extreme message count regardless of tokens, + # breaking the disconnect → no token data → no compression spiral. 5000 clears 1M+ sessions. _needs_compress = _approx_tokens >= _compress_token_threshold or _msg_count >= hs.hard_msg_limit if _needs_compress: @@ -636,9 +623,8 @@ class GatewayTurnMixin: _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). Another attempt would wait up to the 600s - # ceiling behind a commit the fence will refuse, while inbound messages demote to queue. + # A prior compression still holds the durable lock (e.g. a shielded worker left by /stop): + # another attempt would wait up to 600s behind a commit the fence will refuse. logger.info( "Session hygiene: skipping compression for %s; " "another compression is already in flight", session_entry.session_id, @@ -665,16 +651,14 @@ class GatewayTurnMixin: while True: if fence.is_cancelled: raise asyncio.TimeoutError - # Charge the idle budget from the LAST PROGRESS event, not from the start of this wait - # slice — otherwise silence can approach 2x the configured timeout. + # Charge the idle budget from the LAST PROGRESS event, else silence can approach 2x timeout. _hyg_waited = time.monotonic() - attempt.wait_started _slice = min( max(hs.timeout_seconds - fence.seconds_since_progress(), 0.005), max(hs.total_ceiling_seconds - _hyg_waited, 0.005), ) - # Cap the slice at the remaining turn-hold budget so it is re-checked at least that - # often — a continuously-streaming worker would otherwise keep the slice large and hold - # the turn until the ceiling. Budget exhausted → immediate timeout → abandonment path. + # Cap the slice at the remaining turn-hold budget so a continuously-streaming worker can't + # hold the turn until the ceiling. Budget exhausted → immediate timeout → abandonment. _turn_hold_remaining = hs.max_turn_hold_seconds - (time.monotonic() - attempt.wait_started) _slice = 0.005 if _turn_hold_remaining <= 0 else min(_slice, max(_turn_hold_remaining, 0.005)) # Short poll so a /stop or /restart cancel is not stuck behind a full idle window. @@ -688,9 +672,8 @@ class GatewayTurnMixin: raise _hyg_waited = time.monotonic() - attempt.wait_started _idle = fence.seconds_since_progress() - # Never hold the user's TURN past the budget even if the summary model is still - # streaming: the turn-hold path proceeds on the uncompressed transcript so the wire - # never trips a transport idle-timeout. + # Never hold the TURN past the budget even while the summary streams: proceed on the + # uncompressed transcript so the wire never trips a transport idle-timeout. if _hyg_waited >= hs.max_turn_hold_seconds: logger.info( "Session hygiene compression for session %s exceeded the turn-hold " @@ -813,9 +796,8 @@ class GatewayTurnMixin: 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). Flat spacing so sustained traffic does not spawn and - # abandon a fresh compressor every turn. + # Nothing to adopt (summary failed / fence refused / superseded): flat spacing so + # sustained traffic doesn't spawn and abandon a compressor every turn. _record_hygiene_cooldown( _gw, _sid, _HYGIENE_TURNHOLD_RETRY_SECONDS, "hygiene compression deferred: turn-hold budget expired and the " @@ -831,8 +813,7 @@ class GatewayTurnMixin: _adopted = await self._hmwa_hygiene_cancel_or_adopt(attempt, "session hygiene turn-hold") if _adopted is not None: return _adopted - # Short flat retry-after: without it every turn re-spawns a compressor, holds it for - # the budget and cancels it — token burn that never commits. + # Short flat retry-after, else every turn re-spawns, holds and cancels a compressor. _record_hygiene_cooldown( self, session_entry.session_id, _HYGIENE_TURNHOLD_RETRY_SECONDS, "hygiene compression deferred: turn-hold budget expired while the " @@ -862,12 +843,10 @@ class GatewayTurnMixin: _hyg_waited = time.monotonic() - attempt.wait_started _hyg_total_exhausted = _hyg_waited >= hs.total_ceiling_seconds or fence.deadline_exceeded if _hyg_total_exhausted: - # The worker cooperatively checks this deadline between digest calls. Keep its lease - # until it exits so an unchanged session cannot overlap a retry (the release below is - # then a no-op). + # The worker checks this deadline between digest calls; keep its lease until it exits so + # an unchanged session cannot overlap a retry (the release below is then a no-op). 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. + # Capture fence state BEFORE try_cancel (which itself sets is_cancelled). _hyg_fence_cancelled = fence.is_cancelled _adopted = await self._hmwa_hygiene_cancel_or_adopt(attempt, "session hygiene timeout") if _adopted is not None: @@ -955,13 +934,11 @@ class GatewayTurnMixin: a repoint-then-failed-rewrite would point the live entry at an empty session.""" from agent.model_metadata import estimate_messages_tokens_rough _hyg_agent = attempt.agent - # _compress_context ends the old session and creates a new session_id; compressed messages - # go into the NEW session so the old transcript stays intact and searchable. + # _compress_context rotates to a NEW session_id so the old transcript stays intact/searchable. _hyg_new_sid = _hyg_agent.session_id _hyg_rotated = _hyg_new_sid != session_entry.session_id _hyg_in_place = bool(getattr(_hyg_agent, "_last_compaction_in_place", False)) - # Anti-growth guard: refuse a compression that did not shrink the transcript (observed: - # 427K -> 598K). Compare like-for-like rough estimates. + # Anti-growth guard: refuse a compression that did not shrink the transcript (seen 427K→598K). _hyg_in_toks = estimate_messages_tokens_rough(history) _hyg_out_toks = estimate_messages_tokens_rough(_compressed) if _hyg_rotated and _hyg_out_toks > _hyg_in_toks: @@ -984,8 +961,7 @@ class GatewayTurnMixin: _hyg_in_place = False else: session_entry.session_id = _hyg_new_sid - # The held turn lease follows the rotation so an alias key resolving the fresh - # child still serializes against this turn. + # The held turn lease follows the rotation (alias keys still serialize on this turn). self._rebind_turn_lease(_quick_key, run_generation, _hyg_new_sid) await self.async_session_store._save() await asyncio.to_thread( @@ -993,8 +969,7 @@ class GatewayTurnMixin: ) if _hyg_rotated or _hyg_in_place: - # Transcript rewritten (rotation) or already persisted by archive_and_compact() - # (in-place): reset the stored token count to match the new active set. + # Rewritten (rotation) or persisted by archive_and_compact() (in-place): reset token count. session_entry.last_prompt_tokens = 0 attempt.history = _compressed _new_count = len(_compressed) @@ -1028,21 +1003,18 @@ class GatewayTurnMixin: attempt, _compressed, history, plan, session_entry=session_entry, source=source, _quick_key=_quick_key, run_generation=run_generation, ) - # Summary failure aborts the compressor entirely (messages unchanged, nothing dropped). - # Warn the gateway user visibly — agent.log is invisible on TG/Discord/etc. — so they - # know the chat is "frozen" at this size and can /compress to retry or /reset. + # Summary failure aborts the compressor (nothing dropped). Warn the user visibly — agent.log + # is invisible on TG/Discord — so they know the chat is "frozen" and can /compress or /reset. _comp = getattr(attempt.agent, "context_compressor", None) _hyg_aborted = _comp is not None and getattr(_comp, "_last_compress_aborted", False) # A fence-cancelled _compress_context returns the original transcript with - # _last_compress_aborted still False. Treat that no-op as an abort so hygiene records a - # cooldown instead of retrying into the 600s wait. A successful rotate/in-place commit is - # not an abort even if a later invalidation flipped the fence. + # _last_compress_aborted False: treat that no-op as an abort so hygiene records a cooldown + # instead of retrying into the 600s wait. A committed rotate/in-place is never an abort. _hyg_fence_cancelled = bool(attempt.commit_fence.is_cancelled and not outcome.rotated and not outcome.in_place) if _hyg_fence_cancelled: _hyg_aborted = True - # Recovery decision lives in the unit-tested predicate: the degenerate "neither rotated - # nor compacted in place" path reuses the pre-compression counts, so a numbers-only check - # would read a no-op as success and clear the streak. + # Recovery decision lives in the unit-tested predicate: the "neither rotated nor in place" + # path reuses pre-compression counts, so a numbers-only check would read a no-op as success. if not _hyg_aborted and hygiene_compaction_recovered( aborted=_hyg_aborted, rotated=outcome.rotated, in_place=outcome.in_place, msg_count=plan.msg_count, new_count=outcome.new_count, approx_tokens=plan.approx_tokens, @@ -1061,8 +1033,7 @@ class GatewayTurnMixin: ) if not _hyg_fence_cancelled: _err = getattr(_comp, "_last_summary_error", None) or "unknown error" - # Force-redact: provider exception text may contain credentials and this message - # reaches gateway users. + # Force-redact: provider exception text may contain credentials; this reaches users. from agent.redact import redact_sensitive_text _err = redact_sensitive_text(_err, force=True) await self._hmwa_hygiene_notify( @@ -1072,8 +1043,7 @@ class GatewayTurnMixin: "session, or check your auxiliary.compression model configuration.", "compression-failure warning", ) - # If the CONFIGURED aux model failed and we recovered on the main model, tell the user — - # a misconfigured auxiliary.compression.model is something only they can fix. + # Configured aux model failed, recovered on the main model: only the user can fix that config. elif _comp is not None and getattr(_comp, "_last_aux_model_failure_model", None): _aux_model = getattr(_comp, "_last_aux_model_failure_model", "") _aux_err = getattr(_comp, "_last_aux_model_failure_error", None) or "unknown error" @@ -1120,8 +1090,7 @@ class GatewayTurnMixin: "configured providers.", session_entry.session_id, exc, exc_info=True, ) _hyg_session_db = getattr(self._session_db, "_db", self._session_db) - # Hygiene is the same lossy rewrite as normal compression: with - # compression.checkpoint_required on, load the memory provider so the checkpoint exists + # With compression.checkpoint_required on, load the memory provider so the checkpoint exists # before any mutation; otherwise keep the fast path (no provider init). from hermes_cli.config import load_config as _load_cfg from utils import is_truthy_value as _is_truthy @@ -1135,8 +1104,7 @@ class GatewayTurnMixin: session_id=session_entry.session_id, session_db=_hyg_session_db, ) _seed_hygiene_system_prompt(_hyg_agent, _hyg_session_row) - # If compression must rebuild instead of retaining the cached prompt, make the persisted - # result deliberately stale for every real gateway surface. + # A rebuilt (not retained) prompt is deliberately stale for every real gateway surface. _hyg_agent.platform = _GATEWAY_HYGIENE_PLATFORM return _hyg_agent, _hyg_session_db @@ -1190,8 +1158,7 @@ class GatewayTurnMixin: run_generation=run_generation, ) finally: - # Evict the cached agent so the next turn rebuilds its system prompt from current - # SOUL.md, memory, and skills. + # Evict the cached agent so the next turn rebuilds its system prompt. self._evict_cached_agent(session_key) if not attempt.cleanup_deferred: await self._cleanup_agent_resources_off_loop(_hyg_agent, context="session hygiene") @@ -1222,9 +1189,8 @@ class GatewayTurnMixin: if str(_hyg_runtime.get("api_mode") or "").lower() == "codex_app_server": await self._hmwa_hygiene_codex_compaction(hs, plan, history, session_entry, session_key, _hyg_runtime) elif _hyg_runtime.get("api_key"): - # Pass the FULL transcript (tool results included), matching the agent loop: - # filtering to user/assistant starved the compressor — tool results are the bulk - # of context and short histories tripped the protect-first/last early-return. + # Pass the FULL transcript (tool results included) as the agent loop does: filtering + # to user/assistant starved the compressor (tool results are the bulk of context). _hyg_msgs = [m for m in history if m.get("role") in {"user", "assistant", "tool"}] if len(_hyg_msgs) >= 4: await self._hmwa_hygiene_detached_attempt( @@ -1232,8 +1198,7 @@ class GatewayTurnMixin: source, session_entry, session_key, _quick_key, run_generation, ) except HygieneTurnHoldExceeded: - # Availability boundary, not a failure — already logged at INFO by the turn-hold - # handler; the generic warning below made thinking-model deployments read as broken. + # Availability boundary, not a failure — already logged at INFO by the turn-hold handler. pass except Exception as e: logger.warning("Session hygiene auto-compress failed: %s", e) @@ -1252,9 +1217,8 @@ class GatewayTurnMixin: "Briefly introduce yourself and mention that /help shows available commands. " "Keep the introduction concise -- one or two sentences max.]" ) - # Opt-in profile-build path: when onboarding.profile_build is "ask" (default) and not - # yet offered on this install, swap the plain intro for a consent-gated directive that - # offers to build a user profile via memory(target="user"). Fires at most once. + # onboarding.profile_build == "ask" (default) and not yet offered: swap the plain intro for + # a consent-gated profile-build directive. Fires at most once. try: from agent.onboarding import ( PROFILE_BUILD_FLAG, is_seen, mark_seen, profile_build_directive, @@ -1270,14 +1234,12 @@ class GatewayTurnMixin: logger.debug("Profile-build onboarding directive failed, using plain intro: %s", _pb_err) turn_sidecar_notes.append(_intro_note) - # One-time prompt if no home channel is set for this platform. Skipped for webhooks — - # they deliver directly to configured targets (github_comment, etc.). + # One-time prompt if no home channel is set (webhooks deliver to configured targets instead). if not source.platform or source.platform in (Platform.LOCAL, Platform.WEBHOOK): return platform_name = source.platform.value env_key = _home_target_env_var(platform_name) - # Multiplex: the home channel may live only in the profile secret scope / PlatformConfig, - # not process os.environ. + # Multiplex: the home channel may live only in the profile secret scope, not os.environ. home_env = "" if env_key: with suppress(Exception): @@ -1288,8 +1250,7 @@ class GatewayTurnMixin: with suppress(Exception): if not home_env and self.config.get_home_channel(source.platform): home_env = "set" - # Secondary-profile platforms (e.g. Slack on yolo) may only exist under that profile's - # loaded config — re-read live config inside the already-installed scope. + # Secondary-profile platforms may only exist under that profile's config — re-read in scope. if not home_env: with suppress(Exception): from gateway.config import load_gateway_config as _lgc @@ -1297,8 +1258,7 @@ class GatewayTurnMixin: if prof and prof != "default" and _lgc().get_home_channel(source.platform): home_env = "set" if not home_env: - # Slack routes every Hermes command through the single parent slash command - # `/hermes`; bare `/sethome` is not registered and would fail. + # Slack routes every command through the parent `/hermes`; bare `/sethome` would fail. sethome_cmd = "/hermes sethome" if source.platform == Platform.SLACK else "/sethome" await self._deliver_platform_notice( source, f"📬 No home channel is set for {platform_name.title()}. " @@ -1330,8 +1290,7 @@ class GatewayTurnMixin: if _message_timestamps_enabled(_load_gateway_config()): message_text = _render_msg_ts(_clean_message_text, persist_user_timestamp, tz=_evt_tz) else: - # Toggle off: the model sees the clean message; the timestamp is still stored - # as metadata for later opt-in. + # Toggle off: the model sees the clean message; timestamp stored for later opt-in. message_text = _clean_message_text except Exception as _ts_err: logger.debug("Message timestamp injection failed (non-fatal): %s", _ts_err) @@ -1363,9 +1322,8 @@ class GatewayTurnMixin: _sanitize_gateway_final_response, _should_clear_resume_pending_after_turn, ) response = agent_result.get("final_response") or "" - # Hidden-reasoning-only retry exhaustion: the loop's sentinel text doubles as - # final_response and would be delivered verbatim — where peer agents can ingest it as a - # completed assistant turn. + # Hidden-reasoning-only retry exhaustion: the loop's sentinel text doubles as final_response + # and would be delivered verbatim (peer agents would ingest it as a completed turn). if _is_gateway_hidden_reasoning_incomplete_turn(agent_result): response = "" _intentional_silence = self._is_intentional_silence(agent_result, response) @@ -1383,9 +1341,8 @@ class GatewayTurnMixin: time.time() - _msg_start_time, agent_result.get("api_calls", 0), len(response), ) - # Successful turn: clear the stuck-loop counter (accumulates only across CONSECUTIVE - # restarts where the session never completed) and resume_pending (set by drain-timeout - # shutdown) so later messages don't get the restart-interruption system note. + # Successful turn: clear the consecutive-restart stuck-loop counter and resume_pending (set + # by drain-timeout shutdown) so later messages don't get the restart-interruption note. if session_key and _should_clear_resume_pending_after_turn(agent_result): await self._clear_restart_failure_count(session_key) try: @@ -1398,14 +1355,12 @@ class GatewayTurnMixin: response = _normalize_empty_agent_response(agent_result, response, history_len=len(history)) response = _sanitize_gateway_final_response(source.platform, response) - # Ordering contract: the agent thread already updated the contextvar in - # conversation_compression.py; propagate to SessionEntry + _save() — but only if the - # binding still points at the session this run was launched against. + # The agent thread already updated the contextvar; propagate to SessionEntry + _save() only + # if the binding still points at the session this run was launched against. if agent_result.get("session_id") and agent_result["session_id"] != session_entry.session_id: if session_entry.session_id == _run_start_session_id: session_entry.session_id = agent_result["session_id"] - # The held turn lease follows the rotation: transcript persistence writes to the - # NEW id, so the serialization boundary must move with it. + # The held turn lease follows the rotation (persistence writes to the NEW id). self._rebind_turn_lease(_quick_key, run_generation, session_entry.session_id) await self.async_session_store._save() await self.async_session_store._record_gateway_session_peer( @@ -1451,8 +1406,7 @@ class GatewayTurnMixin: display_reasoning = "\n".join(lines[:15]) + f"\n_... ({len(lines) - 15} more lines)_" else: display_reasoning = last_reasoning.strip() - # Render style is per-platform: Discord defaults to "-# " subtext (native small grey - # metadata text); other platforms keep the fenced code block. + # Per-platform render style: Discord defaults to "-# " subtext, others keep the code block. try: from gateway.display_config import resolve_display_setting _reasoning_style = resolve_display_setting( @@ -1496,8 +1450,7 @@ class GatewayTurnMixin: # Pending process watchers (check_interval on background processes) try: from tools.process_registry import process_registry - # Detach the current batch atomically: reassign to a fresh list so a watcher appended - # by a concurrent session during the yield isn't dropped by clear(). + # Detach the batch atomically (reassign, not clear()) so concurrent appends aren't dropped. watchers = process_registry.pending_watchers process_registry.pending_watchers = [] for i, watcher in enumerate(watchers): @@ -1507,9 +1460,8 @@ class GatewayTurnMixin: except Exception as e: logger.error("Process watcher setup error: %s", e) - # Drain watch notifications that arrived during the run. The queue also carries process - # completions (per-process watcher task above) and async-delegation completions (owned by - # _async_delegation_watcher) — inject only watch-type events, leave the rest queued. + # Drain watch notifications that arrived during the run; the queue also carries process / + # async-delegation completions owned elsewhere — inject only watch-type events. try: from tools.process_registry import process_registry as _pr await self._drain_watch_notifications(_pr.completion_queue) @@ -1534,8 +1486,7 @@ class GatewayTurnMixin: agent_failed_early = bool(agent_result.get("failed")) hidden_reasoning_incomplete = _is_gateway_hidden_reasoning_incomplete_turn(agent_result) _err = str(agent_result.get("error", "")).lower() - # Specific multi-word phrases (not bare "exceed"/"token") avoid false positives on - # transient errors such as "rate limit exceeded"; matches run_agent.py's classifier. + # Multi-word phrases (not bare "exceed"/"token") avoid matching "rate limit exceeded". is_context_overflow_failure = agent_failed_early and ( bool(agent_result.get("compression_exhausted")) or any(p in _err for p in self._CONTEXT_OVERFLOW_ERROR_PHRASES) @@ -1576,8 +1527,7 @@ class GatewayTurnMixin: logger.info("Auto-resetting session %s after compression exhaustion.", session_entry.session_id) new_entry = await self.async_session_store.reset_session(session_key) self._evict_cached_agent(session_key) - # Conversation boundary: one funnel call clears every conversation-scoped per-session - # dict (see _CONVERSATION_SCOPED_STATE). + # Conversation boundary: the funnel clears every conversation-scoped per-session dict. self._clear_conversation_scope(session_key, reason="compression_exhausted_reset") if new_entry is not None: # Re-point the Telegram topic binding at the fresh session, or the binding-heal walk @@ -1622,18 +1572,16 @@ class GatewayTurnMixin: store = self.async_session_store sid = session_entry.session_id history = prepared.history - # The agent already persisted this turn's rows via _flush_messages_to_session_db() (the - # codex app-server runtime too: it flushes its own projection and reports - # agent_persisted=True); skip the DB write to avoid duplicates. Default = a session DB - # exists; a non-persisting runtime opts in via False. + # The agent already persisted this turn's rows (codex app-server reports agent_persisted=True + # too); skip the DB write. Default = a session DB exists; non-persisting runtimes pass False. agent_persisted = agent_result.get("agent_persisted", self._session_db is not None) if is_context_overflow_failure: pass # Skip all transcript writes — don't grow a broken session else: if not history: - # Fresh session: write the full tool definitions as the first entry so the transcript - # is self-describing — the same dicts sent as tools=[...] in the API request. + # Fresh session: the tool definitions (as sent in the API request) make the transcript + # self-describing. await store.append_to_transcript(sid, { "role": "session_meta", "tools": agent_result.get("tools", []) or [], @@ -1670,9 +1618,8 @@ class GatewayTurnMixin: skip_db=agent_persisted, ) else: - # Attach the inbound platform message_id to the first user entry written this turn - # so platform-level quote-resolution (e.g. Yuanbao QuoteContextMiddleware's - # transcript fallback) can find earlier @bot messages by their original id. + # Attach the inbound platform message_id to the first user entry so platform-level + # quote-resolution (e.g. Yuanbao) can find earlier @bot messages by original id. _user_msg_id_attached = False for msg in new_messages: if msg.get("role") == "system": @@ -1688,8 +1635,7 @@ class GatewayTurnMixin: _user_msg_id_attached = True await store.append_to_transcript(sid, entry, skip_db=agent_persisted) - # The agent persists token counts and model itself; keep only last_prompt_tokens here for - # context-window tracking and compression decisions. + # The agent persists token counts/model itself; keep only last_prompt_tokens for hygiene. await store.update_session( session_entry.session_key, last_prompt_tokens=agent_result.get("last_prompt_tokens", 0), touch_activity=not bool(getattr(event, "internal", False)), @@ -1706,15 +1652,13 @@ class GatewayTurnMixin: ): """Final delivery decisions: intentional silence, voice reply, streamed-turn media/footer. Returns the text for the adapter to send, or ``None`` when already delivered.""" - # Intentional silence is a delivery decision, not a transcript mutation: the [SILENT] - # assistant turn stays persisted so later turns keep user/assistant alternation. + # Intentional silence is a delivery decision: the [SILENT] turn stays persisted (alternation). if _intentional_silence: logger.info("Suppressing intentional silence marker for session %s", session_entry.session_id) response = "" adapter = self._adapter_for_source(source) - # Auto voice reply (TTS audio before the text) unless streaming TTS already delivered - # audio for this turn. + # Auto voice reply (TTS audio before the text) unless streaming TTS already delivered audio. _streaming_tts_done = adapter is not None and bool( getattr(adapter, "_streaming_tts_turn_completed", lambda *_a, **_k: False)(session_key, run_generation) ) @@ -1723,9 +1667,8 @@ class GatewayTurnMixin: ): await self._send_voice_reply(event, response) - # Streamed responses still need MEDIA: files delivered before returning None (chunks carry - # the tags verbatim and post-processing is skipped when already_sent). Never skip when the - # agent failed: the error text is new content streaming didn't show. + # Streamed responses still need MEDIA: files delivered (chunks carry the tags verbatim). Never + # skip when the agent failed: the error text is new content streaming didn't show. if agent_result.get("already_sent") and not agent_result.get("failed"): if response and adapter: await self._deliver_media_from_response(response, event, adapter) @@ -1735,9 +1678,8 @@ class GatewayTurnMixin: await adapter.send(source.chat_id, _footer_line, metadata=self._event_thread_metadata(event, source)) except Exception as _e: logger.debug("trailing footer send failed: %s", _e) - # Return None so the adapter does not send the body twice; /loop and /goal hooks in - # _handle_message read the return value, so stash the delivered text on the event or - # those hooks never run and a /loop tick stays awaiting. + # Return None so the body isn't sent twice; stash the delivered text on the event for the + # /loop and /goal hooks that read the return value. with suppress(Exception): event._streamed_final_response = str(response or "") return None @@ -1756,9 +1698,8 @@ class GatewayTurnMixin: # Retain Slack thread/workspace routing so a failed turn cannot leave its status visible. await self._hmwa_stop_typing_for_turn(event, source) logger.exception("Agent error in session %s", session_key) - # Crash-resilience for failures before AIAgent enters run_conversation() (e.g. provider/ - # httpx client init): the agent can't persist the inbound turn there, so append the user - # message here once; skip if the latest user row already matches it. + # Failures before run_conversation() (provider/httpx init) can't persist the inbound turn: + # append the user message here once, unless the latest user row already matches it. try: if prepared.message_text is not None and session_entry is not None: try: @@ -1778,8 +1719,7 @@ class GatewayTurnMixin: ) except Exception: logger.debug("Failed to persist inbound user message after agent exception", exc_info=True) - # Log full details server-side only; never expose raw exception types or messages to end - # users (info-leakage risk). + # Never expose raw exception types/messages to end users (info-leakage risk). status_code = getattr(e, "status_code", None) status_hint = self._STATUS_HINTS.get(status_code, "") if status_code == 429: @@ -1798,8 +1738,7 @@ class GatewayTurnMixin: else: status_hint = " Your plan's usage limit has been reached. Please wait until it resets." elif status_code in {400, 500}: - # 400 on a large session is context overflow; 500 on a large session often means the - # payload is too large for the API — treat it the same way. + # 400/500 on a large session: context overflow / payload too large. if len(prepared.history) > 50: return ( "⚠️ Session too large for the model's context window.\nUse /compact to " @@ -1847,9 +1786,8 @@ class GatewayTurnMixin: with suppress(Exception): _redact_pii = bool((_load_gateway_config().get("privacy") or {}).get("redact_pii", False)) - # The context prompt render is pinned per session, keyed by a hash of the exact renderer - # inputs (_ephemeral_change_key): a hit reuses the pinned bytes so the system prompt cannot - # drift turn-over-turn; a miss (thread rename, /sethome, redact_pii flip) re-renders. + # The context prompt render is pinned per session, keyed by a hash of the renderer inputs, so + # the system prompt cannot drift turn-over-turn; a miss (thread rename, /sethome) re-renders. context_prompt = self._pinned_session_context_prompt(context, _redact_pii, session_key) # Per-turn notes ride the user message via the api_content sidecar, NOT context_prompt @@ -1858,17 +1796,15 @@ class GatewayTurnMixin: if _was_auto_reset: await self._hmwa_deliver_auto_reset_notice(session_entry, source, turn_sidecar_notes) - # Auto-load skill(s) for topic/channel bindings (single name or ordered list) — only on NEW - # sessions; ongoing conversations already carry the skill content in their history. + # Auto-load bound skill(s) only on NEW sessions; ongoing ones carry the content in history. _auto = getattr(event, "auto_skill", None) if _is_new_session and _auto: self._hmwa_auto_load_skills(event, _auto, _quick_key, session_key) await self._hmwa_acquire_turn_lease(_quick_key, run_generation, session_entry, _session_env_tokens) - # A turn only becomes durable recovery work after it owns (or has explicitly degraded past) - # the per-session lease. Marking before the await above would falsely recover an alias- - # routed message that never began processing if the gateway died while it was still waiting. + # A turn becomes durable recovery work only after it owns the per-session lease; marking + # earlier would falsely recover a message that never began processing. await self._mark_durable_active_turn(event, session_entry.session_key) # An unreadable store is not an empty conversation: stop before the agent invents continuity @@ -1895,8 +1831,7 @@ class GatewayTurnMixin: if _vc_note: turn_sidecar_notes.append(_vc_note) - # Auto-analyze user images (vision tool eagerly, image media_type only) so the model always - # gets a text description plus the local path for re-examination via vision_analyze. + # Auto-analyze user images so the model gets a description plus the local path. message_text = await self._prepare_profile_scoped_inbound_message_text( event=event, source=source, history=history, session_key=session_key, ) @@ -1907,13 +1842,13 @@ class GatewayTurnMixin: self._hmwa_apply_message_timestamp(event, message_text) ) - # Stage this turn's must-deliver notes (one-shot; consumed in run_sync) AFTER the - # message_text early-out so an aborted turn cannot leak its notes into the next turn. + # Stage the notes (one-shot; consumed in run_sync) AFTER the early-out so an aborted turn + # cannot leak them into the next turn. if turn_sidecar_notes and session_key: self._set_pending_turn_sidecar_notes(session_key, turn_sidecar_notes) - # Bind this run generation to the adapter's active-session event so deferred post-delivery - # callbacks can be released by the same run that registered them. + # Bind this run generation to the adapter so deferred post-delivery callbacks are released + # by the run that registered them. self._bind_adapter_run_generation(self._adapter_for_source(source), session_key, run_generation) return self._PreparedTurn( history, context_prompt, message_text, persist_user_message, persist_user_timestamp, @@ -1962,9 +1897,8 @@ class GatewayTurnMixin: } await self.hooks.emit("agent:start", hook_ctx) - # Capture the session id this run launches against so post-run compression publication - # can be identity-guarded; a /new or another lifecycle transition may move - # session_entry.session_id while the old run is still unwinding. + # Capture the launch session id so post-run compression publication is identity-guarded + # (a /new may move session_entry.session_id while the old run is still unwinding). _run_start_session_id = session_entry.session_id _turn_started_monotonic = time.monotonic() agent_result = await self._run_agent( @@ -2217,8 +2151,8 @@ class GatewayTurnMixin: await adapter.send_image( chat_id=source.chat_id, image_url=image_url, caption=alt_text, metadata=_thread_metadata, ) - # Route each media file by type so a TTS clip arrives as a voice bubble and a clip as a - # video rather than a generic document. Mirrors the streaming + kanban paths. + # Route each media file by type (voice bubble / video / image / document), as the + # streaming + kanban paths do. from gateway.platforms.base import should_send_media_as_audio as _should_send_media_as_audio from gateway.run_notifications import _IMAGE_EXTS, _VIDEO_EXTS for media_path, _is_voice in (media_files or []): @@ -2255,8 +2189,7 @@ class GatewayTurnMixin: _cache_lock = getattr(self, "_agent_cache_lock", None) if _cache_lock is None or not _cache: return - # Multiplex: only this profile's sessions; rebuilding another profile's agent in - # this scope would hand it this profile's tool registry. + # Multiplex: only this profile's sessions (another profile's agent would get this registry). _ns_prefix = _session_key_namespace(profile) + ":" if multiplex else None with _cache_lock: for _sess_key, _entry in list(_cache.items()): @@ -2324,8 +2257,7 @@ class GatewayTurnMixin: self._mcp_reload_refresh_cached_agents(multiplex, event.source.profile) - # Inject a message at the END of the session history so the model knows tools changed - # next turn; appending after all existing messages preserves the prompt-cache prefix. + # Append a note at the END of the history (preserves the prompt-cache prefix). change_parts = [ f"{label} servers: {', '.join(sorted(names))}" for label, names in (("Added", added), ("Removed", removed), ("Reconnected", reconnected)) @@ -2375,12 +2307,11 @@ class GatewayTurnMixin: if not _adapter_supports_edit and not _adapter_supports_native_stream and on_missing_cursor == "raise": raise RuntimeError("skip streaming for non-editable platform") _effective_cursor = scfg.cursor if _adapter_supports_edit else "" - # Some Matrix clients render the cursor as a tofu/white-box artifact: stream text, no cursor. + # Some Matrix clients render the cursor as tofu: stream text, no cursor. _buffer_only = source.platform == Platform.MATRIX if _buffer_only: _effective_cursor = "" - # Fresh-final applies to Telegram only — other platforms edit in place cheaply (Discord, Slack) - # or lack the edit-timestamp-stays-stale problem. + # Fresh-final applies to Telegram only (others edit in place cheaply). _fresh_final_secs = ( float(getattr(scfg, "fresh_final_after_seconds", 0.0) or 0.0) if source.platform == Platform.TELEGRAM else 0.0 @@ -2456,8 +2387,7 @@ class GatewayTurnMixin: if not proxy_url: return self._proxy_error_result("⚠️ Proxy URL not configured (GATEWAY_PROXY_URL or gateway.proxy_url)") - # Scope-aware read: the proxy key is a per-profile credential; under multiplex honor the - # installed scope's verdict (Slack pattern for the unscoped default-profile loop). + # The proxy key is a per-profile credential: honor the installed secret scope under multiplex. proxy_key = os.getenv("GATEWAY_PROXY_KEY", "").strip() with suppress(Exception): # UnscopedSecretError and import failures fall back to the env from agent.secret_scope import get_secret @@ -2475,9 +2405,8 @@ class GatewayTurnMixin: "history_offset": len(history), "session_id": session_id, "response_previewed": False, } - # OpenAI chat format. The remote api_server keeps continuity via X-Hermes-Session-Id and - # loads its own history, so send only the current message plus a compact text-only local - # history for a remote that has none yet (remote replays tools). + # OpenAI chat format. The remote keeps continuity via X-Hermes-Session-Id; send the current + # message plus a compact text-only history for a remote that has none yet. api_messages: List[Dict[str, str]] = [] if context_prompt: api_messages.append({"role": "system", "content": context_prompt}) @@ -2597,8 +2526,7 @@ class GatewayTurnMixin: platform_key = _platform_config_key(source.platform) enabled_toolsets, disabled_toolsets = self._resolve_turn_toolsets(user_config, source, platform_key) adapter = self._adapter_for_source(source) - # Per-platform display settings: display.platforms.., then display. - # global, then built-in platform defaults. + # display.platforms.. → display. → built-in platform defaults. _display_cfg = user_config.get("display", {}) if not isinstance(_display_cfg, dict): _display_cfg = {} @@ -2613,8 +2541,7 @@ class GatewayTurnMixin: _val = resolve_display_setting(user_config, platform_key, _setting, _default) getattr(_agent_display, _setter)(_cast(_val)) - # Tool progress mode — per-platform, with HERMES_TOOL_PROGRESS_MODE winning only when the - # config never set it. + # Tool progress mode; HERMES_TOOL_PROGRESS_MODE wins only when the config never set it. _resolved_tp = resolve_display_setting(user_config, platform_key, "tool_progress") _env_tp = os.getenv("HERMES_TOOL_PROGRESS_MODE") _platform_cfg = (_display_cfg.get("platforms") or {}).get(platform_key) or {} @@ -2658,21 +2585,18 @@ class GatewayTurnMixin: logger.debug("generic status phrase selection failed: %s", _phrase_err) return "still on it" if kind in {"heartbeat", "waiting", "long_running", "status"} else "one sec" - # Webhooks can't edit messages, so tool progress / log mode are off there (each progress - # line would be its own message). + # Webhooks can't edit messages, so tool progress / log mode are off there. is_webhook = source.platform == Platform.WEBHOOK tool_progress_enabled = progress_mode not in {"off", "log"} and not is_webhook - # Live working-state status for text-rendering typing indicators (Slack's assistant status - # line). Independent of tool_progress; rides the existing _keep_typing refresh. + # Live status for text-rendering typing indicators (Slack); independent of tool_progress. _live_status_mode = resolve_display_setting(user_config, platform_key, "live_status", "full") _live_status_adapter = ( adapter if getattr(adapter, "supports_status_text", False) and _live_status_mode != "off" else None ) # "log" mode: tool calls go to ~/.hermes/logs/tool_calls.log instead of the chat. Gateway-only. log_mode_enabled = progress_mode == "log" and not is_webhook - # Natural assistant status messages are independent from tool progress and token streaming. - # thinking_progress is independent too (same queue). Mattermost requires a per-platform - # opt-in for both: global scratch-text display leaks too easily into busy public threads. + # Interim assistant messages and thinking_progress are independent of tool progress (same + # queue). Mattermost requires a per-platform opt-in: scratch text leaks into public threads. interim_assistant_messages_mode = _display_surface_mode( "interim_assistant_messages", default=True, require_platform_override_for={Platform.MATTERMOST}, ) @@ -2680,8 +2604,7 @@ class GatewayTurnMixin: _thinking_enabled = _display_surface_mode( "thinking_progress", default=False, require_platform_override_for={Platform.MATTERMOST}, ) != "off" - # Slack-native task cards render tool progress via chat.startStream, so the progress queue is - # needed even though Slack keeps text tool_progress off by default. + # Slack-native task cards need the progress queue even with text tool_progress off. _native_slack_task_cards = False if source.platform == Platform.SLACK and hasattr(adapter, "native_task_cards_enabled"): try: @@ -2732,9 +2655,8 @@ class GatewayTurnMixin: _voice_ack_guild[0] = _gid break - # Auto-cleanup of temporary progress bubbles (any adapter implementing ``delete_message``; - # failed runs skip cleanup so the bubbles remain as breadcrumbs). getattr on the type, not - # attribute access: a fake adapter without delete_message means "can't delete", not a crash. + # Auto-cleanup of temporary progress bubbles needs a real ``delete_message`` (getattr on the + # type: a fake adapter without it means "can't delete", not a crash). _cleanup_progress = bool( disp.resolve_display_setting(disp.user_config, disp.platform_key, "cleanup_progress") ) @@ -2803,8 +2725,7 @@ class GatewayTurnMixin: if is_buzz: _progress_reply_in_thread = getattr(_adapter, "_reply_to_mode", "first") != "off" else: - # Relay lane: adapter owns mode resolution (nested platforms.relay.extra.slack - # subset, flat-key fallback). Native lane: read the flat extra as before. + # Relay lane: the adapter owns mode resolution; native lane: flat extra key. _mode_fn = getattr(_adapter, "_effective_reply_in_thread", None) _progress_reply_in_thread = bool( _mode_fn() if callable(_mode_fn) else _adapter.config.extra.get("reply_in_thread", True) @@ -2814,8 +2735,7 @@ class GatewayTurnMixin: _progress_thread_id = _resolve_progress_thread_id( source.platform, source.thread_id, event_message_id, reply_in_thread=_progress_reply_in_thread, ) - # Relay Discord auto-thread lane: a channel-initiating message has no thread_id at ingest; - # the connector stamps prospective_thread_id (anchor id == the thread it will create). + # Relay Discord auto-thread lane: the connector stamps prospective_thread_id at ingest. _relay_prospective_thread_id = ( str(getattr(source, "prospective_thread_id", None)) if source.platform == Platform.DISCORD @@ -2838,8 +2758,7 @@ class GatewayTurnMixin: _progress_metadata.setdefault("slack_team_id", source.scope_id) if source.user_id: _progress_metadata.setdefault("recipient_user_id", source.user_id) - # Buzz has no native thread_id; threading is always via reply-to the triggering event id - # (channel clutter otherwise), skipped when the user opted out of threaded replies. + # Buzz has no native thread_id: thread via reply-to unless the user opted out. _progress_reply_to = ( event_message_id if (source.platform in (Platform.FEISHU, Platform.MATTERMOST) and source.thread_id and event_message_id) @@ -2896,8 +2815,7 @@ class GatewayTurnMixin: _relay_prospective_thread_id: Optional[str], ) -> Optional[Dict[str, Any]]: """Thread metadata for status / approval / stream sends; Feishu topics need the triggering - message id (reply API with reply_in_thread=true) so carry it as a fallback. - """ + message id (reply API with reply_in_thread=true) so carry it as a fallback.""" if source.platform == Platform.FEISHU and source.thread_id and event_message_id: return {"thread_id": _progress_thread_id, "reply_to_message_id": event_message_id} return self._thread_metadata_for_progress( @@ -2954,8 +2872,7 @@ class GatewayTurnMixin: async def _run_agent_track_agent(self, turn_ctx: TurnContext) -> None: """Track this agent as running for the session (interrupt support) once it is created — only - if this run is still current, else leave the newer run's slot alone. - """ + if this run is still current, else leave the newer run's slot alone.""" session_key, run_generation, agent_holder = turn_ctx.session_key, turn_ctx.run_generation, turn_ctx.agent_holder while agent_holder[0] is None: await asyncio.sleep(0.05) @@ -3094,15 +3011,13 @@ class GatewayTurnMixin: _agent_timeout = _agent_timeout if _agent_timeout > 0 else None _agent_warning = _agent_warning if _agent_warning > 0 else None - # A background=true process intentionally survives a successful turn, so capture existing IDs - # and reap only children created by THIS turn if it times out. + # background=true processes survive a turn: reap only children created by THIS turn on timeout. _turn_task_id = turn_ctx.session_id or "" _turn_process_baseline = process_registry.snapshot_running_ids(_turn_task_id) turn_ctx.process_task_id = _turn_task_id turn_ctx.process_baseline = _turn_process_baseline - # task_id is session-scoped, not turn-scoped: gate the eventual reap on this exact claim still - # being current, so a replacement turn on the same session that starts before the watchdog - # fires doesn't get its own fresh process killed by this turn's stale baseline. + # task_id is session-scoped: gate the reap on this claim still being current so a replacement + # turn's fresh process isn't killed by this turn's stale baseline. worker = self._RunAgentWorker( agent_timeout=_agent_timeout, agent_warning=_agent_warning, task_id=_turn_task_id, process_baseline=_turn_process_baseline, worker_done=threading.Event(), @@ -3242,9 +3157,8 @@ class GatewayTurnMixin: while True: done, _ = await asyncio.wait({worker.executor_task}, timeout=5.0) if done: - # Prefer the real result even if the watchdog fired in the same window: the completed - # run already persisted its reply, so the "agent inactive" diagnostic would contradict - # the stored transcript. + # Prefer the real result even if the watchdog fired in the same window (the run already + # persisted its reply). return worker.executor_task.result() if worker.agent_timeout is not None: if worker.timeout_fired.is_set(): @@ -3276,9 +3190,8 @@ class GatewayTurnMixin: if _agent is None or not hasattr(_agent, 'model') or _run_failed: return _cfg_model = _resolve_gateway_model() - # Normalize _cfg_model as AIAgent.__init__ does so a vendor-prefixed config value matches - # the agent's stripped model on native providers — otherwise the cached agent is evicted - # every turn, destroying prompt caching. Aggregators keep the vendor slug. + # Normalize as AIAgent.__init__ does (vendor prefix stripped on native providers), else the + # cached agent is evicted every turn, destroying prompt caching. with suppress(Exception): from hermes_cli.model_normalize import _AGGREGATOR_PROVIDERS, normalize_model_for_provider _agent_provider = getattr(_agent, 'provider', '') or '' @@ -3335,8 +3248,7 @@ class GatewayTurnMixin: else: pending = interrupt_message elif pending_event: - # Transcribe audio on the dequeued event BEFORE it becomes the next user turn, so - # queued/interrupting voice messages drain with the real transcript, not a file path. + # Transcribe audio BEFORE it becomes the next user turn (real transcript, not a path). _pending_text = pending_event.text or "" if self._pending_event_audio_paths(pending_event): pending, _ = await self._transcribe_and_echo_pending_voice( @@ -3349,14 +3261,12 @@ class GatewayTurnMixin: if pending: logger.debug("Processing queued message after agent completion: '%s...'", pending[:40]) - # Leftover /steer: a steer arriving after the last tool batch (e.g. during the final API - # call) comes back in result["pending_steer"]; deliver it as the next user turn, not drop it. + # Leftover /steer (arrived after the last tool batch): deliver as the next user turn. if result and not pending and not pending_event and result.get("pending_steer"): pending = result.get("pending_steer") logger.debug("Delivering leftover /steer as next turn: '%s...'", pending[:40]) - # Safety net: a pending slash command (e.g. "/stop", "/new") is discarded — commands must - # never be passed to the agent as user input. + # Safety net: a pending slash command is never passed to the agent as user input. if pending and pending.strip().startswith("/"): _pending_parts = pending.strip().split(None, 1) _pending_cmd_word = _pending_parts[0][1:].lower() if _pending_parts else "" @@ -3391,16 +3301,13 @@ class GatewayTurnMixin: await self._await_stream_task(stream_task) except Exception as e: logger.debug("Stream consumer wait before queued message failed: %s", e) - # The queued branch needs raw ``result`` for interruption, history, and recursion state, but - # delivery must use the finalized task result — it carries empty/failure normalization and - # final-response processing from _run_agent_task. + # Delivery uses the finalized task result (empty/failure normalization), not raw ``result``. _delivery_result = response if isinstance(response, dict) else (result or {}) first_response = _delivery_result.get("final_response", "") _already_streamed = self._run_agent_stream_confirmed_final_delivery( _sc, first_response, previewed=bool(_delivery_result.get("response_previewed")), ) - # Same predicate as the normal completed-turn path: this direct queued-send branch predates - # intentional-silence filtering and would leak the literal marker. + # Same silence predicate as the normal path, else this branch leaks the literal marker. if self._is_intentional_silence(_delivery_result, first_response): logger.info( "Queued follow-up for session %s: suppressing intentional silence marker before continuing.", @@ -3422,8 +3329,7 @@ class GatewayTurnMixin: ) except Exception as e: logger.warning("Failed to send first response before queued message: %s", e) - # Release deferred bg-review notifications now that the first response is delivered: pop from - # the adapter's callback dict (no double-fire in base.py's finally) and call. + # Release deferred bg-review notifications: pop (no double-fire in base.py's finally) and call. _bg_cb = self._pop_post_delivery_callback(adapter, session_key, turn_ctx.run_generation) if callable(_bg_cb): with suppress(Exception): @@ -3445,13 +3351,11 @@ class GatewayTurnMixin: ) logger.debug("Processing pending message: '%s...'", pending[:40]) - # Clear the adapter's interrupt event so the next _run_agent call doesn't re-trigger the - # interrupt before the new agent's first API call (infinite loop otherwise). + # Clear the interrupt event so the recursive _run_agent isn't re-interrupted (infinite loop). if adapter and hasattr(adapter, '_active_sessions') and session_key and session_key in adapter._active_sessions: adapter._active_sessions[session_key].clear() - # Cap recursion depth: resource exhaustion when the user keeps sending messages while the - # agent keeps failing. + # Cap recursion depth (user keeps sending while the agent keeps failing). if _interrupt_depth >= self._MAX_INTERRUPT_DEPTH: logger.warning( "Interrupt recursion depth %d reached for session %s — " @@ -3464,8 +3368,7 @@ class GatewayTurnMixin: adapter.queue_message(session_key, pending) return turn_ctx.result_holder[0] or {"final_response": response, "messages": history} - # Interrupted: discard the response ("Operation interrupted." is noise; the user knows they - # sent a new message). + # Interrupted: discard the response ("Operation interrupted." is noise). if not result.get("interrupted"): await self._run_agent_deliver_first_response(turn_ctx, adapter, response, result, stream_task) @@ -3481,9 +3384,8 @@ class GatewayTurnMixin: session_key or "?", ) return result - # Resolve the follow-up's session key BEFORE preparing the inbound text: - # _prepare_inbound_message_text buffers native image paths under the key given, and the - # recursive _run_agent consumes them under next_session_key — mismatch drops them. + # Resolve the follow-up's session key BEFORE preparing the inbound text: native image + # paths are buffered under the key given and consumed under next_session_key. try: next_session_key = self._session_key_for_source(next_source) except Exception: @@ -3500,8 +3402,7 @@ class GatewayTurnMixin: next_channel_prompt = getattr(pending_event, "channel_prompt", None) next_message_type = getattr(pending_event, "message_type", None) - # Clear the prior logical turn's completed streaming marker so the recursive turn's streaming - # TTS isn't suppressed by that completion. + # Clear the prior turn's streaming-TTS completion marker so the recursive turn isn't suppressed. _clear_adapter = self._adapter_for_source(source) if _clear_adapter is not None and session_key and run_generation is not None: _completed_turns = getattr(_clear_adapter, "_streaming_tts_completed_turns", None) @@ -3511,15 +3412,13 @@ class GatewayTurnMixin: if _pk: _completed_turns.discard(_pk) - # Restart the typing indicator for the follow-up turn; the outer _process_message_background - # typing task is alive but may be stale. + # Restart the typing indicator; the outer typing task may be stale. if _clear_adapter: with suppress(Exception): await _clear_adapter.send_typing(source.chat_id, metadata=_status_thread_metadata) - # Re-baseline the cached agent's message_count before recursing into the /queue follow-up: - # the coherence guard would otherwise rebuild on OUR OWN flushed rows and destroy the - # prompt-cache prefix; _handle_message_with_agent re-baselines only after the chain ends. + # Re-baseline the cached agent's message_count before recursing, else the coherence guard + # rebuilds on OUR OWN flushed rows (the outer handler re-baselines only after the chain). await self._refresh_agent_cache_message_count(session_key, session_id) followup_result = await self._run_agent( @@ -3542,9 +3441,7 @@ class GatewayTurnMixin: task.cancel() if stream_task: - # If the agent never created a stream consumer (non-streaming path, or a test stub returning - # synchronously) there is nothing to flush — cancel now instead of waiting out the 5s - # timeout polling for a consumer that will never arrive. + # No stream consumer was created: nothing to flush, cancel instead of waiting out 5s. if not (stream_consumer_holder and stream_consumer_holder[0] is not None): stream_task.cancel() with suppress(asyncio.CancelledError): @@ -3552,8 +3449,7 @@ class GatewayTurnMixin: else: await self._await_stream_task(stream_task) - # Unconditional abort + bounded wait for the streaming-TTS consumer: covers cancellation / - # exception paths where the normal finalisation block was skipped. + # Abort + bounded wait for streaming TTS: covers paths where normal finalisation was skipped. _stts_finally = turn_ctx.streaming_tts_consumer_holder[0] if _stts_finally is not None and not _stts_finally.done: _stts_finally.abort("cleanup") @@ -3562,8 +3458,8 @@ class GatewayTurnMixin: tracking_task.cancel() if session_key: - # Release the slot only if this run's generation still owns it: a /stop or /new that bumped - # the generation while we unwound already installed its own state; keep it. + # Release the slot only if this run's generation still owns it (/stop or /new may have + # installed its own state). self._release_running_agent_state(session_key, run_generation=turn_ctx.run_generation) if self._draining: self._update_runtime_status("draining") @@ -3575,9 +3471,7 @@ class GatewayTurnMixin: except asyncio.CancelledError: pass except Exception: - # A background task that died of a non-cancellation error (transport drop in a - # progress/card publish) must not abort the cleanup path — final-delivery - # bookkeeping after this loop still runs. + # A background task that died of a real error must not abort the cleanup path. logger.debug("background turn task failed during cleanup", exc_info=True) async def _run_agent_mark_streamed_delivery(self, response: Any, turn_ctx: TurnContext) -> None: @@ -3593,8 +3487,7 @@ class GatewayTurnMixin: return _final = response.get("final_response") or "" _is_empty_sentinel = not _final or _final == "(empty)" - # response_previewed means interim_assistant_callback already saw the final text, but only - # suppress the send if that exact text was delivered — unrelated commentary/progress isn't it. + # response_previewed: only suppress if that EXACT text was delivered, not unrelated commentary. _previewed = bool(response.get("response_previewed")) _content_delivered = bool(_sc and getattr(_sc, "final_content_delivered", False)) _stale_finalized = False @@ -3607,11 +3500,9 @@ class GatewayTurnMixin: _stale_finalized = False if _stale_finalized: _content_delivered = False - # Plugin hooks (e.g. transform_llm_output) may append content after streaming finished — when - # transformed, always send the final version so the appended content reaches the client. + # Plugin hooks may append content after streaming finished — then send the final version. _transformed = bool(response.get("response_transformed")) - # Suppress the normal send only when the actual final reply reached the user (streamed, or - # interim preview of that *exact* text); commentary shown during a compression/split isn't it. + # Suppress the normal send only when the actual final reply reached the user. _streamed = self._run_agent_stream_confirmed_final_delivery(_sc, _final, previewed=_previewed) if _is_empty_sentinel: return @@ -3658,8 +3549,7 @@ class GatewayTurnMixin: _sk, ) elif _transformed and _sc is not None: - # Plugin hooks transformed the response after streaming — edit the existing streamed - # message instead of sending a duplicate. + # Transformed after streaming: edit the streamed message instead of sending a duplicate. _sc_msg_id = _sc.message_id if _sc_msg_id: try: @@ -3674,9 +3564,8 @@ class GatewayTurnMixin: except Exception as _edit_err: logger.warning("Failed to edit streamed message for session %s: %s", _sk, _edit_err) elif _sc is not None: - # DUPLICATE-RISK DIAGNOSTIC: a stream consumer existed but suppression did NOT fire, so the - # normal final-send is about to run. Log the decision inputs so a recurrence can be pinned - # to "signal never set" vs "ack-pending race". + # DUPLICATE-RISK DIAGNOSTIC: a stream consumer existed but suppression did NOT fire; log the + # decision inputs ("signal never set" vs "ack-pending race"). logger.warning( "Normal final-send NOT suppressed despite active stream consumer for session %s: " "streamed=%s previewed=%s content_delivered=%s transformed=%s final_len=%d — " @@ -3735,8 +3624,7 @@ class GatewayTurnMixin: turn_ctx._progress_metadata, turn_ctx._progress_reply_to, _progress_thread_id, _relay_prospective_thread_id, ) = self._run_agent_progress_threading(source, event_message_id, _native_slack_task_cards) - # Bridges: sync step_callback / event_callback / status_callback → async hooks.emit and - # adapter.send (TurnRunner methods bound onto the shared ctx). + # Bridges: sync step/event/status callbacks → async hooks.emit and adapter.send. turn_ctx._loop_for_step = asyncio.get_running_loop() turn_ctx._hooks_ref = self.hooks turn_ctx._step_callback_sync = turn_runner._step_callback_sync @@ -3779,8 +3667,7 @@ class GatewayTurnMixin: ): break _elapsed_mins = int((time.time() - _notify_start) // 60) - # Default heartbeat is terse (elapsed + current tool); the verbose iteration counter is - # gated on busy_ack_detail so users can opt in per platform. + # Terse heartbeat by default; the iteration counter is gated on busy_ack_detail. _status_detail = "" _want_iteration_detail = bool( disp.resolve_display_setting(disp.user_config, disp.platform_key, "busy_ack_detail", True) @@ -3861,12 +3748,10 @@ class GatewayTurnMixin: source, message_type, _status_thread_metadata, turn_ctx.streaming_tts_consumer_holder, ) - # Progress sender gates on needs_progress_queue (tool_progress OR thinking_progress), not - # tool_progress alone: it drains BOTH tool-progress lines and _thinking scratch bubbles. + # Progress sender drains BOTH tool-progress lines and thinking bubbles (needs_progress_queue). progress_task = asyncio.create_task(turn_runner.send_progress_messages()) if disp.needs_progress_queue else None log_task = asyncio.create_task(self._run_agent_write_tool_log(disp.log_queue)) if disp.log_mode_enabled else None - # The stream consumer is created inside run_sync (thread pool) after the agent is constructed; - # this task polls for it. + # The stream consumer is created inside run_sync; this task polls for it. stream_task = asyncio.create_task(self._run_agent_stream_consumer_task(turn_ctx.stream_consumer_holder)) tracking_task = asyncio.create_task(self._run_agent_track_agent(turn_ctx)) _interrupt_detected = asyncio.Event() # shared with backup check