From 98a528a65a01e5467d6c0007817a120dafeec43a Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 15:50:25 -0700 Subject: [PATCH] refactor(gateway): split _handle_message_with_agent into _hmwa_* helpers (2365 -> 242 LOC, max helper 249) --- gateway/run.py | 4202 ++++++++++++++++++++++++++---------------------- 1 file changed, 2248 insertions(+), 1954 deletions(-) diff --git a/gateway/run.py b/gateway/run.py index d05b1d49ec..3d57b02bf1 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -18187,19 +18187,45 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew cached_sources.move_to_end(session_key) return source - async def _handle_message_with_agent(self, event, source, _quick_key: str, run_generation: int): - """Inner handler that runs under the _running_agents sentinel guard.""" - _msg_start_time = time.time() - _platform_name = source.platform.value if hasattr(source.platform, "value") else str(source.platform) - _msg_preview = (event.text or "")[:80].replace("\n", " ") - _reply_id = getattr(event, "reply_to_message_id", None) - _reply_txt = (getattr(event, "reply_to_text", None) or "")[:80].replace("\n", " ") - logger.info( - "inbound message: platform=%s user=%s chat=%s msg=%r reply_to_id=%s reply_to_text=%r", - _platform_name, source.user_name or source.user_id or "unknown", - source.chat_id or "unknown", _msg_preview, _reply_id, _reply_txt, - ) + @dataclasses.dataclass + class _HygieneSettings: + """Resolved session-hygiene configuration for one inbound turn.""" + model: str + threshold_pct: float + compression_enabled: bool + hard_msg_limit: int + timeout_seconds: float + total_ceiling_seconds: float + max_turn_hold_seconds: float + failure_cooldown_seconds: float + config_context_length: Optional[int] + provider: Optional[str] + base_url: Optional[str] + api_key: Optional[str] + data: Any + + @dataclasses.dataclass + class _HygieneAttempt: + """One detached hygiene compression attempt (agent, worker future, commit fence). + + ``cleanup_deferred`` is shared mutable state: the wait handlers set it on their raise + paths and the owning ``finally`` reads it to decide whether to clean the agent up now. + """ + + agent: Any + meta: Any + commit_fence: Any = None + future: Any = None + wait_started: float = 0.0 + cleanup_deferred: bool = False + history: Any = None + + + async def _hmwa_resolve_session(self, event, source): + """Resolve ``source`` to its session entry (topic recovery, internal-route guards, Telegram + topic-binding heal). Returns ``(source, session_entry, session_key)`` or ``None`` to drop + the event.""" # Get or create session 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 across sessions. @@ -18318,6 +18344,11 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew await asyncio.to_thread(self._record_telegram_topic_binding, source, session_entry) except Exception: logger.debug("Failed to record Telegram topic binding", exc_info=True) + return source, session_entry, session_key + + 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)``.""" # Capture and consume was_auto_reset immediately so it cannot re-fire on later messages and # wipe model/reasoning overrides set between turns. _was_auto_reset = getattr(session_entry, "was_auto_reset", False) @@ -18349,140 +18380,103 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew "session_id": session_entry.session_id, "session_key": session_key, }) + return _was_auto_reset, _is_new_session - # Build session context - context = build_session_context(source, self.config, session_entry) - - # Set session context variables for tools (task-local, concurrency-safe) - _session_env_tokens = self._set_session_env(context) - - # Read privacy.redact_pii from config (re-read per message) - _redact_pii = False - persist_user_message = None - persist_user_timestamp = None - # Synthetic self-injected turns (batch completions, watch notifications, resume wake-ups) - # arrive as MessageEvent(internal=True). Persist with display_kind="internal_notification" - # so UIs render timeline notices, not user bubbles. display_kind is a DB-only sidecar - # stripped from every provider-bound payload; role/content untouched. - persist_user_display_kind = ( - "internal_notification" if getattr(event, "internal", False) else None - ) + async def _hmwa_deliver_auto_reset_notice(self, session_entry, source, turn_sidecar_notes): + """Stage the auto-reset sidecar note for the agent and notify the user (policy-gated).""" + 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"]) + # Slack/Discord channels/threads are long-lived: point the agent at the specific prior + # same-channel session so it recalls that context via session_search instead of an + # unrelated recent session. Deterministic — no extra API/DB calls. try: - _pcfg = _load_gateway_config() - _redact_pii = bool((_pcfg.get("privacy") or {}).get("redact_pii", False)) + continuity_note = build_channel_continuity_note(session_entry, source) except Exception: - pass + continuity_note = None + if continuity_note: + context_note = context_note + "\n\n" + continuity_note + turn_sidecar_notes.append(context_note) - # Build the context prompt. The 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. - context_prompt = self._pinned_session_context_prompt( - context, _redact_pii, session_key - ) - - # Per-turn must-deliver notes ride the user message via the api_content sidecar (staged - # below, consumed in run_sync → build_turn_context), NOT context_prompt: appending them to - # the ephemeral system prompt guaranteed a turn1→turn2 diff and a full agent rebuild. - turn_sidecar_notes: List[str] = [] - - # If the previous session expired and was auto-reset, deliver a notice - # so the agent knows this is a fresh conversation (not an intentional /reset). - if _was_auto_reset: - 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"]) - # Slack/Discord channels/threads are long-lived: point the agent at the specific prior - # same-channel session so it recalls that context via session_search instead of an - # unrelated recent session. Deterministic — no extra API/DB calls. - try: - continuity_note = build_channel_continuity_note(session_entry, source) - except Exception: - continuity_note = None - if continuity_note: - context_note = context_note + "\n\n" + continuity_note - turn_sidecar_notes.append(context_note) - - # Notify the user about the reset unless notifications are disabled in config, the - # platform is excluded (e.g. api_server, webhook), or the expired session had no - # activity. - try: - policy = self.session_store.config.get_reset_policy( - platform=source.platform, - session_type=getattr(source, 'chat_type', 'dm'), - ) - platform_name = source.platform.value if source.platform else "" - had_activity = getattr(session_entry, 'reset_had_activity', False) - # Suspended and restart-recovery-expired sessions always notify regardless of - # policy.notify — the user had an active session that was silently replaced, so they - # need to know they can /resume it. Idle/daily resets respect the policy flag. - should_notify = reset_reason in {"suspended", "resume_pending_expired"} or ( - policy.notify - and had_activity - and platform_name not in policy.notify_exclude_platforms - ) - if should_notify: - adapter = self._adapter_for_source(source) - if adapter: - reason_text = _auto_reset_reason_text(reset_reason, policy) - notice = ( - f"◐ Session automatically reset ({reason_text}). " - f"Conversation history cleared.\n" - f"Use /resume to browse and restore a previous session.\n" - f"Adjust reset timing in config.yaml under session_reset." - ) - try: - session_info = await asyncio.to_thread( - self._reset_notice_session_info, source - ) - if session_info: - notice = f"{notice}\n\n{session_info}" - except Exception: - pass - await adapter.send( - source.chat_id, notice, - metadata=self._thread_metadata_for_source(source), - ) - except Exception as e: - logger.debug("Auto-reset notification failed (non-fatal): %s", e) - - # was_auto_reset is already consumed in the cleanup block above - # (single source of truth); only the reset reason needs clearing here. - session_entry.auto_reset_reason = None - - # Auto-load skill(s) for topic/channel bindings (Telegram DM Topics, Discord - # channel_skill_bindings). Supports a single name or ordered list. Only inject on NEW - # sessions — ongoing conversations already carry the skill content in their history. - _auto = getattr(event, "auto_skill", None) - if _is_new_session and _auto: - _skill_names = [_auto] if isinstance(_auto, str) else list(_auto) - try: - from agent.skill_commands import _load_skill_payload, _build_skill_message - _combined_parts: list[str] = [] - _loaded_names: list[str] = [] - for _sname in _skill_names: - _loaded = _load_skill_payload(_sname, task_id=_quick_key) - if _loaded: - _loaded_skill, _skill_dir, _display_name = _loaded - _note = ( - f'[IMPORTANT: The "{_display_name}" skill is auto-loaded. ' - f"Follow its instructions for this session.]" - ) - _part = _build_skill_message(_loaded_skill, _skill_dir, _note) - if _part: - _combined_parts.append(_part) - _loaded_names.append(_sname) - else: - logger.warning("[Gateway] Auto-skill '%s' not found", _sname) - if _combined_parts: - # Append the user's original text after all skill payloads - _combined_parts.append(event.text) - event.text = "\n\n".join(_combined_parts) - logger.info( - "[Gateway] Auto-loaded skill(s) %s for session %s", - _loaded_names, session_key, + # Notify the user about the reset unless notifications are disabled in config, the + # platform is excluded (e.g. api_server, webhook), or the expired session had no + # activity. + try: + policy = self.session_store.config.get_reset_policy( + platform=source.platform, + session_type=getattr(source, 'chat_type', 'dm'), + ) + platform_name = source.platform.value if source.platform else "" + had_activity = getattr(session_entry, 'reset_had_activity', False) + # Suspended and restart-recovery-expired sessions always notify regardless of + # policy.notify — the user had an active session that was silently replaced, so they + # need to know they can /resume it. Idle/daily resets respect the policy flag. + should_notify = reset_reason in {"suspended", "resume_pending_expired"} or ( + policy.notify + and had_activity + and platform_name not in policy.notify_exclude_platforms + ) + if should_notify: + adapter = self._adapter_for_source(source) + if adapter: + reason_text = _auto_reset_reason_text(reset_reason, policy) + notice = ( + f"◐ Session automatically reset ({reason_text}). " + f"Conversation history cleared.\n" + f"Use /resume to browse and restore a previous session.\n" + f"Adjust reset timing in config.yaml under session_reset." ) - except Exception as e: - logger.warning("[Gateway] Failed to auto-load skill(s) %s: %s", _skill_names, e) + try: + session_info = await asyncio.to_thread( + self._reset_notice_session_info, source + ) + if session_info: + notice = f"{notice}\n\n{session_info}" + except Exception: + pass + await adapter.send( + source.chat_id, notice, + metadata=self._thread_metadata_for_source(source), + ) + except Exception as e: + logger.debug("Auto-reset notification failed (non-fatal): %s", e) + # was_auto_reset is already consumed in the cleanup block above + # (single source of truth); only the reset reason needs clearing here. + session_entry.auto_reset_reason = None + + def _hmwa_auto_load_skills(self, event, _auto, _quick_key, session_key): + """Prepend topic/channel-bound skill payload(s) to ``event.text`` on a new session.""" + _skill_names = [_auto] if isinstance(_auto, str) else list(_auto) + try: + from agent.skill_commands import _load_skill_payload, _build_skill_message + _combined_parts: list[str] = [] + _loaded_names: list[str] = [] + for _sname in _skill_names: + _loaded = _load_skill_payload(_sname, task_id=_quick_key) + if _loaded: + _loaded_skill, _skill_dir, _display_name = _loaded + _note = ( + f'[IMPORTANT: The "{_display_name}" skill is auto-loaded. ' + f"Follow its instructions for this session.]" + ) + _part = _build_skill_message(_loaded_skill, _skill_dir, _note) + if _part: + _combined_parts.append(_part) + _loaded_names.append(_sname) + else: + logger.warning("[Gateway] Auto-skill '%s' not found", _sname) + if _combined_parts: + # Append the user's original text after all skill payloads + _combined_parts.append(event.text) + event.text = "\n\n".join(_combined_parts) + logger.info( + "[Gateway] Auto-loaded skill(s) %s for session %s", + _loaded_names, session_key, + ) + except Exception as e: + logger.warning("[Gateway] Failed to auto-load skill(s) %s: %s", _skill_names, e) + + async def _hmwa_acquire_turn_lease(self, _quick_key, run_generation, session_entry, _session_env_tokens): # ── Turn lease: session resolution is FINAL here. Serialize [load history → run → flush] # per resolved SESSION_ID: another routing key mapped to the same session_id waits for the # prior flush instead of loading a stale base and interleaving writes (same-key messages @@ -18513,1205 +18507,1324 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew _lease_state.lease_token = _lease_token _lease_state.lease_generation = run_generation - # 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. - await self._mark_durable_active_turn(event, session_entry.session_key) - - # Load conversation history from transcript. An unreadable canonical store is not an empty - # conversation: stop before the agent can invent continuity from a plausible-looking []. - # This return happens before the broad cleanup finally below, so restore task-local context - # here; the outer dispatch still clears the durable marker and turn lease. + async def _hmwa_hygiene_settings(self, source, session_key): + """Resolve model/provider/context-length + hygiene knobs for the pre-agent compression + safety net (fail-soft: any config/runtime error keeps the defaults).""" + # Read model + compression config. Hygiene threshold is intentionally HIGHER than the + # agent's own compressor (0.85 vs 0.50): it is a pre-agent safety net for sessions that + # grew between turns; at 0.50 it compressed prematurely on every turn in long sessions. + _hyg_model = "anthropic/claude-sonnet-4.6" + _hyg_threshold_pct = 0.85 + _hyg_compression_enabled = True + _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. The compressor keeps + # running detached; its commit is fenced (revoke_commit_admission) so a stale result can + # never clobber later turns. Kept well below transport idle-timeouts (Telegram ~30s). + _hyg_max_turn_hold_seconds = 10.0 + _hyg_failure_cooldown_seconds = 300.0 + _hyg_config_context_length = None + _hyg_provider = None + _hyg_base_url = None + _hyg_api_key = None + _hyg_configured_model = None + _hyg_configured_provider = None + _hyg_configured_base_url = None + _hyg_data = {} try: - history = await self.async_session_store.load_transcript( - session_entry.session_id + _hyg_data = _load_gateway_config() + if _hyg_data: + # Resolve model name (same logic as run_sync) + _model_cfg = _hyg_data.get("model", {}) + if isinstance(_model_cfg, str): + _hyg_model = _model_cfg + elif isinstance(_model_cfg, dict): + _hyg_model = _model_cfg.get("default") or _model_cfg.get("model") or _hyg_model + # Read explicit context_length override from model config + # (same as run_agent.py lines 995-1005) + _raw_ctx = _model_cfg.get("context_length") + if _raw_ctx is not None: + with suppress(TypeError, ValueError): + _hyg_config_context_length = int(_raw_ctx) + # Read provider for accurate context detection + _hyg_provider = _model_cfg.get("provider") or None + _hyg_base_url = _model_cfg.get("base_url") or None + + # Only the enabled flag is read; hygiene's threshold is deliberately separate + # from the agent's compression.threshold (hygiene runs higher). + _comp_cfg = _hyg_data.get("compression", {}) + if isinstance(_comp_cfg, dict): + _hyg_compression_enabled = str( + _comp_cfg.get("enabled", True) + ).lower() in {"true", "1", "yes"} + _raw_hard_limit = _comp_cfg.get("hygiene_hard_message_limit") + if _raw_hard_limit is not None: + try: + _parsed = int(_raw_hard_limit) + if _parsed > 0: + _hyg_hard_msg_limit = _parsed + except (TypeError, ValueError): + pass + _raw_timeout = _comp_cfg.get("hygiene_timeout_seconds") + if _raw_timeout is not None: + try: + _parsed = float(_raw_timeout) + if _parsed > 0: + _hyg_timeout_seconds = _parsed + except (TypeError, ValueError): + pass + _raw_ceiling = _comp_cfg.get("hygiene_total_ceiling_seconds") + if _raw_ceiling is not None: + try: + _parsed = float(_raw_ceiling) + if _parsed > 0: + _hyg_total_ceiling_seconds = _parsed + except (TypeError, ValueError): + pass + # The ceiling can never be tighter than one idle + # window, or the extension loop would be dead code. + _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: + _parsed = float(_raw_cooldown) + if _parsed >= 0: + _hyg_failure_cooldown_seconds = _parsed + except (TypeError, ValueError): + pass + + _hyg_configured_model = _hyg_model + _hyg_configured_provider = _hyg_provider + _hyg_configured_base_url = _hyg_base_url + + try: + _hyg_model, _hyg_runtime = self._resolve_session_agent_runtime( + source=source, + session_key=session_key, + user_config=_hyg_data if isinstance(_hyg_data, dict) else None, + ) + _hyg_provider = _hyg_runtime.get("provider") or _hyg_provider + _hyg_base_url = _hyg_runtime.get("base_url") or _hyg_base_url + _hyg_api_key = _hyg_runtime.get("api_key") or _hyg_api_key + except Exception: + pass + + if _hyg_config_context_length is not None: + try: + from hermes_cli.route_identity import should_clear_context_pin_async + + if await should_clear_context_pin_async( + _hyg_configured_model, + _hyg_model, + _hyg_configured_base_url, + _hyg_base_url, + _hyg_configured_provider, + _hyg_provider, + ): + _hyg_config_context_length = None + except Exception: + _hyg_config_context_length = None + + # custom_providers per-model context_length fallback (as in run_agent.py); must run + # after runtime resolution so _hyg_base_url is set. + if _hyg_config_context_length is None and _hyg_base_url: + try: + try: + from hermes_cli.config import ( + get_compatible_custom_providers as _gw_gcp, + get_custom_provider_context_length as _gw_gccl, + ) + _hyg_custom_providers = _gw_gcp(_hyg_data) + except Exception: + _hyg_custom_providers = _hyg_data.get("custom_providers") + if not isinstance(_hyg_custom_providers, list): + _hyg_custom_providers = [] + _hyg_custom_ctx = _gw_gccl( + model=_hyg_model, + base_url=_hyg_base_url, + custom_providers=_hyg_custom_providers, + ) + if _hyg_custom_ctx: + _hyg_config_context_length = int(_hyg_custom_ctx) + except (TypeError, ValueError): + pass + except Exception: + pass + return self._HygieneSettings( + model=_hyg_model, + threshold_pct=_hyg_threshold_pct, + compression_enabled=_hyg_compression_enabled, + hard_msg_limit=_hyg_hard_msg_limit, + timeout_seconds=_hyg_timeout_seconds, + total_ceiling_seconds=_hyg_total_ceiling_seconds, + max_turn_hold_seconds=_hyg_max_turn_hold_seconds, + failure_cooldown_seconds=_hyg_failure_cooldown_seconds, + config_context_length=_hyg_config_context_length, + provider=_hyg_provider, + base_url=_hyg_base_url, + api_key=_hyg_api_key, + data=_hyg_data, + ) + + async def _hmwa_hygiene_plan(self, hs, history, session_entry, session_key): + """Decide whether hygiene compression fires this turn (token/message thresholds, DB-backed + failure cooldown, in-flight compression). Returns + ``(_needs_compress, _approx_tokens, _msg_count, _warn_token_threshold)``.""" + from agent.model_metadata import ( + estimate_messages_tokens_rough, + get_model_context_length_async, + ) + + _hyg_model = hs.model + _hyg_threshold_pct = hs.threshold_pct + _hyg_hard_msg_limit = hs.hard_msg_limit + _hyg_config_context_length = hs.config_context_length + _hyg_provider = hs.provider + _hyg_base_url = hs.base_url + _hyg_api_key = hs.api_key + + _hyg_context_length = await get_model_context_length_async( + _hyg_model, + base_url=_hyg_base_url or "", + api_key=_hyg_api_key or "", + config_context_length=_hyg_config_context_length, + provider=_hyg_provider or "", + ) + _compress_token_threshold = int( + _hyg_context_length * _hyg_threshold_pct + ) + _warn_token_threshold = int(_hyg_context_length * 0.95) + + _msg_count = len(history) + + # Prefer actual API-reported tokens from the last turn + # (stored in session entry) over the rough char-based estimate. + _stored_tokens = session_entry.last_prompt_tokens + if _stored_tokens > 0: + _approx_tokens = _stored_tokens + _token_source = "actual" + else: + _approx_tokens = estimate_messages_tokens_rough(history) + _token_source = "estimated" + # Rough estimates run 30-50% high on code/JSON-heavy sessions, which only makes + # hygiene fire early (safe). Do NOT compensate with a threshold multiplier: 85% + # * 1.4 = 119% of context kept hygiene from ever firing for ~200K models. + + # 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 (those compress on tokens). Config: compression.hygiene_hard_message_limit. + _HARD_MSG_LIMIT = _hyg_hard_msg_limit + _needs_compress = ( + _approx_tokens >= _compress_token_threshold + or _msg_count >= _HARD_MSG_LIMIT + ) + + if _needs_compress: + # Use the persistent DB-backed cooldown (same as the in-conversation compression + # path in context_compressor.py) so the cooldown survives gateway restarts. The + # in-memory dict reset on every restart, re-triggering the same failing + # compression and wedging session storage. + _session_db = getattr(self, "_session_db", None) + if _session_db is not None: + _session_db = getattr(_session_db, "_db", _session_db) + _getter = getattr(_session_db, "get_compression_failure_cooldown", None) + if _getter is not None: + try: + _cooldown_state = _getter(session_entry.session_id) + except Exception: + _cooldown_state = None + if _cooldown_state and _cooldown_state.get("remaining_seconds", 0) > 0: + logger.info( + "Session hygiene: skipping compression for %s; " + "previous failure cooldown active for %.1fs", + session_entry.session_id, + _cooldown_state["remaining_seconds"], + ) + _needs_compress = False + + if _needs_compress and await self._session_has_compression_in_flight( + session_key + ): + # A prior hygiene/agent compression still holds the durable lock (typically a + # shielded worker left behind by /stop or /restart). Starting another attempt + # would wait up to the 600s ceiling behind a commit the fence will refuse, while + # inbound messages demote to queue. + logger.info( + "Session hygiene: skipping compression for %s; " + "another compression is already in flight", + session_entry.session_id, ) - except TranscriptReadError: - self._clear_session_env(_session_env_tokens) - return ( - "⚠️ This session's history is temporarily unavailable, so " - "this message was not processed. Ask the operator to inspect " - "state.db, then resend after it is healthy. Use /reset only " - "if you intentionally want to start a new conversation." + _needs_compress = False + + if _needs_compress: + logger.info( + "Session hygiene: %s messages, ~%s tokens (%s) — auto-compressing " + "(threshold: %s%% of %s = %s tokens)", + _msg_count, f"{_approx_tokens:,}", _token_source, + int(_hyg_threshold_pct * 100), + f"{_hyg_context_length:,}", + f"{_compress_token_threshold:,}", + ) + return _needs_compress, _approx_tokens, _msg_count, _warn_token_threshold + + async def _hmwa_hygiene_wait_for_summary(self, attempt, hs, session_entry): + """Progress-aware inline wait for the detached hygiene compressor. Returns the compressed + transcript; raises ``HygieneTurnHoldExceeded`` (turn-hold budget) or + ``asyncio.TimeoutError`` (idle/ceiling/fence cancel) for the caller's handlers.""" + _hyg_commit_fence = attempt.commit_fence + _hyg_future = attempt.future + _hyg_wait_started = attempt.wait_started + _hyg_timeout_seconds = hs.timeout_seconds + _hyg_total_ceiling_seconds = hs.total_ceiling_seconds + _hyg_max_turn_hold_seconds = hs.max_turn_hold_seconds + + # Progress-aware wait: the timeout is an INACTIVITY budget — + # the worker ticks the fence per streamed token, so a slow but + # still-generating model extends the deadline. A hard ceiling + # bounds the total so a trickle stream can't hold the turn. + while True: + if _hyg_commit_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. + _hyg_waited = ( + time.monotonic() - _hyg_wait_started + ) + _slice = min( + max( + _hyg_timeout_seconds + - _hyg_commit_fence.seconds_since_progress(), + 0.005, + ), + max( + _hyg_total_ceiling_seconds + - _hyg_waited, + 0.005, + ), + ) + # Bounded turn-hold: cap this slice at the remaining + # turn-hold budget so it is re-checked against + # _hyg_max_turn_hold_seconds at least that often — + # otherwise a continuously-streaming worker keeps the + # slice large and holds the turn until the ceiling. + _turn_hold_remaining = ( + _hyg_max_turn_hold_seconds + - (time.monotonic() - _hyg_wait_started) + ) + if _turn_hold_remaining <= 0: + # Budget exhausted: force an immediate timeout so + # the abandonment path below runs. + _slice = 0.005 + else: + _slice = min( + _slice, + max(_turn_hold_remaining, 0.005), + ) + # Re-check the fence on a short poll so a + # /stop or /restart cancel is not stuck + # behind a full idle window (#96953). + _idle_left = max( + _hyg_timeout_seconds + - _hyg_commit_fence.seconds_since_progress(), + 0.005, + ) + _slice = min(_slice, 0.25) + try: + _compressed, _ = await asyncio.wait_for( + asyncio.shield(_hyg_future), + timeout=_slice, + ) + break + except asyncio.TimeoutError: + if _hyg_commit_fence.is_cancelled: + raise + _hyg_waited = time.monotonic() - _hyg_wait_started + _idle = _hyg_commit_fence.seconds_since_progress() + # Bounded turn-hold: never hold the user's TURN past + # _hyg_max_turn_hold_seconds even if the summary + # model is still streaming; fall through to the + # timeout path, which revokes commit admission and + # proceeds on the uncompressed transcript, so the + # wire never trips 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 HygieneTurnHoldExceeded( + f"turn-hold budget {_hyg_max_turn_hold_seconds:.1f}s " + f"elapsed after {_hyg_waited:.1f}s" + ) + if hygiene_wait_should_extend( + idle=_idle, + timeout=_hyg_timeout_seconds, + waited=_hyg_waited, + ceiling=_hyg_total_ceiling_seconds, + fence_cancelled=_hyg_commit_fence.is_cancelled, + ): + if _slice >= _idle_left - 1e-9: + logger.info( + "Session hygiene compression for " + "session %s still streaming after " + "%.0fs (last progress %.1fs ago) — " + "extending wait (ceiling %.0fs)", + session_entry.session_id, + _hyg_waited, _idle, + _hyg_total_ceiling_seconds, + ) + continue + raise + return _compressed + + async def _hmwa_hygiene_on_turn_hold(self, attempt, hs, session_entry, session_key, source): + """``except HygieneTurnHoldExceeded`` body: keep or cancel the worker's commit admission, + notify the user, and re-raise; returns the compressed transcript only when the worker + was already committing.""" + _hyg_agent = attempt.agent + _hyg_meta = attempt.meta + _hyg_commit_fence = attempt.commit_fence + _hyg_future = attempt.future + _hyg_wait_started = attempt.wait_started + _hyg_max_turn_hold_seconds = hs.max_turn_hold_seconds + + # Turn-hold expiry is an availability boundary, not a + # failure: the compressor is healthy and still streaming; we + # just can't hold the turn longer. Share the safe mechanics + # (fence, release, defer, proceed uncompressed) with + # distinct provenance / user message and NO failure-cooldown + # increment. Decouple the TURN from the COMPRESSION: when + # the worker's commit is watermark-fenced (rows appended + # after compression start, this turn included, survive its + # late commit as cloned concurrent tail) the attempt KEEPS + # commit admission — the turn proceeds uncompressed NOW and + # the summary is adopted at the worker's own fenced commit + # (archive_and_compact / rotation publish); always + # cancelling burned every attempt for thinking summary + # models whose reasoning prefix alone exceeds the hold. If + # NOT watermark-fenced (no session_db, capture failed, + # legacy lock API) a late commit could clobber newer turns, + # so cancel. + _hyg_keep_admission = bool( + getattr( + _hyg_commit_fence, + "commit_watermark_fenced", + False, + ) + ) and not _hyg_commit_fence.is_cancelled + if _hyg_keep_admission: + self._defer_agent_cleanup_until_future_done( + _hyg_future, + _hyg_agent, + context="session hygiene turn-hold", + ) + attempt.cleanup_deferred = True + # NO retry-after here: the attempt is still running + # toward a real commit, and the flat 60s retry-after + # would also block the agent-side preflight compressor + # (same-session cooldown). Re-attempt spacing comes from + # the durable compression lock instead: the next turn's + # hygiene pre-check skips while this worker's lease is + # held (_session_has_compression_in_flight). The + # done-callback below records the flat retry-after ONLY + # if the worker ends without committing anything. + _hyg_deferred_sid = session_entry.session_id + _hyg_deferred_key = session_key + _hyg_deferred_agent = _hyg_agent + + def _hyg_adopt_or_space_retry( + _fut, + _gw=self, + _sid=_hyg_deferred_sid, + _skey=_hyg_deferred_key, + _agent=_hyg_deferred_agent, + ): + try: + _exc = _fut.exception() + except ( + asyncio.CancelledError, + Exception, + ): + _exc = None + _committed = False + else: + _committed = _exc is None and ( + bool( + getattr( + _agent, + "_last_compaction_in_place", + False, + ) + ) + or getattr( + _agent, "session_id", _sid + ) + != _sid + ) + if _committed: + logger.info( + "Session hygiene compression for " + "session %s finished after the " + "turn-hold was released — summary " + "adopted at the watermark-fenced " + "commit boundary (#97963)", + _sid, + ) + try: + _reset_hygiene_failure_streak( + _gw, _skey + ) + except Exception as _rs_err: + logger.debug( + "hygiene streak reset after " + "deferred adoption failed: %s", + _rs_err, + ) + else: + # Nothing to adopt (summary failed, fence + # refused the commit, or the attempt was + # superseded). Record flat spacing so sustained + # traffic does not spawn and abandon a fresh + # compressor every turn; non-escalating — the + # failure streak must not advance for a deferral. + _record_hygiene_cooldown( + _gw, _sid, + _HYGIENE_TURNHOLD_RETRY_SECONDS, + "hygiene compression deferred: " + "turn-hold budget expired and the " + "detached attempt did not commit", + ) + + _hyg_future.add_done_callback( + _hyg_adopt_or_space_retry + ) + from agent.session_activity import ( + ActivityProvenance, + ) + _stamp_hygiene_compression_provenance( + _hyg_agent, + "session hygiene compression turn-hold", + ActivityProvenance.AGENT_COMPRESSION_TURNHOLD, + "hygiene compression turn-hold " + "activity stamp failed", + ) + logger.info( + "Session hygiene compression for session %s " + "exceeded turn-hold budget (%.1fs); " + "proceeding without compression this turn — " + "the watermark-fenced worker keeps its " + "commit admission and the summary will be " + "adopted when it finishes", + session_entry.session_id, + time.monotonic() - _hyg_wait_started, + ) + _turnhold_msg = t( + "gateway.compress.turnhold_deferred" + ) + try: + _adapter = self._adapter_for_source(source) + if _adapter and source.chat_id: + await _adapter.send( + source.chat_id, + _turnhold_msg, + metadata=_hyg_meta, + ) + except Exception as _werr: + logger.warning( + "Failed to deliver compression-turnhold " + "notice to user: %s", + _werr, + ) + raise + _cancelled = None + while _cancelled is None: + if _hyg_commit_fence.commit_in_flight: + _cancelled = False + break + _cancelled = ( + _hyg_commit_fence.try_cancel_before_commit() + ) + if _cancelled is None: + await asyncio.sleep(0.025) + if not _cancelled: + # NOTE: bounded overshoot by design: the turn can be held + # past _hyg_max_turn_hold_seconds by up to the commit + # duration. Aborting mid-commit would corrupt the + # message-store transaction — the overshoot is the + # cheaper failure mode. Do NOT "fix" this into a + # mid-commit cancellation. + _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", + ) + attempt.cleanup_deferred = True + # Short, NON-escalating retry-after. Without it every + # turn re-spawns a compressor, holds it for the turn-hold + # budget and cancels it — token burn that never commits. + # Deliberately NOT _hygiene_cooldown_for_failure: the + # compressor is healthy, so the failure streak must not + # advance; only flat retry spacing is recorded. + _record_hygiene_cooldown( + self, session_entry.session_id, + _HYGIENE_TURNHOLD_RETRY_SECONDS, + "hygiene compression deferred: " + "turn-hold budget expired while the " + "summary was still streaming", + ) + 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 = t( + "gateway.compress.turnhold_deferred" + ) + try: + _adapter = self._adapter_for_source(source) + if _adapter and source.chat_id: + await _adapter.send( + source.chat_id, + _turnhold_msg, + metadata=_hyg_meta, + ) + except Exception as _werr: + logger.warning( + "Failed to deliver compression-turnhold " + "notice to user: %s", + _werr, + ) + raise + return _compressed + + async def _hmwa_hygiene_on_timeout(self, attempt, hs, session_entry, session_key, source): + """``except asyncio.TimeoutError`` body: cancel at the commit fence, record the failure + cooldown, warn the user, and re-raise; returns the compressed transcript only when the + worker crossed the commit boundary first.""" + _hyg_agent = attempt.agent + _hyg_meta = attempt.meta + _hyg_commit_fence = attempt.commit_fence + _hyg_future = attempt.future + _hyg_wait_started = attempt.wait_started + _hyg_timeout_seconds = hs.timeout_seconds + _hyg_total_ceiling_seconds = hs.total_ceiling_seconds + _hyg_failure_cooldown_seconds = hs.failure_cooldown_seconds + + _hyg_waited = time.monotonic() - _hyg_wait_started + _hyg_total_exhausted = ( + _hyg_waited >= _hyg_total_ceiling_seconds + or _hyg_commit_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. + _hyg_commit_fence.retain_compression_lock_until_worker_done() + # Capture fence state BEFORE try_cancel — that call itself + # sets is_cancelled, which would mis-label a genuine idle + # timeout as a fence cancel. + _hyg_fence_cancelled = ( + _hyg_commit_fence.is_cancelled + ) + _cancelled = None + while _cancelled is None: + # #76354 F1: a hung commit retains the fence lock; the + # lock-free phase marker keeps this loop from spinning + # forever while the commit blocks. + if _hyg_commit_fence.commit_in_flight: + _cancelled = False + break + _cancelled = ( + _hyg_commit_fence.try_cancel_before_commit() + ) + if _cancelled is None: + # Round-2 #5: transient lock-setup windows ride + # write patience for seconds; 25ms keeps sub-tick + # latency without 1kHz spin. + await asyncio.sleep(0.025) + if not _cancelled: + # The worker crossed the commit boundary just before the + # timeout; the fence poll waited for it to finish, so + # consume the result instead of treating a successful + # compaction as a timeout. + _compressed, _ = await _hyg_future + else: + # Release an inactivity-timed-out worker's holder- + # qualified lease promptly. Total-ceiling attempts + # retained it above, so this is a no-op for them. + _hyg_commit_fence.release_cancelled_compression_lock() + self._defer_agent_cleanup_until_future_done( + _hyg_future, + _hyg_agent, + context="session hygiene timeout", + ) + attempt.cleanup_deferred = True + _hyg_timeout_error = ( + "session hygiene compression " + "cancelled at commit fence" + if _hyg_fence_cancelled + else ( + "session hygiene compression " + "timed out with no output from " + "the summary model" + ) + ) + if _hyg_failure_cooldown_seconds >= 0: + _hyg_cooldown = await asyncio.to_thread( + _hygiene_cooldown_for_failure, + self, + session_key, + _hyg_failure_cooldown_seconds, + ) + _timeout_reason = ( + _hyg_timeout_error + if _hyg_fence_cancelled + else ( + "session hygiene compression total " + "ceiling exhausted" + if _hyg_total_exhausted + else "session hygiene compression " + "timed out with no output from the " + "summary model" + ) + ) + _record_hygiene_cooldown( + self, session_entry.session_id, + _hyg_cooldown, + _timeout_reason, + ) + from agent.session_activity import ( + ActivityProvenance, + ) + _stamp_hygiene_compression_provenance( + _hyg_agent, + ( + "session hygiene compression " + "cancelled at commit fence" + if _hyg_fence_cancelled + else "session hygiene compression timed out" + ), + ActivityProvenance.AGENT_COMPRESSION_TIMEOUT, + "hygiene compression timeout " + "activity stamp failed", + ) + if _hyg_fence_cancelled: + logger.warning( + "Session hygiene compression for " + "session %s was cancelled at the " + "commit fence; continuing without " + "compression", + session_entry.session_id, + ) + else: + _hyg_elapsed = ( + time.monotonic() - _hyg_wait_started + ) + if _hyg_total_exhausted: + logger.warning( + "Session hygiene compression for session %s " + "reached its total ceiling after %.1fs " + "(progress observed=%s); continuing without " + "compression", + session_entry.session_id, + _hyg_elapsed, + _hyg_commit_fence.progress_observed, + ) + else: + logger.warning( + "Session hygiene compression for session %s " + "made no progress for %.1fs (total wait " + "%.1fs, ceiling %.1fs); continuing without " + "compression", + session_entry.session_id, + _hyg_commit_fence.seconds_since_progress(), + _hyg_elapsed, + _hyg_total_ceiling_seconds, + ) + _timeout_msg = ( + _hygiene_compression_timeout_message( + total_exhausted=_hyg_total_exhausted, + elapsed=_hyg_elapsed, + idle_timeout=_hyg_timeout_seconds, + progress_observed=( + _hyg_commit_fence.progress_observed + ), + ) + ) + try: + _adapter = self._adapter_for_source(source) + if _adapter and source.chat_id: + await _adapter.send( + source.chat_id, + _timeout_msg, + metadata=_hyg_meta, + ) + except Exception as _werr: + logger.warning( + "Failed to deliver compression-timeout " + "warning to user: %s", + _werr, + ) + raise + return _compressed + + def _hmwa_hygiene_on_unwind(self, attempt, hs, session_entry, session_key): + """``except BaseException`` body (caller re-raises): revoke commit admission and record a + cooldown so the next turn does not immediately re-arm hygiene.""" + _hyg_agent = attempt.agent + _hyg_commit_fence = attempt.commit_fence + _hyg_future = attempt.future + _hyg_failure_cooldown_seconds = hs.failure_cooldown_seconds + + # #76354 F2: non-timeout unwind while the detached hygiene + # worker may still run — KeyboardInterrupt, task + # cancellation, or any unexpected error. Revoke commit + # admission (and release the worker's durable lease) BEFORE + # the host unwinds so the worker can never commit later. + _hyg_commit_fence.revoke_commit_admission() + if not attempt.cleanup_deferred: + self._defer_agent_cleanup_until_future_done( + _hyg_future, + _hyg_agent, + context="session hygiene unwind", + ) + attempt.cleanup_deferred = True + # restart drain / task cancel must record a cooldown, or the + # next turn immediately re-arms hygiene and waits up to 600s + # behind a fence that would refuse the commit again. + if _hyg_failure_cooldown_seconds >= 0: + try: + _hyg_cooldown = _hygiene_cooldown_for_failure( + self, + session_key, + _hyg_failure_cooldown_seconds, + ) + _record_hygiene_cooldown( + self, session_entry.session_id, + _hyg_cooldown, + "session hygiene compression " + "cancelled at commit fence", + ) + except Exception as _cd_err: + logger.debug( + "hygiene unwind cooldown " + "record failed: %s", + _cd_err, + ) + + async def _hmwa_hygiene_apply_result( + self, attempt, hs, _compressed, history, *, + _approx_tokens, _msg_count, _warn_token_threshold, + session_entry, session_key, source, _quick_key, run_generation, + ): + """Adopt a finished hygiene compression (rotation / in-place / refused), rebind the session + + turn lease, record streak/cooldown, and warn the user on abort. Publishes the + (possibly replaced) transcript on ``attempt.history``.""" + from agent.model_metadata import estimate_messages_tokens_rough + + _hyg_agent = attempt.agent + _hyg_meta = attempt.meta + _hyg_commit_fence = attempt.commit_fence + _hyg_failure_cooldown_seconds = hs.failure_cooldown_seconds + + # _compress_context ends the old session and creates a new + # session_id. Write compressed messages into the NEW session so + # the old transcript stays intact and 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. + _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: + logger.warning( + "Gateway hygiene compression for session %s " + "would grow transcript (~%s -> ~%s tokens); " + "keeping the original transcript unchanged", + session_entry.session_id, + f"{_hyg_in_toks:,}", + f"{_hyg_out_toks:,}", + ) + _hyg_rotated = False + _compressed = history + # Rewrite the transcript only when rotation produced a NEW + # session id. In-place compaction needs none: + # archive_and_compact() already soft-archived the previous + # active rows and inserted the compacted set, and + # rewrite_transcript() would run + # replace_messages(active_only=False) and DELETE the archived + # turns. Likewise a summary with neither rotation nor a + # completed archive_and_compact() (unchanged session_id) signals + # FAILURE; an unconditional rewrite would replace the originals + # with only the summary (permanent data loss). + # Write-before-repoint (mirrors manual /compress): if + # session_entry were repointed to the child SID and + # rewrite_transcript then failed (lock/ENOSPC), the live entry + # would reference an empty session — the conversation silently + # vanishes. Persist the child transcript first, then rebind. + if _hyg_rotated: + if not await self.async_session_store.rewrite_transcript( + _hyg_new_sid, _compressed + ): + logger.error( + "Session hygiene: failed to persist " + "compressed transcript for rotated " + "session %s → %s; keeping the live " + "entry on the original session so the " + "conversation is not dropped", + session_entry.session_id, + _hyg_new_sid, + ) + # Fail closed: treat like no rotation. + _hyg_rotated = False + _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. + self._rebind_turn_lease( + _quick_key, run_generation, _hyg_new_sid + ) + await self.async_session_store._save() + await asyncio.to_thread( + self._sync_telegram_topic_binding, + source, session_entry, + reason="hygiene-compression", + ) + + if _hyg_rotated: + # Reset stored token count — transcript rewritten + session_entry.last_prompt_tokens = 0 + attempt.history = _compressed + _new_count = len(_compressed) + _new_tokens = estimate_messages_tokens_rough( + _compressed + ) + elif _hyg_in_place: + # archive_and_compact() already persisted the + # compacted transcript inside _compress_context. + # Reset counts to match the new active set. + session_entry.last_prompt_tokens = 0 + attempt.history = _compressed + _new_count = len(_compressed) + _new_tokens = estimate_messages_tokens_rough( + _compressed + ) + else: + # No rewrite happened — the transcript is unchanged, so the + # post-compression counts equal the pre-compression ones. + _new_count = _msg_count + _new_tokens = _approx_tokens + logger.warning( + "Gateway hygiene compression for session %s " + "did not rotate or compact in place " + "(no session_db on the hygiene agent) — " + "preserving the original transcript instead " + "of overwriting it with the summary (#21301).", + session_entry.session_id, ) + logger.info( + "Session hygiene: compressed %s → %s msgs, " + "~%s → ~%s tokens", + _msg_count, _new_count, + f"{_approx_tokens:,}", f"{_new_tokens:,}", + ) + + if _new_tokens >= _warn_token_threshold: + logger.warning( + "Session hygiene: still ~%s tokens after " + "compression", + f"{_new_tokens:,}", + ) + + # 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 to start fresh. + _comp = getattr(_hyg_agent, "context_compressor", None) + _hyg_aborted = _comp is not None and getattr( + _comp, "_last_compress_aborted", False + ) + # Fence-cancelled _compress_context returns the original + # transcript with _last_compress_aborted still False + # (failure_class=commit_fence_cancelled, chunk_count=0). Treat + # that no-op as an abort so hygiene records a cooldown instead + # of retrying into the 600s wait. A successful rotate/in-place + # commit is not an abort even if a later invalidation flipped + # the fence. + _hyg_fence_cancelled = bool( + _hyg_commit_fence.is_cancelled + and not _hyg_rotated + and not _hyg_in_place + ) + if _hyg_fence_cancelled: + _hyg_aborted = True + if not _hyg_aborted: + # Recovery decision lives in the unit-tested predicate: the + # degenerate "neither rotated nor compacted in place" path + # sets both flags False and reuses the pre-compression + # counts, so a numbers-only check would read a no-op as + # success and clear the streak. + if hygiene_compaction_recovered( + aborted=_hyg_aborted, + rotated=_hyg_rotated, + in_place=_hyg_in_place, + msg_count=_msg_count, + new_count=_new_count, + approx_tokens=_approx_tokens, + new_tokens=_new_tokens, + ): + await asyncio.to_thread( + _reset_hygiene_failure_streak, + self, + session_key, + ) + if _hyg_aborted: + if _hyg_failure_cooldown_seconds >= 0: + _hyg_cooldown = await asyncio.to_thread( + _hygiene_cooldown_for_failure, + self, + session_key, + _hyg_failure_cooldown_seconds, + ) + _record_hygiene_cooldown( + self, session_entry.session_id, + _hyg_cooldown, + ( + "session hygiene compression " + "cancelled at commit fence" + if _hyg_fence_cancelled + else getattr( + _comp, "_last_summary_error", None + ) + ), + ) + from agent.session_activity import ( + ActivityProvenance, + ) + _stamp_hygiene_compression_provenance( + _hyg_agent, + "session hygiene compression aborted", + ActivityProvenance.AGENT_COMPRESSION_COOLDOWN, + "hygiene compression abort " + "activity stamp failed", + ) + 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. + from agent.redact import redact_sensitive_text + _err = redact_sensitive_text(_err, force=True) + _warn_msg = ( + "⚠️ Context compression aborted " + f"({_err}). No messages were dropped — " + "conversation is unchanged. Run /compress " + "to retry, /reset for a clean session, or " + "check your auxiliary.compression model " + "configuration." + ) + try: + _adapter = self._adapter_for_source(source) + if _adapter and source.chat_id: + await _adapter.send(source.chat_id, _warn_msg, metadata=_hyg_meta) + except Exception as _werr: + logger.warning( + "Failed to deliver compression-failure warning to user: %s", + _werr, + ) + # Separately: if the user's CONFIGURED aux model failed and we + # recovered by falling back to the main model, tell them — a + # misconfigured auxiliary.compression.model is something only + # they can fix, and silent recovery would hide it. + 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" + _aux_msg = ( + f"ℹ️ Configured compression model `{_aux_model}` " + f"failed ({_aux_err}). Recovered using your main " + "model — context is intact — but you may want to " + "check `auxiliary.compression.model` in config.yaml." + ) + try: + _adapter = self._adapter_for_source(source) + if _adapter and source.chat_id: + await _adapter.send(source.chat_id, _aux_msg, metadata=_hyg_meta) + except Exception as _werr: + logger.warning( + "Failed to deliver aux-model-fallback notice to user: %s", + _werr, + ) + + async def _hmwa_run_session_hygiene( + self, event, source, session_entry, session_key, history, _quick_key, run_generation, + ): # Session hygiene: auto-compress pathologically large transcripts before the agent starts so # oversized histories don't cause repeated truncation/context failures. Token source: the # API's prompt_tokens from the last turn (session_entry.last_prompt_tokens), else a char/4 # estimate (30-50% high on code-heavy sessions, so hygiene merely fires a bit early). - if history and len(history) >= 4: - from agent.model_metadata import ( - estimate_messages_tokens_rough, - get_model_context_length_async, + if not history or len(history) < 4: + return history + + hs = await self._hmwa_hygiene_settings(source, session_key) + if not hs.compression_enabled: + return history + _needs_compress, _approx_tokens, _msg_count, _warn_token_threshold = ( + await self._hmwa_hygiene_plan(hs, history, session_entry, session_key) + ) + if not _needs_compress: + return history + + _hyg_total_ceiling_seconds = hs.total_ceiling_seconds + _hyg_failure_cooldown_seconds = hs.failure_cooldown_seconds + _hyg_data = hs.data + + _hyg_meta = self._thread_metadata_for_source(source, self._reply_anchor_for_event(event)) + + try: + from agent.conversation_compression import CompressionCommitFence + from run_agent import AIAgent + + _hyg_model, _hyg_runtime = self._resolve_session_agent_runtime( + source=source, + session_key=session_key, + user_config=_hyg_data if isinstance(_hyg_data, dict) else None, ) - - # Read model + compression config. Hygiene threshold is intentionally HIGHER than the - # agent's own compressor (0.85 vs 0.50): it is a pre-agent safety net for sessions that - # grew between turns; at 0.50 it compressed prematurely on every turn in long sessions. - _hyg_model = "anthropic/claude-sonnet-4.6" - _hyg_threshold_pct = 0.85 - _hyg_compression_enabled = True - _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. The compressor keeps - # running detached; its commit is fenced (revoke_commit_admission) so a stale result can - # never clobber later turns. Kept well below transport idle-timeouts (Telegram ~30s). - _hyg_max_turn_hold_seconds = 10.0 - _hyg_failure_cooldown_seconds = 300.0 - _hyg_config_context_length = None - _hyg_provider = None - _hyg_base_url = None - _hyg_api_key = None - _hyg_configured_model = None - _hyg_configured_provider = None - _hyg_configured_base_url = None - _hyg_data = {} - try: - _hyg_data = _load_gateway_config() - if _hyg_data: - # Resolve model name (same logic as run_sync) - _model_cfg = _hyg_data.get("model", {}) - if isinstance(_model_cfg, str): - _hyg_model = _model_cfg - elif isinstance(_model_cfg, dict): - _hyg_model = _model_cfg.get("default") or _model_cfg.get("model") or _hyg_model - # Read explicit context_length override from model config - # (same as run_agent.py lines 995-1005) - _raw_ctx = _model_cfg.get("context_length") - if _raw_ctx is not None: - with suppress(TypeError, ValueError): - _hyg_config_context_length = int(_raw_ctx) - # Read provider for accurate context detection - _hyg_provider = _model_cfg.get("provider") or None - _hyg_base_url = _model_cfg.get("base_url") or None - - # Only the enabled flag is read; hygiene's threshold is deliberately separate - # from the agent's compression.threshold (hygiene runs higher). - _comp_cfg = _hyg_data.get("compression", {}) - if isinstance(_comp_cfg, dict): - _hyg_compression_enabled = str( - _comp_cfg.get("enabled", True) - ).lower() in {"true", "1", "yes"} - _raw_hard_limit = _comp_cfg.get("hygiene_hard_message_limit") - if _raw_hard_limit is not None: - try: - _parsed = int(_raw_hard_limit) - if _parsed > 0: - _hyg_hard_msg_limit = _parsed - except (TypeError, ValueError): - pass - _raw_timeout = _comp_cfg.get("hygiene_timeout_seconds") - if _raw_timeout is not None: - try: - _parsed = float(_raw_timeout) - if _parsed > 0: - _hyg_timeout_seconds = _parsed - except (TypeError, ValueError): - pass - _raw_ceiling = _comp_cfg.get("hygiene_total_ceiling_seconds") - if _raw_ceiling is not None: - try: - _parsed = float(_raw_ceiling) - if _parsed > 0: - _hyg_total_ceiling_seconds = _parsed - except (TypeError, ValueError): - pass - # The ceiling can never be tighter than one idle - # window, or the extension loop would be dead code. - _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: - _parsed = float(_raw_cooldown) - if _parsed >= 0: - _hyg_failure_cooldown_seconds = _parsed - except (TypeError, ValueError): - pass - - _hyg_configured_model = _hyg_model - _hyg_configured_provider = _hyg_provider - _hyg_configured_base_url = _hyg_base_url - - try: - _hyg_model, _hyg_runtime = self._resolve_session_agent_runtime( - source=source, - session_key=session_key, - user_config=_hyg_data if isinstance(_hyg_data, dict) else None, - ) - _hyg_provider = _hyg_runtime.get("provider") or _hyg_provider - _hyg_base_url = _hyg_runtime.get("base_url") or _hyg_base_url - _hyg_api_key = _hyg_runtime.get("api_key") or _hyg_api_key - except Exception: - pass - - if _hyg_config_context_length is not None: - try: - from hermes_cli.route_identity import should_clear_context_pin_async - - if await should_clear_context_pin_async( - _hyg_configured_model, - _hyg_model, - _hyg_configured_base_url, - _hyg_base_url, - _hyg_configured_provider, - _hyg_provider, - ): - _hyg_config_context_length = None - except Exception: - _hyg_config_context_length = None - - # custom_providers per-model context_length fallback (as in run_agent.py); must run - # after runtime resolution so _hyg_base_url is set. - if _hyg_config_context_length is None and _hyg_base_url: - try: - try: - from hermes_cli.config import ( - get_compatible_custom_providers as _gw_gcp, - get_custom_provider_context_length as _gw_gccl, - ) - _hyg_custom_providers = _gw_gcp(_hyg_data) - except Exception: - _hyg_custom_providers = _hyg_data.get("custom_providers") - if not isinstance(_hyg_custom_providers, list): - _hyg_custom_providers = [] - _hyg_custom_ctx = _gw_gccl( - model=_hyg_model, - base_url=_hyg_base_url, - custom_providers=_hyg_custom_providers, - ) - if _hyg_custom_ctx: - _hyg_config_context_length = int(_hyg_custom_ctx) - except (TypeError, ValueError): - pass - except Exception: - pass - - if _hyg_compression_enabled: - _hyg_context_length = await get_model_context_length_async( - _hyg_model, - base_url=_hyg_base_url or "", - api_key=_hyg_api_key or "", - config_context_length=_hyg_config_context_length, - provider=_hyg_provider or "", + _hyg_api_mode = str( + _hyg_runtime.get("api_mode") or "" + ).lower() + if _hyg_api_mode == "codex_app_server": + # codex app-server runtime: the real context is the server-side thread, + # not the transcript mirror. The detached-agent block below would only + # rewrite the mirror and its finally-eviction would destroy the live + # thread (next turn starts blank). Use the cached agent's + # thread/compact/start and KEEP it cached. + _hyg_codex_auto = "native" + _hyg_comp_cfg = ( + _hyg_data.get("compression") + if isinstance(_hyg_data, dict) + else None ) - _compress_token_threshold = int( - _hyg_context_length * _hyg_threshold_pct - ) - _warn_token_threshold = int(_hyg_context_length * 0.95) - - _msg_count = len(history) - - # Prefer actual API-reported tokens from the last turn - # (stored in session entry) over the rough char-based estimate. - _stored_tokens = session_entry.last_prompt_tokens - if _stored_tokens > 0: - _approx_tokens = _stored_tokens - _token_source = "actual" - else: - _approx_tokens = estimate_messages_tokens_rough(history) - _token_source = "estimated" - # Rough estimates run 30-50% high on code/JSON-heavy sessions, which only makes - # hygiene fire early (safe). Do NOT compensate with a threshold multiplier: 85% - # * 1.4 = 119% of context kept hygiene from ever firing for ~200K models. - - # 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 (those compress on tokens). Config: compression.hygiene_hard_message_limit. - _HARD_MSG_LIMIT = _hyg_hard_msg_limit - _needs_compress = ( - _approx_tokens >= _compress_token_threshold - or _msg_count >= _HARD_MSG_LIMIT - ) - - if _needs_compress: - # Use the persistent DB-backed cooldown (same as the in-conversation compression - # path in context_compressor.py) so the cooldown survives gateway restarts. The - # in-memory dict reset on every restart, re-triggering the same failing - # compression and wedging session storage. - _session_db = getattr(self, "_session_db", None) - if _session_db is not None: - _session_db = getattr(_session_db, "_db", _session_db) - _getter = getattr(_session_db, "get_compression_failure_cooldown", None) - if _getter is not None: - try: - _cooldown_state = _getter(session_entry.session_id) - except Exception: - _cooldown_state = None - if _cooldown_state and _cooldown_state.get("remaining_seconds", 0) > 0: - logger.info( - "Session hygiene: skipping compression for %s; " - "previous failure cooldown active for %.1fs", - session_entry.session_id, - _cooldown_state["remaining_seconds"], - ) - _needs_compress = False - - if _needs_compress and await self._session_has_compression_in_flight( - session_key - ): - # A prior hygiene/agent compression still holds the durable lock (typically a - # shielded worker left behind by /stop or /restart). Starting another attempt - # would wait up to the 600s ceiling behind a commit the fence will refuse, while - # inbound messages demote to queue. - logger.info( - "Session hygiene: skipping compression for %s; " - "another compression is already in flight", - session_entry.session_id, - ) - _needs_compress = False - - if _needs_compress: - logger.info( - "Session hygiene: %s messages, ~%s tokens (%s) — auto-compressing " - "(threshold: %s%% of %s = %s tokens)", - _msg_count, f"{_approx_tokens:,}", _token_source, - int(_hyg_threshold_pct * 100), - f"{_hyg_context_length:,}", - f"{_compress_token_threshold:,}", - ) - - _hyg_meta = self._thread_metadata_for_source(source, self._reply_anchor_for_event(event)) - - try: - from agent.conversation_compression import CompressionCommitFence - from run_agent import AIAgent - - _hyg_model, _hyg_runtime = self._resolve_session_agent_runtime( - source=source, - session_key=session_key, - user_config=_hyg_data if isinstance(_hyg_data, dict) else None, + if isinstance(_hyg_comp_cfg, dict): + _hyg_codex_auto = str( + _hyg_comp_cfg.get( + "codex_app_server_auto", "native" ) - _hyg_api_mode = str( - _hyg_runtime.get("api_mode") or "" - ).lower() - if _hyg_api_mode == "codex_app_server": - # codex app-server runtime: the real context is the server-side thread, - # not the transcript mirror. The detached-agent block below would only - # rewrite the mirror and its finally-eviction would destroy the live - # thread (next turn starts blank). Use the cached agent's - # thread/compact/start and KEEP it cached. - _hyg_codex_auto = "native" - _hyg_comp_cfg = ( - _hyg_data.get("compression") - if isinstance(_hyg_data, dict) - else None - ) - if isinstance(_hyg_comp_cfg, dict): - _hyg_codex_auto = str( - _hyg_comp_cfg.get( - "codex_app_server_auto", "native" - ) - or "native" - ) - _hyg_codex_outcome = await run_codex_hygiene_compaction( - self, - session_key, - session_entry.session_id, - auto_mode=_hyg_codex_auto, - history=history, - approx_tokens=_approx_tokens, - timeout_seconds=_hyg_total_ceiling_seconds, - failure_cooldown_seconds=_hyg_failure_cooldown_seconds, - ) - logger.info( - "Session hygiene (codex app-server): %s " - "(session=%s, mode=%s, ~%s tokens)", - _hyg_codex_outcome, - session_entry.session_id, - _hyg_codex_auto, - f"{_approx_tokens:,}", - ) - 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 so nothing compressed. - _hyg_msgs = [ - m for m in history - if m.get("role") in {"user", "assistant", "tool"} - ] + or "native" + ) + _hyg_codex_outcome = await run_codex_hygiene_compaction( + self, + session_key, + session_entry.session_id, + auto_mode=_hyg_codex_auto, + history=history, + approx_tokens=_approx_tokens, + timeout_seconds=_hyg_total_ceiling_seconds, + failure_cooldown_seconds=_hyg_failure_cooldown_seconds, + ) + logger.info( + "Session hygiene (codex app-server): %s " + "(session=%s, mode=%s, ~%s tokens)", + _hyg_codex_outcome, + session_entry.session_id, + _hyg_codex_auto, + f"{_approx_tokens:,}", + ) + 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 so nothing compressed. + _hyg_msgs = [ + m for m in history + if m.get("role") in {"user", "assistant", "tool"} + ] - if len(_hyg_msgs) >= 4: - try: - _hyg_session_row = await self._session_db.get_session( - session_entry.session_id - ) - except Exception as exc: - _hyg_session_row = None - logger.warning( - "Session hygiene could not restore the system " - "prompt for session %s: %s. Preserving an empty " - "prompt so the live turn rebuilds it with its " - "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: when - # compression.checkpoint_required is on, load the memory provider so - # the checkpoint exists before any mutation; otherwise keep the - # fast path (no provider init, no best-effort hook). - from hermes_cli.config import load_config as _load_cfg - from utils import is_truthy_value as _is_truthy - - _hyg_checkpoint_required = _is_truthy( - ((_load_cfg() or {}).get("compression") or {}).get( - "checkpoint_required" - ), - default=False, - ) - _hyg_agent = AIAgent( - **_hyg_runtime, - model=_hyg_model, - max_iterations=4, - quiet_mode=True, - skip_memory=not _hyg_checkpoint_required, - enabled_toolsets=["memory"], - 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. - _hyg_agent.platform = _GATEWAY_HYGIENE_PLATFORM - _hyg_cleanup_deferred = False - try: - # Hygiene runs before the turn and owns the session binding, so - # prefer in-place compaction: archive old rows under the same - # session id rather than minting a continuation child that must - # be published back to SessionStore/topic bindings. Without a - # SessionDB this stays False and the guard below preserves it. - _hyg_agent.compression_in_place = True - _bind_hyg_state = getattr( - getattr(_hyg_agent, "context_compressor", None), - "bind_session_state", - None, - ) - if callable(_bind_hyg_state): - _bind_hyg_state( - _hyg_session_db, - session_entry.session_id, - ) - # It must never finalize on close() — close() - # would end the live gateway session row. - _hyg_agent._end_session_on_close = False - _hyg_agent._print_fn = lambda *a, **kw: None - - loop = asyncio.get_running_loop() - _hyg_commit_fence = CompressionCommitFence( - total_ceiling_seconds=_hyg_total_ceiling_seconds - ) - # Default executor (NOT self._get_executor): a fence-cancelled - # hung summary must never occupy an agent-work slot. But it MUST - # run in the caller's contextvars: under multiplex_profiles the - # secret scope / HERMES_HOME live in ContextVars, and an empty - # Context makes get_secret() fail closed → lossy truncation. - _hyg_future = loop.run_in_executor( - None, - copy_context().run, - lambda: _hyg_agent._compress_context( - _hyg_msgs, "", - approx_tokens=_approx_tokens, - commit_fence=_hyg_commit_fence, - ), - ) - try: - # Progress-aware wait: the timeout is an INACTIVITY budget — - # the worker ticks the fence per streamed token, so a slow but - # still-generating model extends the deadline. A hard ceiling - # bounds the total so a trickle stream can't hold the turn. - _hyg_wait_started = time.monotonic() - while True: - if _hyg_commit_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. - _hyg_waited = ( - time.monotonic() - _hyg_wait_started - ) - _slice = min( - max( - _hyg_timeout_seconds - - _hyg_commit_fence.seconds_since_progress(), - 0.005, - ), - max( - _hyg_total_ceiling_seconds - - _hyg_waited, - 0.005, - ), - ) - # Bounded turn-hold: cap this slice at the remaining - # turn-hold budget so it is re-checked against - # _hyg_max_turn_hold_seconds at least that often — - # otherwise a continuously-streaming worker keeps the - # slice large and holds the turn until the ceiling. - _turn_hold_remaining = ( - _hyg_max_turn_hold_seconds - - (time.monotonic() - _hyg_wait_started) - ) - if _turn_hold_remaining <= 0: - # Budget exhausted: force an immediate timeout so - # the abandonment path below runs. - _slice = 0.005 - else: - _slice = min( - _slice, - max(_turn_hold_remaining, 0.005), - ) - # Re-check the fence on a short poll so a - # /stop or /restart cancel is not stuck - # behind a full idle window (#96953). - _idle_left = max( - _hyg_timeout_seconds - - _hyg_commit_fence.seconds_since_progress(), - 0.005, - ) - _slice = min(_slice, 0.25) - try: - _compressed, _ = await asyncio.wait_for( - asyncio.shield(_hyg_future), - timeout=_slice, - ) - break - except asyncio.TimeoutError: - if _hyg_commit_fence.is_cancelled: - raise - _hyg_waited = time.monotonic() - _hyg_wait_started - _idle = _hyg_commit_fence.seconds_since_progress() - # Bounded turn-hold: never hold the user's TURN past - # _hyg_max_turn_hold_seconds even if the summary - # model is still streaming; fall through to the - # timeout path, which revokes commit admission and - # proceeds on the uncompressed transcript, so the - # wire never trips 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 HygieneTurnHoldExceeded( - f"turn-hold budget {_hyg_max_turn_hold_seconds:.1f}s " - f"elapsed after {_hyg_waited:.1f}s" - ) - if hygiene_wait_should_extend( - idle=_idle, - timeout=_hyg_timeout_seconds, - waited=_hyg_waited, - ceiling=_hyg_total_ceiling_seconds, - fence_cancelled=_hyg_commit_fence.is_cancelled, - ): - if _slice >= _idle_left - 1e-9: - logger.info( - "Session hygiene compression for " - "session %s still streaming after " - "%.0fs (last progress %.1fs ago) — " - "extending wait (ceiling %.0fs)", - session_entry.session_id, - _hyg_waited, _idle, - _hyg_total_ceiling_seconds, - ) - continue - raise - except HygieneTurnHoldExceeded: - # Turn-hold expiry is an availability boundary, not a - # failure: the compressor is healthy and still streaming; we - # just can't hold the turn longer. Share the safe mechanics - # (fence, release, defer, proceed uncompressed) with - # distinct provenance / user message and NO failure-cooldown - # increment. Decouple the TURN from the COMPRESSION: when - # the worker's commit is watermark-fenced (rows appended - # after compression start, this turn included, survive its - # late commit as cloned concurrent tail) the attempt KEEPS - # commit admission — the turn proceeds uncompressed NOW and - # the summary is adopted at the worker's own fenced commit - # (archive_and_compact / rotation publish); always - # cancelling burned every attempt for thinking summary - # models whose reasoning prefix alone exceeds the hold. If - # NOT watermark-fenced (no session_db, capture failed, - # legacy lock API) a late commit could clobber newer turns, - # so cancel. - _hyg_keep_admission = bool( - getattr( - _hyg_commit_fence, - "commit_watermark_fenced", - False, - ) - ) and not _hyg_commit_fence.is_cancelled - if _hyg_keep_admission: - self._defer_agent_cleanup_until_future_done( - _hyg_future, - _hyg_agent, - context="session hygiene turn-hold", - ) - _hyg_cleanup_deferred = True - # NO retry-after here: the attempt is still running - # toward a real commit, and the flat 60s retry-after - # would also block the agent-side preflight compressor - # (same-session cooldown). Re-attempt spacing comes from - # the durable compression lock instead: the next turn's - # hygiene pre-check skips while this worker's lease is - # held (_session_has_compression_in_flight). The - # done-callback below records the flat retry-after ONLY - # if the worker ends without committing anything. - _hyg_deferred_sid = session_entry.session_id - _hyg_deferred_key = session_key - _hyg_deferred_agent = _hyg_agent - - def _hyg_adopt_or_space_retry( - _fut, - _gw=self, - _sid=_hyg_deferred_sid, - _skey=_hyg_deferred_key, - _agent=_hyg_deferred_agent, - ): - try: - _exc = _fut.exception() - except ( - asyncio.CancelledError, - Exception, - ): - _exc = None - _committed = False - else: - _committed = _exc is None and ( - bool( - getattr( - _agent, - "_last_compaction_in_place", - False, - ) - ) - or getattr( - _agent, "session_id", _sid - ) - != _sid - ) - if _committed: - logger.info( - "Session hygiene compression for " - "session %s finished after the " - "turn-hold was released — summary " - "adopted at the watermark-fenced " - "commit boundary (#97963)", - _sid, - ) - try: - _reset_hygiene_failure_streak( - _gw, _skey - ) - except Exception as _rs_err: - logger.debug( - "hygiene streak reset after " - "deferred adoption failed: %s", - _rs_err, - ) - else: - # Nothing to adopt (summary failed, fence - # refused the commit, or the attempt was - # superseded). Record flat spacing so sustained - # traffic does not spawn and abandon a fresh - # compressor every turn; non-escalating — the - # failure streak must not advance for a deferral. - _record_hygiene_cooldown( - _gw, _sid, - _HYGIENE_TURNHOLD_RETRY_SECONDS, - "hygiene compression deferred: " - "turn-hold budget expired and the " - "detached attempt did not commit", - ) - - _hyg_future.add_done_callback( - _hyg_adopt_or_space_retry - ) - from agent.session_activity import ( - ActivityProvenance, - ) - _stamp_hygiene_compression_provenance( - _hyg_agent, - "session hygiene compression turn-hold", - ActivityProvenance.AGENT_COMPRESSION_TURNHOLD, - "hygiene compression turn-hold " - "activity stamp failed", - ) - logger.info( - "Session hygiene compression for session %s " - "exceeded turn-hold budget (%.1fs); " - "proceeding without compression this turn — " - "the watermark-fenced worker keeps its " - "commit admission and the summary will be " - "adopted when it finishes", - session_entry.session_id, - time.monotonic() - _hyg_wait_started, - ) - _turnhold_msg = t( - "gateway.compress.turnhold_deferred" - ) - try: - _adapter = self._adapter_for_source(source) - if _adapter and source.chat_id: - await _adapter.send( - source.chat_id, - _turnhold_msg, - metadata=_hyg_meta, - ) - except Exception as _werr: - logger.warning( - "Failed to deliver compression-turnhold " - "notice to user: %s", - _werr, - ) - raise - _cancelled = None - while _cancelled is None: - if _hyg_commit_fence.commit_in_flight: - _cancelled = False - break - _cancelled = ( - _hyg_commit_fence.try_cancel_before_commit() - ) - if _cancelled is None: - await asyncio.sleep(0.025) - if not _cancelled: - # NOTE: bounded overshoot by design: the turn can be held - # past _hyg_max_turn_hold_seconds by up to the commit - # duration. Aborting mid-commit would corrupt the - # message-store transaction — the overshoot is the - # cheaper failure mode. Do NOT "fix" this into a - # mid-commit cancellation. - _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 - # Short, NON-escalating retry-after. Without it every - # turn re-spawns a compressor, holds it for the turn-hold - # budget and cancels it — token burn that never commits. - # Deliberately NOT _hygiene_cooldown_for_failure: the - # compressor is healthy, so the failure streak must not - # advance; only flat retry spacing is recorded. - _record_hygiene_cooldown( - self, session_entry.session_id, - _HYGIENE_TURNHOLD_RETRY_SECONDS, - "hygiene compression deferred: " - "turn-hold budget expired while the " - "summary was still streaming", - ) - 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 = t( - "gateway.compress.turnhold_deferred" - ) - try: - _adapter = self._adapter_for_source(source) - if _adapter and source.chat_id: - await _adapter.send( - source.chat_id, - _turnhold_msg, - metadata=_hyg_meta, - ) - except Exception as _werr: - logger.warning( - "Failed to deliver compression-turnhold " - "notice to user: %s", - _werr, - ) - raise - except asyncio.TimeoutError: - _hyg_waited = time.monotonic() - _hyg_wait_started - _hyg_total_exhausted = ( - _hyg_waited >= _hyg_total_ceiling_seconds - or _hyg_commit_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. - _hyg_commit_fence.retain_compression_lock_until_worker_done() - # Capture fence state BEFORE try_cancel — that call itself - # sets is_cancelled, which would mis-label a genuine idle - # timeout as a fence cancel. - _hyg_fence_cancelled = ( - _hyg_commit_fence.is_cancelled - ) - _cancelled = None - while _cancelled is None: - # #76354 F1: a hung commit retains the fence lock; the - # lock-free phase marker keeps this loop from spinning - # forever while the commit blocks. - if _hyg_commit_fence.commit_in_flight: - _cancelled = False - break - _cancelled = ( - _hyg_commit_fence.try_cancel_before_commit() - ) - if _cancelled is None: - # Round-2 #5: transient lock-setup windows ride - # write patience for seconds; 25ms keeps sub-tick - # latency without 1kHz spin. - await asyncio.sleep(0.025) - if not _cancelled: - # The worker crossed the commit boundary just before the - # timeout; the fence poll waited for it to finish, so - # consume the result instead of treating a successful - # compaction as a timeout. - _compressed, _ = await _hyg_future - else: - # Release an inactivity-timed-out worker's holder- - # qualified lease promptly. Total-ceiling attempts - # retained it above, so this is a no-op for them. - _hyg_commit_fence.release_cancelled_compression_lock() - self._defer_agent_cleanup_until_future_done( - _hyg_future, - _hyg_agent, - context="session hygiene timeout", - ) - _hyg_cleanup_deferred = True - _hyg_timeout_error = ( - "session hygiene compression " - "cancelled at commit fence" - if _hyg_fence_cancelled - else ( - "session hygiene compression " - "timed out with no output from " - "the summary model" - ) - ) - if _hyg_failure_cooldown_seconds >= 0: - _hyg_cooldown = await asyncio.to_thread( - _hygiene_cooldown_for_failure, - self, - session_key, - _hyg_failure_cooldown_seconds, - ) - _timeout_reason = ( - _hyg_timeout_error - if _hyg_fence_cancelled - else ( - "session hygiene compression total " - "ceiling exhausted" - if _hyg_total_exhausted - else "session hygiene compression " - "timed out with no output from the " - "summary model" - ) - ) - _record_hygiene_cooldown( - self, session_entry.session_id, - _hyg_cooldown, - _timeout_reason, - ) - from agent.session_activity import ( - ActivityProvenance, - ) - _stamp_hygiene_compression_provenance( - _hyg_agent, - ( - "session hygiene compression " - "cancelled at commit fence" - if _hyg_fence_cancelled - else "session hygiene compression timed out" - ), - ActivityProvenance.AGENT_COMPRESSION_TIMEOUT, - "hygiene compression timeout " - "activity stamp failed", - ) - if _hyg_fence_cancelled: - logger.warning( - "Session hygiene compression for " - "session %s was cancelled at the " - "commit fence; continuing without " - "compression", - session_entry.session_id, - ) - else: - _hyg_elapsed = ( - time.monotonic() - _hyg_wait_started - ) - if _hyg_total_exhausted: - logger.warning( - "Session hygiene compression for session %s " - "reached its total ceiling after %.1fs " - "(progress observed=%s); continuing without " - "compression", - session_entry.session_id, - _hyg_elapsed, - _hyg_commit_fence.progress_observed, - ) - else: - logger.warning( - "Session hygiene compression for session %s " - "made no progress for %.1fs (total wait " - "%.1fs, ceiling %.1fs); continuing without " - "compression", - session_entry.session_id, - _hyg_commit_fence.seconds_since_progress(), - _hyg_elapsed, - _hyg_total_ceiling_seconds, - ) - _timeout_msg = ( - _hygiene_compression_timeout_message( - total_exhausted=_hyg_total_exhausted, - elapsed=_hyg_elapsed, - idle_timeout=_hyg_timeout_seconds, - progress_observed=( - _hyg_commit_fence.progress_observed - ), - ) - ) - try: - _adapter = self._adapter_for_source(source) - if _adapter and source.chat_id: - await _adapter.send( - source.chat_id, - _timeout_msg, - metadata=_hyg_meta, - ) - except Exception as _werr: - logger.warning( - "Failed to deliver compression-timeout " - "warning to user: %s", - _werr, - ) - raise - except BaseException: - # #76354 F2: non-timeout unwind while the detached hygiene - # worker may still run — KeyboardInterrupt, task - # cancellation, or any unexpected error. Revoke commit - # admission (and release the worker's durable lease) BEFORE - # the host unwinds so the worker can never commit later. - _hyg_commit_fence.revoke_commit_admission() - if not _hyg_cleanup_deferred: - self._defer_agent_cleanup_until_future_done( - _hyg_future, - _hyg_agent, - context="session hygiene unwind", - ) - _hyg_cleanup_deferred = True - # restart drain / task cancel must record a cooldown, or the - # next turn immediately re-arms hygiene and waits up to 600s - # behind a fence that would refuse the commit again. - if _hyg_failure_cooldown_seconds >= 0: - try: - _hyg_cooldown = _hygiene_cooldown_for_failure( - self, - session_key, - _hyg_failure_cooldown_seconds, - ) - _record_hygiene_cooldown( - self, session_entry.session_id, - _hyg_cooldown, - "session hygiene compression " - "cancelled at commit fence", - ) - except Exception as _cd_err: - logger.debug( - "hygiene unwind cooldown " - "record failed: %s", - _cd_err, - ) - raise - - # _compress_context ends the old session and creates a new - # session_id. Write compressed messages into the NEW session so - # the old transcript stays intact and 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. - _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: - logger.warning( - "Gateway hygiene compression for session %s " - "would grow transcript (~%s -> ~%s tokens); " - "keeping the original transcript unchanged", - session_entry.session_id, - f"{_hyg_in_toks:,}", - f"{_hyg_out_toks:,}", - ) - _hyg_rotated = False - _compressed = history - # Rewrite the transcript only when rotation produced a NEW - # session id. In-place compaction needs none: - # archive_and_compact() already soft-archived the previous - # active rows and inserted the compacted set, and - # rewrite_transcript() would run - # replace_messages(active_only=False) and DELETE the archived - # turns. Likewise a summary with neither rotation nor a - # completed archive_and_compact() (unchanged session_id) signals - # FAILURE; an unconditional rewrite would replace the originals - # with only the summary (permanent data loss). - # Write-before-repoint (mirrors manual /compress): if - # session_entry were repointed to the child SID and - # rewrite_transcript then failed (lock/ENOSPC), the live entry - # would reference an empty session — the conversation silently - # vanishes. Persist the child transcript first, then rebind. - if _hyg_rotated: - if not await self.async_session_store.rewrite_transcript( - _hyg_new_sid, _compressed - ): - logger.error( - "Session hygiene: failed to persist " - "compressed transcript for rotated " - "session %s → %s; keeping the live " - "entry on the original session so the " - "conversation is not dropped", - session_entry.session_id, - _hyg_new_sid, - ) - # Fail closed: treat like no rotation. - _hyg_rotated = False - _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. - self._rebind_turn_lease( - _quick_key, run_generation, _hyg_new_sid - ) - await self.async_session_store._save() - await asyncio.to_thread( - self._sync_telegram_topic_binding, - source, session_entry, - reason="hygiene-compression", - ) - - if _hyg_rotated: - # Reset stored token count — transcript rewritten - session_entry.last_prompt_tokens = 0 - history = _compressed - _new_count = len(_compressed) - _new_tokens = estimate_messages_tokens_rough( - _compressed - ) - elif _hyg_in_place: - # archive_and_compact() already persisted the - # compacted transcript inside _compress_context. - # Reset counts to match the new active set. - session_entry.last_prompt_tokens = 0 - history = _compressed - _new_count = len(_compressed) - _new_tokens = estimate_messages_tokens_rough( - _compressed - ) - else: - # No rewrite happened — the transcript is unchanged, so the - # post-compression counts equal the pre-compression ones. - _new_count = _msg_count - _new_tokens = _approx_tokens - logger.warning( - "Gateway hygiene compression for session %s " - "did not rotate or compact in place " - "(no session_db on the hygiene agent) — " - "preserving the original transcript instead " - "of overwriting it with the summary (#21301).", - session_entry.session_id, - ) - - logger.info( - "Session hygiene: compressed %s → %s msgs, " - "~%s → ~%s tokens", - _msg_count, _new_count, - f"{_approx_tokens:,}", f"{_new_tokens:,}", - ) - - if _new_tokens >= _warn_token_threshold: - logger.warning( - "Session hygiene: still ~%s tokens after " - "compression", - f"{_new_tokens:,}", - ) - - # 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 to start fresh. - _comp = getattr(_hyg_agent, "context_compressor", None) - _hyg_aborted = _comp is not None and getattr( - _comp, "_last_compress_aborted", False - ) - # Fence-cancelled _compress_context returns the original - # transcript with _last_compress_aborted still False - # (failure_class=commit_fence_cancelled, chunk_count=0). Treat - # that no-op as an abort so hygiene records a cooldown instead - # of retrying into the 600s wait. A successful rotate/in-place - # commit is not an abort even if a later invalidation flipped - # the fence. - _hyg_fence_cancelled = bool( - _hyg_commit_fence.is_cancelled - and not _hyg_rotated - and not _hyg_in_place - ) - if _hyg_fence_cancelled: - _hyg_aborted = True - if not _hyg_aborted: - # Recovery decision lives in the unit-tested predicate: the - # degenerate "neither rotated nor compacted in place" path - # sets both flags False and reuses the pre-compression - # counts, so a numbers-only check would read a no-op as - # success and clear the streak. - if hygiene_compaction_recovered( - aborted=_hyg_aborted, - rotated=_hyg_rotated, - in_place=_hyg_in_place, - msg_count=_msg_count, - new_count=_new_count, - approx_tokens=_approx_tokens, - new_tokens=_new_tokens, - ): - await asyncio.to_thread( - _reset_hygiene_failure_streak, - self, - session_key, - ) - if _hyg_aborted: - if _hyg_failure_cooldown_seconds >= 0: - _hyg_cooldown = await asyncio.to_thread( - _hygiene_cooldown_for_failure, - self, - session_key, - _hyg_failure_cooldown_seconds, - ) - _record_hygiene_cooldown( - self, session_entry.session_id, - _hyg_cooldown, - ( - "session hygiene compression " - "cancelled at commit fence" - if _hyg_fence_cancelled - else getattr( - _comp, "_last_summary_error", None - ) - ), - ) - from agent.session_activity import ( - ActivityProvenance, - ) - _stamp_hygiene_compression_provenance( - _hyg_agent, - "session hygiene compression aborted", - ActivityProvenance.AGENT_COMPRESSION_COOLDOWN, - "hygiene compression abort " - "activity stamp failed", - ) - 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. - from agent.redact import redact_sensitive_text - _err = redact_sensitive_text(_err, force=True) - _warn_msg = ( - "⚠️ Context compression aborted " - f"({_err}). No messages were dropped — " - "conversation is unchanged. Run /compress " - "to retry, /reset for a clean session, or " - "check your auxiliary.compression model " - "configuration." - ) - try: - _adapter = self._adapter_for_source(source) - if _adapter and source.chat_id: - await _adapter.send(source.chat_id, _warn_msg, metadata=_hyg_meta) - except Exception as _werr: - logger.warning( - "Failed to deliver compression-failure warning to user: %s", - _werr, - ) - # Separately: if the user's CONFIGURED aux model failed and we - # recovered by falling back to the main model, tell them — a - # misconfigured auxiliary.compression.model is something only - # they can fix, and silent recovery would hide it. - 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" - _aux_msg = ( - f"ℹ️ Configured compression model `{_aux_model}` " - f"failed ({_aux_err}). Recovered using your main " - "model — context is intact — but you may want to " - "check `auxiliary.compression.model` in config.yaml." - ) - try: - _adapter = self._adapter_for_source(source) - if _adapter and source.chat_id: - await _adapter.send(source.chat_id, _aux_msg, metadata=_hyg_meta) - except Exception as _werr: - logger.warning( - "Failed to deliver aux-model-fallback notice to user: %s", - _werr, - ) - finally: - # Evict the cached agent so the next turn rebuilds its system - # prompt from current SOUL.md, memory, and skills. - self._evict_cached_agent(session_key) - if not _hyg_cleanup_deferred: - await self._cleanup_agent_resources_off_loop( - _hyg_agent, context="session hygiene" - ) - - except HygieneTurnHoldExceeded: - # Availability boundary, not a failure — already logged at INFO by the turn- - # hold handler. Must not hit the generic "auto-compress failed" warning - # below: that log made thinking-model deployments read as permanently broken. - pass - except Exception as e: + if len(_hyg_msgs) >= 4: + try: + _hyg_session_row = await self._session_db.get_session( + session_entry.session_id + ) + except Exception as exc: + _hyg_session_row = None logger.warning( - "Session hygiene auto-compress failed: %s", e + "Session hygiene could not restore the system " + "prompt for session %s: %s. Preserving an empty " + "prompt so the live turn rebuilds it with its " + "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: when + # compression.checkpoint_required is on, load the memory provider so + # the checkpoint exists before any mutation; otherwise keep the + # fast path (no provider init, no best-effort hook). + from hermes_cli.config import load_config as _load_cfg + from utils import is_truthy_value as _is_truthy + _hyg_checkpoint_required = _is_truthy( + ((_load_cfg() or {}).get("compression") or {}).get( + "checkpoint_required" + ), + default=False, + ) + _hyg_agent = AIAgent( + **_hyg_runtime, + model=_hyg_model, + max_iterations=4, + quiet_mode=True, + skip_memory=not _hyg_checkpoint_required, + enabled_toolsets=["memory"], + 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. + _hyg_agent.platform = _GATEWAY_HYGIENE_PLATFORM + attempt = self._HygieneAttempt(agent=_hyg_agent, meta=_hyg_meta, history=history) + try: + # Hygiene runs before the turn and owns the session binding, so + # prefer in-place compaction: archive old rows under the same + # session id rather than minting a continuation child that must + # be published back to SessionStore/topic bindings. Without a + # SessionDB this stays False and the guard below preserves it. + _hyg_agent.compression_in_place = True + _bind_hyg_state = getattr( + getattr(_hyg_agent, "context_compressor", None), + "bind_session_state", + None, + ) + if callable(_bind_hyg_state): + _bind_hyg_state( + _hyg_session_db, + session_entry.session_id, + ) + # It must never finalize on close() — close() + # would end the live gateway session row. + _hyg_agent._end_session_on_close = False + _hyg_agent._print_fn = lambda *a, **kw: None + + loop = asyncio.get_running_loop() + _hyg_commit_fence = CompressionCommitFence( + total_ceiling_seconds=_hyg_total_ceiling_seconds + ) + # Default executor (NOT self._get_executor): a fence-cancelled + # hung summary must never occupy an agent-work slot. But it MUST + # run in the caller's contextvars: under multiplex_profiles the + # secret scope / HERMES_HOME live in ContextVars, and an empty + # Context makes get_secret() fail closed → lossy truncation. + _hyg_future = loop.run_in_executor( + None, + copy_context().run, + lambda: _hyg_agent._compress_context( + _hyg_msgs, "", + approx_tokens=_approx_tokens, + commit_fence=_hyg_commit_fence, + ), + ) + attempt.commit_fence = _hyg_commit_fence + attempt.future = _hyg_future + attempt.wait_started = time.monotonic() + try: + _compressed = await self._hmwa_hygiene_wait_for_summary( + attempt, hs, session_entry, + ) + except HygieneTurnHoldExceeded: + _compressed = await self._hmwa_hygiene_on_turn_hold( + attempt, hs, session_entry, session_key, source, + ) + except asyncio.TimeoutError: + _compressed = await self._hmwa_hygiene_on_timeout( + attempt, hs, session_entry, session_key, source, + ) + except BaseException: + self._hmwa_hygiene_on_unwind(attempt, hs, session_entry, session_key) + raise + + await self._hmwa_hygiene_apply_result( + attempt, hs, _compressed, history, + _approx_tokens=_approx_tokens, + _msg_count=_msg_count, + _warn_token_threshold=_warn_token_threshold, + session_entry=session_entry, + session_key=session_key, + source=source, + _quick_key=_quick_key, + run_generation=run_generation, + ) + finally: + history = attempt.history + # Evict the cached agent so the next turn rebuilds its system + # prompt from current SOUL.md, memory, and skills. + self._evict_cached_agent(session_key) + if not attempt.cleanup_deferred: + await self._cleanup_agent_resources_off_loop( + _hyg_agent, context="session hygiene" + ) + except HygieneTurnHoldExceeded: + # Availability boundary, not a failure — already logged at INFO by the turn- + # hold handler. Must not hit the generic "auto-compress failed" warning + # below: that log made thinking-model deployments read as permanently broken. + pass + except Exception as e: + logger.warning( + "Session hygiene auto-compress failed: %s", e + ) + return history + + async def _hmwa_first_contact_notes(self, source, history, turn_sidecar_notes): + """First-ever-message onboarding note + one-time 'no home channel' prompt (both only when + the session has no history).""" # First-message onboarding -- only on the very first interaction ever. Delivered on the # current user message (sidecar), NOT the ephemeral system prompt: present-on-turn-1/absent- # on-turn-2 was a guaranteed system-prompt diff and agent rebuild. @@ -19805,30 +19918,13 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew ) await self._deliver_platform_notice(source, notice) - # Voice channel awareness: deliver voice channel state (who is present / speaking) on the - # user message, ONLY when changed since the previous turn. It differs almost every turn, and - # in the ephemeral system prompt it forced a full agent rebuild + prompt-cache re-key per - # message; the system prompt carries a static pointer line instead (gateway/session.py). - _vc_note = self._voice_channel_sidecar_note(event, source, session_key) - if _vc_note: - turn_sidecar_notes.append(_vc_note) - - # Auto-analyze user images: run the vision tool eagerly so the model always gets a text - # description plus the local path for re-examination via vision_analyze. Filter to image - # media_type so documents/audio in the same message are not sent to the vision tool. - message_text = await self._prepare_profile_scoped_inbound_message_text( - event=event, - source=source, - history=history, - session_key=session_key, - ) - if message_text is None: - return - + def _hmwa_apply_message_timestamp(self, event, message_text): # Capture the platform event time as message metadata and keep the persisted transcript # clean (strip any leading timestamp prefix). This runs regardless of the toggle so storage # stays clean and the send-time is preserved. Only the in-context RENDER (the prefix the # model sees) is gated behind gateway.message_timestamps.enabled — default OFF. + persist_user_message = None + persist_user_timestamp = None try: from hermes_time import get_timezone as _get_evt_tz from gateway.message_timestamps import ( @@ -19858,6 +19954,781 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew message_text = _clean_message_text except Exception as _ts_err: logger.debug("Message timestamp injection failed (non-fatal): %s", _ts_err) + return message_text, persist_user_message, persist_user_timestamp + + async def _hmwa_stop_typing_for_turn(self, event, source): + """Stop the typing indicator (never raises). Slack AI status is scoped to a thread/ + workspace, so preserve the routing metadata used by the response delivery path.""" + try: + _typing_adapter = self._adapter_for_source(source) + _stop_with_metadata = getattr( + type(_typing_adapter), "_stop_typing_with_metadata", None + ) + _stop_typing = getattr(type(_typing_adapter), "stop_typing", None) + if _typing_adapter and callable(_stop_with_metadata): + await _typing_adapter._stop_typing_with_metadata( + source.chat_id, + self._thread_metadata_for_source( + source, self._reply_anchor_for_event(event) + ), + ) + elif _typing_adapter and callable(_stop_typing): + await _typing_adapter.stop_typing(source.chat_id) + except Exception: + pass + + async def _hmwa_shape_agent_response( + self, agent_result, source, history, session_entry, session_key, + _quick_key, run_generation, _run_start_session_id, _platform_name, _msg_start_time, + ): + """Turn the raw agent result into the outbound text: sentinel/silence handling, response + logging, resume-pending clear, empty-response normalization, and identity-guarded + post-compression session_id propagation. Returns + ``(response, _intentional_silence, agent_messages)``.""" + response = agent_result.get("final_response") or "" + # Hidden-reasoning-only retry exhaustion: the loop's sentinel text ("Codex response + # remained incomplete after 3 continuation attempts") doubles as final_response, so it + # would be delivered verbatim into the channel — where peer agents can ingest it as a + # completed assistant turn. + if _is_gateway_hidden_reasoning_incomplete_turn(agent_result): + response = "" + try: + from gateway.response_filters import is_intentional_silence_agent_result + _intentional_silence = is_intentional_silence_agent_result( + agent_result, response, + ) + except Exception: + _intentional_silence = False + + # Convert the agent's internal "(empty)" sentinel into a user-friendly message. + # "(empty)" means the model failed to produce visible content after exhausting all + # retries (nudge, prefill, empty-retry, fallback). + if response == "(empty)" and not _intentional_silence: + response = ( + "⚠️ The model returned no response after processing tool " + "results. This can happen with some models — try again or " + "rephrase your question." + ) + agent_messages = agent_result.get("messages", []) + _response_time = time.time() - _msg_start_time + _api_calls = agent_result.get("api_calls", 0) + _resp_len = len(response) + logger.info( + "response ready: platform=%s chat=%s time=%.1fs api_calls=%d response=%d chars", + _platform_name, source.chat_id or "unknown", + _response_time, _api_calls, _resp_len, + ) + + # The cross-process cache-coherence re-baseline (_refresh_agent_cache_message_count) is + # deferred until AFTER the transcript persistence block below: it must include the + # first-turn `session_meta` marker row and the compression session_id swap. + + # Successful turn: clear the stuck-loop counter (it only accumulates 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. + if session_key and _should_clear_resume_pending_after_turn(agent_result): + await self._clear_restart_failure_count(session_key) + try: + await self.async_session_store.clear_resume_pending(session_key) + except Exception as _e: + logger.debug( + "clear_resume_pending failed for %s: %s", + session_key, _e, + ) + + # Normalize empty responses: surface errors, partial failures, and + # the case where agent did work but returned no text. Fix for #18765. + if not _intentional_silence: + 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(). + 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: the transcript persistence below + # writes to the NEW id, so the serialization boundary must move with it or an + # alias key resolving the fresh child could interleave. + 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( + session_entry.session_id, + session_key, + source, + ) + await asyncio.to_thread( + self._sync_telegram_topic_binding, + source, session_entry, reason="agent-result-compression", + ) + else: + logger.info( + "Skipping agent-result session split sync for %s because " + "the session binding moved from %s to %s before " + "compression finished", + session_key or "?", + _run_start_session_id, + session_entry.session_id, + ) + return response, _intentional_silence, agent_messages + + def _hmwa_prepend_reasoning(self, agent_result, response, source, _intentional_silence): + # Prepend reasoning if display is enabled (per-platform). Mattermost requires explicit + # opt-in because this is scratch text, not ordinary final-answer content. + try: + _show_reasoning_effective = _resolve_gateway_display_bool( + _load_gateway_config(), + _platform_config_key(source.platform), + "show_reasoning", + default=bool(getattr(self, "_show_reasoning", False)), + platform=source.platform, + require_platform_override_for={Platform.MATTERMOST}, + ) + except Exception: + _show_reasoning_effective = ( + False + if source.platform == Platform.MATTERMOST + else getattr(self, "_show_reasoning", False) + ) + if _show_reasoning_effective and response and not _intentional_silence: + last_reasoning = agent_result.get("last_reasoning") + if last_reasoning: + from gateway.stream_consumer import escape_code_fences_for_display + # Collapse long reasoning to keep messages readable + lines = last_reasoning.strip().splitlines() + if len(lines) > 15: + display_reasoning = "\n".join(lines[:15]) + display_reasoning += 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. + try: + from gateway.display_config import resolve_display_setting + _reasoning_style = resolve_display_setting( + _load_gateway_config(), + _platform_config_key(source.platform), + "reasoning_style", + "code", + ) + except Exception: + _reasoning_style = "code" + if _reasoning_style == "subtext": + _quoted = "\n".join( + f"-# {ln}" if ln else "-#" for ln in display_reasoning.splitlines() + ) + response = f"-# 💭 Reasoning\n{_quoted}\n\n{response}" + elif _reasoning_style == "blockquote": + _quoted = "\n".join( + f"> {ln}" if ln else ">" for ln in display_reasoning.splitlines() + ) + response = f"> 💭 **Reasoning:**\n{_quoted}\n\n{response}" + else: + # Escape ``` inside reasoning so inner fences don't + # break the outer code block used to render it. + display_reasoning = escape_code_fences_for_display(display_reasoning) + response = f"💭 **Reasoning:**\n```\n{display_reasoning}\n```\n\n{response}" + return response + + def _hmwa_runtime_footer_line(self, agent_result, source, _turn_seconds): + # Runtime-metadata footer — only on the FINAL message of the turn. Off by default + # (display.runtime_footer.enabled=false). When streaming already delivered the body, we + # can't mutate the sent text, so we fire a separate trailing send below. + _footer_line = "" + try: + from gateway.runtime_footer import build_footer_line as _bfl + _footer_line = _bfl( + user_config=_load_gateway_config(), + platform_key=_platform_config_key(source.platform), + model=agent_result.get("model"), + context_tokens=agent_result.get("last_prompt_tokens", 0) or 0, + context_length=agent_result.get("context_length") or None, + cwd=_terminal_scope_cwd(""), + turn_seconds=_turn_seconds, + ) + except Exception as _footer_err: + logger.debug("runtime_footer build failed: %s", _footer_err) + _footer_line = "" + return _footer_line + + async def _hmwa_post_turn_hooks(self, hook_ctx, agent_result, response): + """agent:end hook, process-watcher scheduling, and watch-notification drain.""" + # Emit agent:end hook + await self.hooks.emit("agent:end", { + **hook_ctx, + "response": (response or "")[:500], + "model": agent_result.get("model", ""), + "provider": agent_result.get("provider", ""), + }) + + # Check for pending process watchers (check_interval on background processes) + try: + from tools.process_registry import process_registry + # Detach the current batch atomically (see crash-recovery drain + # above): reassign to a fresh list so a watcher appended by a + # concurrent session during the yield isn't dropped by clear(). + watchers = process_registry.pending_watchers + process_registry.pending_watchers = [] + for i, watcher in enumerate(watchers): + asyncio.create_task(self._run_process_watcher(watcher)) + if i % 100 == 99: + await asyncio.sleep(0) + 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 (handled by the per-process watcher task above) and async-delegation + # completions (owned by _async_delegation_watcher, the single consumer for idle and + # post-turn cases) — inject only watch-type events and leave the rest on the queue. + try: + from tools.process_registry import process_registry as _pr + await self._drain_watch_notifications(_pr.completion_queue) + except Exception as e: + logger.debug("Watch queue drain error: %s", e) + + def _hmwa_classify_turn_failure(self, agent_result, history, session_entry): + """Classify a finished turn for transcript persistence. Returns + ``(agent_failed_early, hidden_reasoning_incomplete, is_context_overflow_failure)``.""" + # Persist the full agent loop (tool calls, results, reasoning) so sessions resume with + # full context. IMPORTANT: on context-overflow failures (compression exhausted, generic + # 400 on large sessions) do NOT persist the user message — it would grow the session and + # reproduce the failure forever. Transient failures (429, timeout, connection error, + # 5xx) are different: the session is not oversized and dropping the user turn causes + # severe context loss on retry, so persist it. + agent_failed_early = bool(agent_result.get("failed")) + hidden_reasoning_incomplete = _is_gateway_hidden_reasoning_incomplete_turn( + agent_result + ) + _err_str_for_classify = str(agent_result.get("error", "")).lower() + # Use specific multi-word phrases (not bare "exceed"/"token") to avoid false positives + # on transient errors such as "rate limit exceeded"; matches run_agent.py's classifier. + is_context_overflow_failure = agent_failed_early and ( + bool(agent_result.get("compression_exhausted")) + or any(p in _err_str_for_classify for p in ( + "context length", "context size", "context window", + "maximum context", "token limit", "too many tokens", + "reduce the length", "exceeds the limit", + "request entity too large", "prompt is too long", + "payload too large", "input is too long", + )) + or ("400" in _err_str_for_classify and len(history) > 50) + ) + if is_context_overflow_failure: + logger.info( + "Skipping transcript persistence for context-overflow " + "failure in session %s to prevent session growth loop.", + session_entry.session_id, + ) + elif agent_failed_early: + logger.info( + "Transient agent failure in session %s — persisting user " + "message so conversation context is preserved on retry.", + session_entry.session_id, + ) + elif hidden_reasoning_incomplete: + logger.warning( + "Suppressing hidden-reasoning-only incomplete gateway turn " + "for session %s: %s", + session_entry.session_id, + agent_result.get("error", "processing incomplete"), + ) + return agent_failed_early, hidden_reasoning_incomplete, is_context_overflow_failure + + async def _hmwa_compression_exhaustion_reset( + self, agent_result, response, session_entry, session_key, source, + ): + """Auto-reset a permanently oversized session (never on a lock-contended defer). Returns + ``(response, session_entry)``.""" + # Compression exhausted = permanently too large: auto-reset so the next message starts + # fresh instead of replaying the oversized context forever. A lock-contended defer is + # the OPPOSITE case (a concurrent path holds the lock and is shrinking it): never wipe + # for that. + if agent_result.get("compression_deferred"): + logger.info( + "Compression deferred for session %s — the compression " + "lock is held by a concurrent compressor. Keeping the " + "session intact; the next message retries normally.", + session_entry.session_id if session_entry else "?", + ) + elif agent_result.get("compression_exhausted") and session_entry and session_key: + 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). + 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: compression rotated + # session_entry.session_id to the bloated child earlier this turn and that _sync + # also rewrote the (chat_id, thread_id) binding. Without a re-sync the + # binding-heal walk switches the next inbound message back onto the child and + # re-triggers exhaustion forever. No-op on non-topic lanes. + session_entry = new_entry + await asyncio.to_thread( + self._sync_telegram_topic_binding, + source, session_entry, reason="compression-exhausted-reset", + ) + response = (response or "") + ( + "\n\n🔄 Session auto-reset — the conversation exceeded the " + "maximum context size and could not be compressed further. " + "Your next message will start a fresh session." + ) + return response, session_entry + + async def _hmwa_persist_turn_transcript( + self, *, event, source, session_entry, session_key, agent_result, agent_messages, + history, response, message_text, persist_user_message, persist_user_timestamp, + persist_user_display_kind, agent_failed_early, hidden_reasoning_incomplete, + is_context_overflow_failure, + ): + """Persist this turn to the transcript (session_meta on first turn, user-only on transient + failure, nothing on context overflow), update last_prompt_tokens, and re-baseline the + cached agent's message count.""" + ts = time.time() # Unix epoch float — consistent with DB storage + + # Fresh session (no history): 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. + if is_context_overflow_failure: + pass # Skip all transcript writes — don't grow a broken session + elif not history: + tool_defs = agent_result.get("tools", []) + await self.async_session_store.append_to_transcript( + session_entry.session_id, + { + "role": "session_meta", + "tools": tool_defs or [], + "model": _resolve_gateway_model(), + "platform": source.platform.value if source.platform else "", + "timestamp": ts, + } + ) + + # The agent already persisted these via _flush_messages_to_session_db(); skip the DB + # write to avoid duplicates. Holds for the codex app-server runtime too (it flushes its + # own projected messages before returning and reports agent_persisted=True). Reading the + # flag (default = self._session_db is not None) keeps the contract explicit; a + # non-persisting runtime opts in via False. + agent_persisted = agent_result.get("agent_persisted", self._session_db is not None) + + # Only the NEW messages from this turn: use history_offset (what the agent saw), not + # len(history), which counts session_meta entries stripped before the agent saw them. + if is_context_overflow_failure: + pass # handled above — skip all transcript writes + elif agent_failed_early or hidden_reasoning_incomplete: + # Transient failure (429/timeout/5xx): persist only the user message so the next + # message can load a transcript that reflects what was said. Skip the assistant + # error text since it's a gateway-generated hint, not model output. Hidden-reasoning + # incomplete turns follow the same rule so peer-agent channels don't ingest them. + _user_entry = { + "role": "user", + "content": ( + persist_user_message + if persist_user_message is not None + else message_text + ), + "timestamp": ( + persist_user_timestamp + if persist_user_timestamp is not None + else ts + ), + } + if persist_user_display_kind: + _user_entry["display_kind"] = persist_user_display_kind + if event.message_id: + _user_entry["message_id"] = str(event.message_id) + # Dedupe: skip if this platform message_id is already in the transcript (prevents + # duplicate user turns on Telegram retries after transient failures). + _skip_persist = ( + event.message_id + and await self.async_session_store.has_platform_message_id( + session_entry.session_id, str(event.message_id) + ) + ) + if _skip_persist: + logger.info( + "Skipping duplicate user turn " + "(message_id=%s) in session %s", + event.message_id, session_entry.session_id, + ) + else: + await self.async_session_store.append_to_transcript( + session_entry.session_id, + _user_entry, + skip_db=agent_persisted, + ) + else: + history_len = agent_result.get("history_offset", len(history)) + new_messages = agent_messages[history_len:] if len(agent_messages) > history_len else [] + + # If no new messages found (edge case), fall back to simple user/assistant + if not new_messages: + _user_entry = { + "role": "user", + "content": ( + persist_user_message + if persist_user_message is not None + else message_text + ), + "timestamp": ( + persist_user_timestamp + if persist_user_timestamp is not None + else ts + ), + } + if persist_user_display_kind: + _user_entry["display_kind"] = persist_user_display_kind + if event.message_id: + _user_entry["message_id"] = str(event.message_id) + await self.async_session_store.append_to_transcript( + session_entry.session_id, + _user_entry, + skip_db=agent_persisted, + ) + if response: + await self.async_session_store.append_to_transcript( + session_entry.session_id, + {"role": "assistant", "content": response, "timestamp": ts}, + 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. + _user_msg_id_attached = False + for msg in new_messages: + # Skip system messages (they're rebuilt each run) + if msg.get("role") == "system": + continue + # Add timestamp to each message for debugging + entry = {**msg, "timestamp": ts} + if ( + not _user_msg_id_attached + and msg.get("role") == "user" + and event.message_id + and "message_id" not in entry + ): + entry["message_id"] = str(event.message_id) + _user_msg_id_attached = True + await self.async_session_store.append_to_transcript( + session_entry.session_id, 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. + await self.async_session_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)), + ) + + # Re-baseline the cached agent's message_count snapshot now that ALL of this turn's + # transcript writes are done (flushed rows AND the first-turn `session_meta` marker). + # The cross-process coherence guard snapshots at agent-BUILD time and never refreshes + # on reuse, so our own writes would trigger a rebuild next turn (destroying prompt + # caching). MUST run after the session_meta append (that row bumps the count too). + await self._refresh_agent_cache_message_count( + session_key, session_entry.session_id + ) + + async def _hmwa_deliver_turn_response( + self, event, source, session_entry, session_key, run_generation, + agent_result, agent_messages, response, _footer_line, _intentional_silence, + ): + """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; only + # the outbound delivery is suppressed. + if _intentional_silence: + logger.info( + "Suppressing intentional silence marker for session %s", + session_entry.session_id, + ) + response = "" + + # Auto voice reply: send TTS audio before the text response + _already_sent = bool(agent_result.get("already_sent")) + # Skip when streaming TTS already delivered audio for this turn (#60671). + _stts_adapter = self._adapter_for_source(source) + _streaming_tts_done = ( + _stts_adapter is not None + and bool(getattr(_stts_adapter, "_streaming_tts_turn_completed", lambda *_a, **_k: False)(session_key, run_generation)) + ) + if ( + not _streaming_tts_done + and self._should_send_voice_reply(event, response, agent_messages, already_sent=_already_sent) + ): + 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. + if agent_result.get("already_sent") and not agent_result.get("failed"): + if response: + _media_adapter = self._adapter_for_source(source) + if _media_adapter: + await self._deliver_media_from_response( + response, event, _media_adapter, + ) + # Streaming already delivered the body text, but the footer was intentionally held + # back (see the `not already_sent` gate above). + if _footer_line: + try: + _foot_adapter = self._adapter_for_source(source) + if _foot_adapter: + await _foot_adapter.send( + source.chat_id, + _footer_line, + metadata=self._thread_metadata_for_source(source, self._reply_anchor_for_event(event)), + ) + except Exception as _e: + logger.debug("trailing footer send failed: %s", _e) + # This branch returns 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. + with suppress(Exception): + event._streamed_final_response = str(response or "") + return None + + return response + + async def _hmwa_agent_error_reply( + self, e, event, source, session_entry, session_key, history, message_text, + persist_user_message, persist_user_timestamp, persist_user_display_kind, + ): + """``except Exception`` body of the agent turn: stop typing, log, persist the inbound user + turn once, and build the sanitized user-facing error reply.""" + # Stop typing indicator on error too, retaining 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; if the agent already reached turn-start persistence the latest user + # row matches and we skip the duplicate. + try: + if message_text is not None and session_entry is not None: + _already_persisted = False + try: + _recent_transcript = await self.async_session_store.load_transcript(session_entry.session_id) + except Exception: + _recent_transcript = [] + for _msg in reversed(_recent_transcript[-10:]): + if _msg.get("role") == "user": + _expected_user_content = ( + persist_user_message + if persist_user_message is not None + else message_text + ) + _already_persisted = (_msg.get("content") == _expected_user_content) + break + if not _already_persisted: + _user_entry = { + "role": "user", + "content": ( + persist_user_message + if persist_user_message is not None + else message_text + ), + "timestamp": ( + persist_user_timestamp + if persist_user_timestamp is not None + else time.time() + ), + } + if persist_user_display_kind: + _user_entry["display_kind"] = persist_user_display_kind + if getattr(event, "message_id", None): + _user_entry["message_id"] = str(event.message_id) + await self.async_session_store.append_to_transcript( + session_entry.session_id, + _user_entry, + ) + 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). + status_hint = "" + status_code = getattr(e, "status_code", None) + _hist_len = len(history) + if status_code == 401: + status_hint = " Check your API key or run `claude /login` to refresh OAuth credentials." + elif status_code == 402: + status_hint = " Your API balance or quota is exhausted. Check your provider dashboard." + elif status_code == 429: + # Check if this is a plan usage limit (resets on a schedule) vs a transient rate limit + _err_body = getattr(e, "response", None) + _err_json = {} + try: + if _err_body is not None: + _err_json = _err_body.json().get("error", {}) + if not isinstance(_err_json, dict): + _err_json = {} + except Exception: + pass + if _err_json.get("type") == "usage_limit_reached": + _resets_in = _err_json.get("resets_in_seconds") + if _resets_in and _resets_in > 0: + import math + _hours = math.ceil(_resets_in / 3600) + status_hint = f" Your plan's usage limit has been reached. It resets in ~{_hours}h." + else: + status_hint = " Your plan's usage limit has been reached. Please wait until it resets." + else: + status_hint = " You are being rate-limited. Please wait a moment and try again." + elif status_code == 529: + status_hint = " The API is temporarily overloaded. Please try again shortly." + 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. + if _hist_len > 50: + return ( + "⚠️ Session too large for the model's context window.\n" + "Use /compact to compress the conversation, or " + "/reset to start fresh." + ) + elif status_code == 400: + status_hint = " The request was rejected by the API." + return ( + f"Sorry, I encountered an unexpected error.{status_hint}\n" + "Try again or use /reset to start a fresh session." + ) + + async def _handle_message_with_agent(self, event, source, _quick_key: str, run_generation: int): + """Inner handler that runs under the _running_agents sentinel guard.""" + _msg_start_time = time.time() + _platform_name = source.platform.value if hasattr(source.platform, "value") else str(source.platform) + _msg_preview = (event.text or "")[:80].replace("\n", " ") + _reply_id = getattr(event, "reply_to_message_id", None) + _reply_txt = (getattr(event, "reply_to_text", None) or "")[:80].replace("\n", " ") + logger.info( + "inbound message: platform=%s user=%s chat=%s msg=%r reply_to_id=%s reply_to_text=%r", + _platform_name, source.user_name or source.user_id or "unknown", + source.chat_id or "unknown", _msg_preview, _reply_id, _reply_txt, + ) + + resolved = await self._hmwa_resolve_session(event, source) + if resolved is None: + return + source, session_entry, session_key = resolved + _was_auto_reset, _is_new_session = await self._hmwa_open_session( + session_entry, session_key, source + ) + + # Build session context + context = build_session_context(source, self.config, session_entry) + + # Set session context variables for tools (task-local, concurrency-safe) + _session_env_tokens = self._set_session_env(context) + + # Read privacy.redact_pii from config (re-read per message) + _redact_pii = False + # Synthetic self-injected turns (batch completions, watch notifications, resume wake-ups) + # arrive as MessageEvent(internal=True). Persist with display_kind="internal_notification" + # so UIs render timeline notices, not user bubbles. display_kind is a DB-only sidecar + # stripped from every provider-bound payload; role/content untouched. + persist_user_display_kind = ( + "internal_notification" if getattr(event, "internal", False) else None + ) + try: + _pcfg = _load_gateway_config() + _redact_pii = bool((_pcfg.get("privacy") or {}).get("redact_pii", False)) + except Exception: + pass + + # Build the context prompt. The 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. + context_prompt = self._pinned_session_context_prompt( + context, _redact_pii, session_key + ) + + # Per-turn must-deliver notes ride the user message via the api_content sidecar (staged + # below, consumed in run_sync → build_turn_context), NOT context_prompt: appending them to + # the ephemeral system prompt guaranteed a turn1→turn2 diff and a full agent rebuild. + turn_sidecar_notes: List[str] = [] + + # If the previous session expired and was auto-reset, deliver a notice + # so the agent knows this is a fresh conversation (not an intentional /reset). + 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 (Telegram DM Topics, Discord + # channel_skill_bindings). Supports a single name or ordered list. Only inject on NEW + # sessions — ongoing conversations already carry the skill content in their history. + _auto = getattr(event, "auto_skill", None) + if _is_new_session and _auto: + self._hmwa_auto_load_skills(event, _auto, _quick_key, session_key) + + # Turn lease: session resolution is FINAL here (see _hmwa_acquire_turn_lease). + 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. + await self._mark_durable_active_turn(event, session_entry.session_key) + + # Load conversation history from transcript. An unreadable canonical store is not an empty + # conversation: stop before the agent can invent continuity from a plausible-looking []. + # This return happens before the broad cleanup finally below, so restore task-local context + # here; the outer dispatch still clears the durable marker and turn lease. + try: + history = await self.async_session_store.load_transcript( + session_entry.session_id + ) + except TranscriptReadError: + self._clear_session_env(_session_env_tokens) + return ( + "⚠️ This session's history is temporarily unavailable, so " + "this message was not processed. Ask the operator to inspect " + "state.db, then resend after it is healthy. Use /reset only " + "if you intentionally want to start a new conversation." + ) + + # Session hygiene: auto-compress pathologically large transcripts before the agent starts. + history = await self._hmwa_run_session_hygiene( + event, source, session_entry, session_key, history, _quick_key, run_generation, + ) + + await self._hmwa_first_contact_notes(source, history, turn_sidecar_notes) + + # Voice channel awareness: deliver voice channel state (who is present / speaking) on the + # user message, ONLY when changed since the previous turn. It differs almost every turn, and + # in the ephemeral system prompt it forced a full agent rebuild + prompt-cache re-key per + # message; the system prompt carries a static pointer line instead (gateway/session.py). + _vc_note = self._voice_channel_sidecar_note(event, source, session_key) + if _vc_note: + turn_sidecar_notes.append(_vc_note) + + # Auto-analyze user images: run the vision tool eagerly so the model always gets a text + # description plus the local path for re-examination via vision_analyze. Filter to image + # media_type so documents/audio in the same message are not sent to the vision tool. + message_text = await self._prepare_profile_scoped_inbound_message_text( + event=event, + source=source, + history=history, + session_key=session_key, + ) + if message_text is None: + return + + # Capture the platform event time as message metadata and keep the persisted transcript + # clean; only the in-context render is gated behind gateway.message_timestamps.enabled. + message_text, persist_user_message, persist_user_timestamp = ( + 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. @@ -19911,25 +20782,7 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew ) _turn_seconds = time.monotonic() - _turn_started_monotonic - # Stop the typing indicator. Slack AI status is scoped to a thread/workspace, so - # preserve the routing metadata used by the response delivery path. - try: - _typing_adapter = self._adapter_for_source(source) - _stop_with_metadata = getattr( - type(_typing_adapter), "_stop_typing_with_metadata", None - ) - _stop_typing = getattr(type(_typing_adapter), "stop_typing", None) - if _typing_adapter and callable(_stop_with_metadata): - await _typing_adapter._stop_typing_with_metadata( - source.chat_id, - self._thread_metadata_for_source( - source, self._reply_anchor_for_event(event) - ), - ) - elif _typing_adapter and callable(_stop_typing): - await _typing_adapter.stop_typing(source.chat_id) - except Exception: - pass + await self._hmwa_stop_typing_for_turn(event, source) if not self._is_session_run_current(_quick_key, run_generation): logger.info( @@ -19947,607 +20800,48 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew _stale_adapter._post_delivery_callbacks.pop(_quick_key, None) return None - response = agent_result.get("final_response") or "" - # Hidden-reasoning-only retry exhaustion: the loop's sentinel text ("Codex response - # remained incomplete after 3 continuation attempts") doubles as final_response, so it - # would be delivered verbatim into the channel — where peer agents can ingest it as a - # completed assistant turn. - if _is_gateway_hidden_reasoning_incomplete_turn(agent_result): - response = "" - try: - from gateway.response_filters import is_intentional_silence_agent_result - _intentional_silence = is_intentional_silence_agent_result( - agent_result, response, - ) - except Exception: - _intentional_silence = False - - # Convert the agent's internal "(empty)" sentinel into a user-friendly message. - # "(empty)" means the model failed to produce visible content after exhausting all - # retries (nudge, prefill, empty-retry, fallback). - if response == "(empty)" and not _intentional_silence: - response = ( - "⚠️ The model returned no response after processing tool " - "results. This can happen with some models — try again or " - "rephrase your question." - ) - agent_messages = agent_result.get("messages", []) - _response_time = time.time() - _msg_start_time - _api_calls = agent_result.get("api_calls", 0) - _resp_len = len(response) - logger.info( - "response ready: platform=%s chat=%s time=%.1fs api_calls=%d response=%d chars", - _platform_name, source.chat_id or "unknown", - _response_time, _api_calls, _resp_len, + response, _intentional_silence, agent_messages = await self._hmwa_shape_agent_response( + agent_result, source, history, session_entry, session_key, + _quick_key, run_generation, _run_start_session_id, _platform_name, _msg_start_time, ) - - # The cross-process cache-coherence re-baseline (_refresh_agent_cache_message_count) is - # deferred until AFTER the transcript persistence block below: it must include the - # first-turn `session_meta` marker row and the compression session_id swap. - - # Successful turn: clear the stuck-loop counter (it only accumulates 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. - if session_key and _should_clear_resume_pending_after_turn(agent_result): - await self._clear_restart_failure_count(session_key) - try: - await self.async_session_store.clear_resume_pending(session_key) - except Exception as _e: - logger.debug( - "clear_resume_pending failed for %s: %s", - session_key, _e, - ) - - # Normalize empty responses: surface errors, partial failures, and - # the case where agent did work but returned no text. Fix for #18765. - if not _intentional_silence: - 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(). - 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: the transcript persistence below - # writes to the NEW id, so the serialization boundary must move with it or an - # alias key resolving the fresh child could interleave. - 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( - session_entry.session_id, - session_key, - source, - ) - await asyncio.to_thread( - self._sync_telegram_topic_binding, - source, session_entry, reason="agent-result-compression", - ) - else: - logger.info( - "Skipping agent-result session split sync for %s because " - "the session binding moved from %s to %s before " - "compression finished", - session_key or "?", - _run_start_session_id, - session_entry.session_id, - ) - - # Prepend reasoning if display is enabled (per-platform). Mattermost requires explicit - # opt-in because this is scratch text, not ordinary final-answer content. - try: - _show_reasoning_effective = _resolve_gateway_display_bool( - _load_gateway_config(), - _platform_config_key(source.platform), - "show_reasoning", - default=bool(getattr(self, "_show_reasoning", False)), - platform=source.platform, - require_platform_override_for={Platform.MATTERMOST}, - ) - except Exception: - _show_reasoning_effective = ( - False - if source.platform == Platform.MATTERMOST - else getattr(self, "_show_reasoning", False) - ) - if _show_reasoning_effective and response and not _intentional_silence: - last_reasoning = agent_result.get("last_reasoning") - if last_reasoning: - from gateway.stream_consumer import escape_code_fences_for_display - # Collapse long reasoning to keep messages readable - lines = last_reasoning.strip().splitlines() - if len(lines) > 15: - display_reasoning = "\n".join(lines[:15]) - display_reasoning += 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. - try: - from gateway.display_config import resolve_display_setting - _reasoning_style = resolve_display_setting( - _load_gateway_config(), - _platform_config_key(source.platform), - "reasoning_style", - "code", - ) - except Exception: - _reasoning_style = "code" - if _reasoning_style == "subtext": - _quoted = "\n".join( - f"-# {ln}" if ln else "-#" for ln in display_reasoning.splitlines() - ) - response = f"-# 💭 Reasoning\n{_quoted}\n\n{response}" - elif _reasoning_style == "blockquote": - _quoted = "\n".join( - f"> {ln}" if ln else ">" for ln in display_reasoning.splitlines() - ) - response = f"> 💭 **Reasoning:**\n{_quoted}\n\n{response}" - else: - # Escape ``` inside reasoning so inner fences don't - # break the outer code block used to render it. - display_reasoning = escape_code_fences_for_display(display_reasoning) - response = f"💭 **Reasoning:**\n```\n{display_reasoning}\n```\n\n{response}" - - # Runtime-metadata footer — only on the FINAL message of the turn. Off by default - # (display.runtime_footer.enabled=false). When streaming already delivered the body, we - # can't mutate the sent text, so we fire a separate trailing send below. - _footer_line = "" - try: - from gateway.runtime_footer import build_footer_line as _bfl - _footer_line = _bfl( - user_config=_load_gateway_config(), - platform_key=_platform_config_key(source.platform), - model=agent_result.get("model"), - context_tokens=agent_result.get("last_prompt_tokens", 0) or 0, - context_length=agent_result.get("context_length") or None, - cwd=_terminal_scope_cwd(""), - turn_seconds=_turn_seconds, - ) - except Exception as _footer_err: - logger.debug("runtime_footer build failed: %s", _footer_err) - _footer_line = "" + response = self._hmwa_prepend_reasoning(agent_result, response, source, _intentional_silence) + _footer_line = self._hmwa_runtime_footer_line(agent_result, source, _turn_seconds) if _footer_line and response and not agent_result.get("already_sent") and not _intentional_silence: response = f"{response}\n\n{_footer_line}" + await self._hmwa_post_turn_hooks(hook_ctx, agent_result, response) - # Emit agent:end hook - await self.hooks.emit("agent:end", { - **hook_ctx, - "response": (response or "")[:500], - "model": agent_result.get("model", ""), - "provider": agent_result.get("provider", ""), - }) - - # Check for pending process watchers (check_interval on background processes) - try: - from tools.process_registry import process_registry - # Detach the current batch atomically (see crash-recovery drain - # above): reassign to a fresh list so a watcher appended by a - # concurrent session during the yield isn't dropped by clear(). - watchers = process_registry.pending_watchers - process_registry.pending_watchers = [] - for i, watcher in enumerate(watchers): - asyncio.create_task(self._run_process_watcher(watcher)) - if i % 100 == 99: - await asyncio.sleep(0) - 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 (handled by the per-process watcher task above) and async-delegation - # completions (owned by _async_delegation_watcher, the single consumer for idle and - # post-turn cases) — inject only watch-type events and leave the rest on the queue. - try: - from tools.process_registry import process_registry as _pr - await self._drain_watch_notifications(_pr.completion_queue) - except Exception as e: - logger.debug("Watch queue drain error: %s", e) - - # NOTE: Dangerous command approvals are now handled inline by the blocking gateway - # approval mechanism in tools/approval.py. - - # Persist the full agent loop (tool calls, results, reasoning) so sessions resume with - # full context. IMPORTANT: on context-overflow failures (compression exhausted, generic - # 400 on large sessions) do NOT persist the user message — it would grow the session and - # reproduce the failure forever. Transient failures (429, timeout, connection error, - # 5xx) are different: the session is not oversized and dropping the user turn causes - # severe context loss on retry, so persist it. - agent_failed_early = bool(agent_result.get("failed")) - hidden_reasoning_incomplete = _is_gateway_hidden_reasoning_incomplete_turn( - agent_result + agent_failed_early, hidden_reasoning_incomplete, is_context_overflow_failure = ( + self._hmwa_classify_turn_failure(agent_result, history, session_entry) ) - _err_str_for_classify = str(agent_result.get("error", "")).lower() - # Use specific multi-word phrases (not bare "exceed"/"token") to avoid false positives - # on transient errors such as "rate limit exceeded"; matches run_agent.py's classifier. - is_context_overflow_failure = agent_failed_early and ( - bool(agent_result.get("compression_exhausted")) - or any(p in _err_str_for_classify for p in ( - "context length", "context size", "context window", - "maximum context", "token limit", "too many tokens", - "reduce the length", "exceeds the limit", - "request entity too large", "prompt is too long", - "payload too large", "input is too long", - )) - or ("400" in _err_str_for_classify and len(history) > 50) + response, session_entry = await self._hmwa_compression_exhaustion_reset( + agent_result, response, session_entry, session_key, source, ) - if is_context_overflow_failure: - logger.info( - "Skipping transcript persistence for context-overflow " - "failure in session %s to prevent session growth loop.", - session_entry.session_id, - ) - elif agent_failed_early: - logger.info( - "Transient agent failure in session %s — persisting user " - "message so conversation context is preserved on retry.", - session_entry.session_id, - ) - elif hidden_reasoning_incomplete: - logger.warning( - "Suppressing hidden-reasoning-only incomplete gateway turn " - "for session %s: %s", - session_entry.session_id, - agent_result.get("error", "processing incomplete"), - ) - - # Compression exhausted = permanently too large: auto-reset so the next message starts - # fresh instead of replaying the oversized context forever. A lock-contended defer is - # the OPPOSITE case (a concurrent path holds the lock and is shrinking it): never wipe - # for that. - if agent_result.get("compression_deferred"): - logger.info( - "Compression deferred for session %s — the compression " - "lock is held by a concurrent compressor. Keeping the " - "session intact; the next message retries normally.", - session_entry.session_id if session_entry else "?", - ) - elif agent_result.get("compression_exhausted") and session_entry and session_key: - 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). - 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: compression rotated - # session_entry.session_id to the bloated child earlier this turn and that _sync - # also rewrote the (chat_id, thread_id) binding. Without a re-sync the - # binding-heal walk switches the next inbound message back onto the child and - # re-triggers exhaustion forever. No-op on non-topic lanes. - session_entry = new_entry - await asyncio.to_thread( - self._sync_telegram_topic_binding, - source, session_entry, reason="compression-exhausted-reset", - ) - response = (response or "") + ( - "\n\n🔄 Session auto-reset — the conversation exceeded the " - "maximum context size and could not be compressed further. " - "Your next message will start a fresh session." - ) - - ts = time.time() # Unix epoch float — consistent with DB storage - - # Fresh session (no history): 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. - if is_context_overflow_failure: - pass # Skip all transcript writes — don't grow a broken session - elif not history: - tool_defs = agent_result.get("tools", []) - await self.async_session_store.append_to_transcript( - session_entry.session_id, - { - "role": "session_meta", - "tools": tool_defs or [], - "model": _resolve_gateway_model(), - "platform": source.platform.value if source.platform else "", - "timestamp": ts, - } - ) - - # The agent already persisted these via _flush_messages_to_session_db(); skip the DB - # write to avoid duplicates. Holds for the codex app-server runtime too (it flushes its - # own projected messages before returning and reports agent_persisted=True). Reading the - # flag (default = self._session_db is not None) keeps the contract explicit; a - # non-persisting runtime opts in via False. - agent_persisted = agent_result.get("agent_persisted", self._session_db is not None) - - # Only the NEW messages from this turn: use history_offset (what the agent saw), not - # len(history), which counts session_meta entries stripped before the agent saw them. - if is_context_overflow_failure: - pass # handled above — skip all transcript writes - elif agent_failed_early or hidden_reasoning_incomplete: - # Transient failure (429/timeout/5xx): persist only the user message so the next - # message can load a transcript that reflects what was said. Skip the assistant - # error text since it's a gateway-generated hint, not model output. Hidden-reasoning - # incomplete turns follow the same rule so peer-agent channels don't ingest them. - _user_entry = { - "role": "user", - "content": ( - persist_user_message - if persist_user_message is not None - else message_text - ), - "timestamp": ( - persist_user_timestamp - if persist_user_timestamp is not None - else ts - ), - } - if persist_user_display_kind: - _user_entry["display_kind"] = persist_user_display_kind - if event.message_id: - _user_entry["message_id"] = str(event.message_id) - # Dedupe: skip if this platform message_id is already in the transcript (prevents - # duplicate user turns on Telegram retries after transient failures). - _skip_persist = ( - event.message_id - and await self.async_session_store.has_platform_message_id( - session_entry.session_id, str(event.message_id) - ) - ) - if _skip_persist: - logger.info( - "Skipping duplicate user turn " - "(message_id=%s) in session %s", - event.message_id, session_entry.session_id, - ) - else: - await self.async_session_store.append_to_transcript( - session_entry.session_id, - _user_entry, - skip_db=agent_persisted, - ) - else: - history_len = agent_result.get("history_offset", len(history)) - new_messages = agent_messages[history_len:] if len(agent_messages) > history_len else [] - - # If no new messages found (edge case), fall back to simple user/assistant - if not new_messages: - _user_entry = { - "role": "user", - "content": ( - persist_user_message - if persist_user_message is not None - else message_text - ), - "timestamp": ( - persist_user_timestamp - if persist_user_timestamp is not None - else ts - ), - } - if persist_user_display_kind: - _user_entry["display_kind"] = persist_user_display_kind - if event.message_id: - _user_entry["message_id"] = str(event.message_id) - await self.async_session_store.append_to_transcript( - session_entry.session_id, - _user_entry, - skip_db=agent_persisted, - ) - if response: - await self.async_session_store.append_to_transcript( - session_entry.session_id, - {"role": "assistant", "content": response, "timestamp": ts}, - 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. - _user_msg_id_attached = False - for msg in new_messages: - # Skip system messages (they're rebuilt each run) - if msg.get("role") == "system": - continue - # Add timestamp to each message for debugging - entry = {**msg, "timestamp": ts} - if ( - not _user_msg_id_attached - and msg.get("role") == "user" - and event.message_id - and "message_id" not in entry - ): - entry["message_id"] = str(event.message_id) - _user_msg_id_attached = True - await self.async_session_store.append_to_transcript( - session_entry.session_id, 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. - await self.async_session_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)), + await self._hmwa_persist_turn_transcript( + event=event, + source=source, + session_entry=session_entry, + session_key=session_key, + agent_result=agent_result, + agent_messages=agent_messages, + history=history, + response=response, + message_text=message_text, + persist_user_message=persist_user_message, + persist_user_timestamp=persist_user_timestamp, + persist_user_display_kind=persist_user_display_kind, + agent_failed_early=agent_failed_early, + hidden_reasoning_incomplete=hidden_reasoning_incomplete, + is_context_overflow_failure=is_context_overflow_failure, ) - - # Re-baseline the cached agent's message_count snapshot now that ALL of this turn's - # transcript writes are done (flushed rows AND the first-turn `session_meta` marker). - # The cross-process coherence guard snapshots at agent-BUILD time and never refreshes - # on reuse, so our own writes would trigger a rebuild next turn (destroying prompt - # caching). MUST run after the session_meta append (that row bumps the count too). - await self._refresh_agent_cache_message_count( - session_key, session_entry.session_id + return await self._hmwa_deliver_turn_response( + event, source, session_entry, session_key, run_generation, + agent_result, agent_messages, response, _footer_line, _intentional_silence, ) - # Intentional silence is a delivery decision, not a transcript mutation: the [SILENT] - # assistant turn stays persisted so later turns keep user/assistant alternation; only - # the outbound delivery is suppressed. - if _intentional_silence: - logger.info( - "Suppressing intentional silence marker for session %s", - session_entry.session_id, - ) - response = "" - - # Auto voice reply: send TTS audio before the text response - _already_sent = bool(agent_result.get("already_sent")) - # Skip when streaming TTS already delivered audio for this turn (#60671). - _stts_adapter = self._adapter_for_source(source) - _streaming_tts_done = ( - _stts_adapter is not None - and bool(getattr(_stts_adapter, "_streaming_tts_turn_completed", lambda *_a, **_k: False)(session_key, run_generation)) - ) - if ( - not _streaming_tts_done - and self._should_send_voice_reply(event, response, agent_messages, already_sent=_already_sent) - ): - 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. - if agent_result.get("already_sent") and not agent_result.get("failed"): - if response: - _media_adapter = self._adapter_for_source(source) - if _media_adapter: - await self._deliver_media_from_response( - response, event, _media_adapter, - ) - # Streaming already delivered the body text, but the footer was intentionally held - # back (see the `not already_sent` gate above). - if _footer_line: - try: - _foot_adapter = self._adapter_for_source(source) - if _foot_adapter: - await _foot_adapter.send( - source.chat_id, - _footer_line, - metadata=self._thread_metadata_for_source(source, self._reply_anchor_for_event(event)), - ) - except Exception as _e: - logger.debug("trailing footer send failed: %s", _e) - # This branch returns 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. - with suppress(Exception): - event._streamed_final_response = str(response or "") - return None - - return response - except Exception as e: - # Stop typing indicator on error too, retaining Slack thread/workspace - # routing so a failed turn cannot leave its status visible. - try: - _err_adapter = self._adapter_for_source(source) - _stop_with_metadata = getattr( - type(_err_adapter), "_stop_typing_with_metadata", None - ) - _stop_typing = getattr(type(_err_adapter), "stop_typing", None) - if _err_adapter and callable(_stop_with_metadata): - await _err_adapter._stop_typing_with_metadata( - source.chat_id, - self._thread_metadata_for_source( - source, self._reply_anchor_for_event(event) - ), - ) - elif _err_adapter and callable(_stop_typing): - await _err_adapter.stop_typing(source.chat_id) - except Exception: - pass - 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; if the agent already reached turn-start persistence the latest user - # row matches and we skip the duplicate. - try: - if 'message_text' in locals() and message_text is not None and session_entry is not None: - _already_persisted = False - try: - _recent_transcript = await self.async_session_store.load_transcript(session_entry.session_id) - except Exception: - _recent_transcript = [] - for _msg in reversed(_recent_transcript[-10:]): - if _msg.get("role") == "user": - _expected_user_content = ( - persist_user_message - if persist_user_message is not None - else message_text - ) - _already_persisted = (_msg.get("content") == _expected_user_content) - break - if not _already_persisted: - _user_entry = { - "role": "user", - "content": ( - persist_user_message - if persist_user_message is not None - else message_text - ), - "timestamp": ( - persist_user_timestamp - if persist_user_timestamp is not None - else time.time() - ), - } - if 'persist_user_display_kind' in locals() and persist_user_display_kind: - _user_entry["display_kind"] = persist_user_display_kind - if getattr(event, "message_id", None): - _user_entry["message_id"] = str(event.message_id) - await self.async_session_store.append_to_transcript( - session_entry.session_id, - _user_entry, - ) - 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). - status_hint = "" - status_code = getattr(e, "status_code", None) - _hist_len = len(history) if 'history' in locals() else 0 - if status_code == 401: - status_hint = " Check your API key or run `claude /login` to refresh OAuth credentials." - elif status_code == 402: - status_hint = " Your API balance or quota is exhausted. Check your provider dashboard." - elif status_code == 429: - # Check if this is a plan usage limit (resets on a schedule) vs a transient rate limit - _err_body = getattr(e, "response", None) - _err_json = {} - try: - if _err_body is not None: - _err_json = _err_body.json().get("error", {}) - if not isinstance(_err_json, dict): - _err_json = {} - except Exception: - pass - if _err_json.get("type") == "usage_limit_reached": - _resets_in = _err_json.get("resets_in_seconds") - if _resets_in and _resets_in > 0: - import math - _hours = math.ceil(_resets_in / 3600) - status_hint = f" Your plan's usage limit has been reached. It resets in ~{_hours}h." - else: - status_hint = " Your plan's usage limit has been reached. Please wait until it resets." - else: - status_hint = " You are being rate-limited. Please wait a moment and try again." - elif status_code == 529: - status_hint = " The API is temporarily overloaded. Please try again shortly." - 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. - if _hist_len > 50: - return ( - "⚠️ Session too large for the model's context window.\n" - "Use /compact to compress the conversation, or " - "/reset to start fresh." - ) - elif status_code == 400: - status_hint = " The request was rejected by the API." - return ( - f"Sorry, I encountered an unexpected error.{status_hint}\n" - "Try again or use /reset to start a fresh session." + return await self._hmwa_agent_error_reply( + e, event, source, session_entry, session_key, history, message_text, + persist_user_message, persist_user_timestamp, persist_user_display_kind, ) finally: # Restore session context variables to their pre-handler state