diff --git a/gateway/slash_commands.py b/gateway/slash_commands.py index 51fb8c71a0..7f8a2243de 100644 --- a/gateway/slash_commands.py +++ b/gateway/slash_commands.py @@ -28,8 +28,7 @@ from gateway.session import AsyncSessionStore from gateway.slash_commands_goals import GatewayGoalCommandsMixin from gateway.slash_commands_model import ( # noqa: F401 — _model_switch_skew_guard re-exported for tests GatewayModelCommandsMixin, - _model_switch_skew_guard, -) + _model_switch_skew_guard) from gateway.slash_commands_session import GatewaySessionCommandsMixin from gateway.slash_commands_status import GatewayStatusCommandsMixin from hermes_cli.config import atomic_config_write, cfg_get @@ -42,8 +41,7 @@ logger = logging.getLogger("gateway.run") _ROLLBACK_SKIP_LINES = ( ("skipped_user_edits", "gateway.rollback.kept_user_edits"), ("skipped_oversize", "gateway.rollback.kept_oversize"), - ("failed_deletes", "gateway.rollback.failed_deletes"), -) + ("failed_deletes", "gateway.rollback.failed_deletes")) # /busy input modes -> (status-card behavior, set-confirmation behavior). _BUSY_MODE_BEHAVIOR = { @@ -57,34 +55,29 @@ _BUSY_MODE_BEHAVIOR = { _DIFF_MODE_BY_ARG = { "staged": "staged", "--staged": "staged", "cached": "staged", "--cached": "staged", "all": "all", "--all": "all", "head": "all", - "session": "session", -} + "session": "session"} # /voice subcommand -> stored mode (None = auto-TTS disabled), confirmation i18n key. _VOICE_MODE_BY_ARG = { **dict.fromkeys(("on", "enable"), ("voice_only", "gateway.voice.enabled_voice_only")), **dict.fromkeys(("off", "disable"), ("off", "gateway.voice.disabled_text")), - "tts": ("all", "gateway.voice.tts_enabled"), -} + "tts": ("all", "gateway.voice.tts_enabled")} # /footer argument -> new enabled state ("" toggles; anything else is a usage error). _FOOTER_STATE_BY_ARG = { **dict.fromkeys(("on", "enable", "true", "1"), True), - **dict.fromkeys(("off", "disable", "false", "0"), False), -} + **dict.fromkeys(("off", "disable", "false", "0"), False)} # /approve modifier tokens -> approval choice (default "once"). _APPROVE_CHOICE_BY_ARG = { **dict.fromkeys(("always", "permanent", "permanently"), "always"), - **dict.fromkeys(("session", "ses"), "session"), -} + **dict.fromkeys(("session", "ses"), "session")} _PLATFORM_USAGE = ( "Usage: /platform [name]\n" " /platform list — show platform status\n" " /platform pause — stop retrying a failing platform\n" - " /platform resume — re-queue a paused platform" -) + " /platform resume — re-queue a paused platform") _WINDOWS_UPDATE_HELPER = """ import os, subprocess, sys @@ -155,8 +148,7 @@ def _spawn_detached_update(hermes_cmd, output_path, exit_code_path) -> None: subprocess.Popen( [sys.executable, "-c", _WINDOWS_UPDATE_HELPER, str(output_path), str(exit_code_path), sys.executable, "-m", "hermes_cli.main", "update", "--gateway"], - stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, **windows_detach_popen_kwargs(), - ) + stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, **windows_detach_popen_kwargs()) return hermes_cmd_str = " ".join(shlex.quote(part) for part in hermes_cmd) update_cmd = ( @@ -164,8 +156,7 @@ def _spawn_detached_update(hermes_cmd, output_path, exit_code_path) -> None: f" > {shlex.quote(str(output_path))} 2>&1; " # Avoid `status=$?`: `status` is read-only in zsh and this template is reused in # macOS/zsh operator wrappers, so keep it zsh-safe even though bash runs it here. - f"rc=$?; printf '%s' \"$rc\" > {shlex.quote(str(exit_code_path))}" - ) + f"rc=$?; printf '%s' \"$rc\" > {shlex.quote(str(exit_code_path))}") # Preferred: setsid creates a new session, fully detached; fallback start_new_session=True # calls os.setsid() in the child. setsid_bin = shutil.which("setsid") @@ -192,8 +183,7 @@ class GatewaySlashCommandsMixin( GatewayModelCommandsMixin, GatewaySessionCommandsMixin, GatewayStatusCommandsMixin, - GatewayGoalCommandsMixin, -): + GatewayGoalCommandsMixin): """In-session slash-command handlers for GatewayRunner (plus the helpers the sibling mixins share).""" async_session_store: AsyncSessionStore @@ -282,8 +272,7 @@ class GatewaySlashCommandsMixin( try: await adapter.send( source.chat_id, confirmation_text, reply_to=event.message_id, - metadata={"is_approval_prompt": True, "force_proactive_send": True}, - ) + metadata={"is_approval_prompt": True, "force_proactive_send": True}) except Exception as exc: logger.warning( "Failed to send /%s confirmation to %s: %s", verb, source.chat_id, exc, exc_info=True, @@ -331,8 +320,7 @@ class GatewaySlashCommandsMixin( reply = _execute("profile", options={"profile_name": profile_name, "home_display": display}) return "\n".join([ t("gateway.profile.header", profile=reply.data["profile"]), - t("gateway.profile.home", home=reply.data["home"]), - ]) + t("gateway.profile.home", home=reply.data["home"])]) async def _handle_whoami_command(self, event: MessageEvent) -> str: """Handle /whoami — platform, DM-vs-group scope, tier and runnable commands (always allowed).""" @@ -424,8 +412,7 @@ class GatewaySlashCommandsMixin( user_id_alt=_field("user_id_alt"), notifier_profile=getattr(self, "_kanban_notifier_profile", None) or self._active_profile_name(), # Subscribing from chat: deliver the passive message and wake the destination agent. - delivery_mode="notify+wake", delivery_metadata=delivery_metadata, - ) + delivery_mode="notify+wake", delivery_metadata=delivery_metadata) finally: conn.close() await asyncio.to_thread(_sub) @@ -465,8 +452,7 @@ class GatewaySlashCommandsMixin( await _stop(sibling_key, "stop_command_thread_sibling") logger.info( "STOP (thread sibling) by %s — interrupted %d run(s) in thread: %s", - session_key, len(sibling_keys), ", ".join(sibling_keys), - ) + session_key, len(sibling_keys), ", ".join(sibling_keys)) return EphemeralReply(t("gateway.stop.stopped")) # No running agent anywhere for this scope. A platform status indicator can still be stuck — @@ -536,8 +522,7 @@ class GatewaySlashCommandsMixin( "Ignoring redelivered /restart (platform=%s, update_id=%s) — " "already processed by a previous gateway instance.", event.source.platform.value if event.source and event.source.platform else "?", - event.platform_update_id, - ) + event.platform_update_id) return "" if self._restart_requested or self._draining: count = self._running_agent_count() @@ -607,16 +592,14 @@ class GatewaySlashCommandsMixin( source.platform in {None, Platform.LOCAL, Platform.RELAY} or not getattr(source, "user_id", None) or not callable(fronts_platform) - or not fronts_platform(source.platform) - ): + or not fronts_platform(source.platform)): return t("gateway.set_home.save_failed", error="Relay does not authenticate this logical home target") thread_id = _home_thread_from_source(source) home = HomeChannel( platform=source.platform, chat_id=str(chat_id), name=chat_name, thread_id=thread_id, user_id=str(source.user_id) if getattr(source, "user_id", None) else None, - scope_id=str(source.scope_id) if getattr(source, "scope_id", None) else None, - ) + scope_id=str(source.scope_id) if getattr(source, "scope_id", None) else None) # config.yaml is canonical because it can persist the authenticated logical-target # provenance required by Relay after a restart. try: @@ -797,8 +780,7 @@ class GatewaySlashCommandsMixin( prompt, event.source, task_id, event_message_id=self._reply_anchor_for_event(event), # Forward image/audio attachments so the background agent can see them. media_urls=list(event.media_urls) if event.media_urls else [], - media_types=list(event.media_types) if event.media_types else [], - )) + media_types=list(event.media_types) if event.media_types else [])) return t("gateway.background.started", preview=_preview(prompt), task_id=task_id) async def _handle_btw_command(self, event: MessageEvent) -> str: @@ -837,8 +819,7 @@ class GatewaySlashCommandsMixin( try: answer = await asyncio.to_thread( answer_side_question, question, history_snapshot, - parent_agent=parent_agent, main_runtime=main_runtime, - ) + parent_agent=parent_agent, main_runtime=main_runtime) reply = t("gateway.btw.answer", preview=preview, answer=answer or "") except Exception as e: logger.warning("/btw side question failed: %s", e) @@ -859,8 +840,7 @@ class GatewaySlashCommandsMixin( # the store persists to the same MEMORY/USER.md and honors the configured char limits). out = handle_pending_subcommand( wa.MEMORY, event.get_command_args().strip().split(), memory_store=load_on_disk_store(), - set_mode_fn=self._write_approval_setter("memory", event), - ) + set_mode_fn=self._write_approval_setter("memory", event)) return out if out is not None else ( "Unknown /memory subcommand. Use: pending, approve , reject , approval ." ) @@ -879,8 +859,7 @@ class GatewaySlashCommandsMixin( "Enable it with /skills approval on, then review staged " "writes here with /skills pending.") out = handle_pending_subcommand( - wa.SKILLS, args, set_mode_fn=self._write_approval_setter("skills", event) - ) + wa.SKILLS, args, set_mode_fn=self._write_approval_setter("skills", event)) if out is None: return ("Unknown /skills subcommand on this platform. Use: pending, " "approve , reject , diff , approval . " @@ -929,8 +908,7 @@ class GatewaySlashCommandsMixin( try: user_config = _load_gateway_config() gate_enabled = is_truthy_value( - cfg_get(user_config, "display", "tool_progress_command"), default=False - ) + cfg_get(user_config, "display", "tool_progress_command"), default=False) except Exception: gate_enabled = False if not gate_enabled: @@ -958,13 +936,11 @@ class GatewaySlashCommandsMixin( return EphemeralReply( f"**Busy input mode: `{mode}`" + "\n" f"Messages while busy: _{behavior}_" + "\n" - f"Change with `/busy queue`, `/busy steer`, or `/busy interrupt`." - ) + f"Change with `/busy queue`, `/busy steer`, or `/busy interrupt`.") if arg not in _BUSY_MODE_BEHAVIOR: return EphemeralReply( - f"Unknown mode `{arg}`. Use `/busy queue`, `/busy steer`, or `/busy interrupt`." - ) + f"Unknown mode `{arg}`. Use `/busy queue`, `/busy steer`, or `/busy interrupt`.") # Persist before mutate from cli import save_config_value @@ -1029,8 +1005,7 @@ class GatewaySlashCommandsMixin( # Show a preview using current agent state if available. preview = format_runtime_footer( model=_resolve_gateway_model(user_config) or None, context_tokens=0, context_length=None, - fields=effective.get("fields") or ["model", "context_pct", "cwd"], - ) + fields=effective.get("fields") or ["model", "context_pct", "cwd"]) if preview: example = t("gateway.footer.example_line", preview=preview) return t("gateway.footer.saved", state=_state(new_state), example=example) @@ -1067,8 +1042,7 @@ class GatewaySlashCommandsMixin( return result return await self._request_slash_confirm( event=event, command="reload-mcp", title="/reload-mcp", - message=t("gateway.reload_mcp.confirm_prompt"), handler=_on_confirm, - ) + message=t("gateway.reload_mcp.confirm_prompt"), handler=_on_confirm) async def _handle_reload_skills_command(self, event: MessageEvent) -> str: """Handle /reload-skills — rescan skills dir, queue a note for next turn. Skills are invoked at @@ -1114,8 +1088,7 @@ class GatewaySlashCommandsMixin( sections = ["[USER INITIATED SKILLS RELOAD:"] for i18n_key, note_header, items in ( ("gateway.reload_skills.added_header", "Added Skills:", added), - ("gateway.reload_skills.removed_header", "Removed Skills:", removed), - ): + ("gateway.reload_skills.removed_header", "Removed Skills:", removed)): if items: formatted = [_fmt_line(item) for item in items] lines += [t(i18n_key)] + formatted @@ -1145,8 +1118,7 @@ class GatewaySlashCommandsMixin( "No skill bundles installed.\n" "Create one on the host with:\n" " `hermes bundles create --skill --skill `\n" - f"Directory: `{reply.data['dir']}`" - ) + f"Directory: `{reply.data['dir']}`") lines = [f"**Skill Bundles** ({len(bundles)} installed):", ""] for info in bundles: @@ -1172,8 +1144,7 @@ class GatewaySlashCommandsMixin( signalling the event resumes them so the command executes inline (same flow as the CLI).""" from tools.approval import resolve_gateway_approval session_key, stale = self._blocking_approval_or_stale( - event, "gateway.approval_expired", "gateway.approve.no_pending" - ) + event, "gateway.approval_expired", "gateway.approve.no_pending") if stale: return stale @@ -1193,8 +1164,7 @@ class GatewaySlashCommandsMixin( the CLI. ``/deny`` denies the oldest; ``/deny all`` denies everything.""" from tools.approval import resolve_gateway_approval session_key, stale = self._blocking_approval_or_stale( - event, "gateway.deny.stale", "gateway.deny.no_pending" - ) + event, "gateway.deny.stale", "gateway.deny.no_pending") if stale: return stale @@ -1219,8 +1189,7 @@ class GatewaySlashCommandsMixin( protect privacy; ``hermes debug share`` from the CLI does full uploads.""" from hermes_cli.debug import ( _GATEWAY_PRIVACY_NOTICE, _best_effort_sweep_expired_pastes, _capture_dump, _schedule_auto_delete, - collect_debug_report, upload_to_pastebin, - ) + collect_debug_report, upload_to_pastebin) # Run blocking I/O (dump capture, log reads, uploads) in a thread. def _collect_and_upload(): @@ -1274,8 +1243,7 @@ class GatewaySlashCommandsMixin( pending = { "platform": src.platform.value, "chat_id": src.chat_id, "chat_type": src.chat_type, "user_id": src.user_id, "session_key": self._session_key_for_source(src), - "timestamp": datetime.now().isoformat(), - } + "timestamp": datetime.now().isoformat()} pending.update({k: v for k, v in (("thread_id", src.thread_id), ("message_id", event.message_id)) if v}) _tmp_pending = pending_path.with_suffix(".tmp") _tmp_pending.write_text(json.dumps(pending), encoding="utf-8") diff --git a/gateway/slash_commands_session.py b/gateway/slash_commands_session.py index 7591d3cacd..012ec6dea0 100644 --- a/gateway/slash_commands_session.py +++ b/gateway/slash_commands_session.py @@ -32,8 +32,7 @@ _DM_CHAT_TYPES = {"dm", "direct", "private", ""} _BRANCH_COPIED_FIELDS = ( "content", "tool_calls", "tool_call_id", "finish_reason", "reasoning", "reasoning_content", - "reasoning_details", "codex_reasoning_items", "codex_message_items", "timestamp", -) + "reasoning_details", "codex_reasoning_items", "codex_message_items", "timestamp") def _sattr(obj, name: str) -> str: @@ -65,8 +64,7 @@ def _manual_compression_reply_lines(summary: dict, compressor, focus_topic) -> l lines.append(t( "gateway.compress.aux_failed", model=aux_fail_model, - error=(getattr(compressor, "_last_aux_model_failure_error", None) or "unknown error"), - )) + error=(getattr(compressor, "_last_aux_model_failure_error", None) or "unknown error"))) return lines @@ -78,11 +76,9 @@ def _compress_preview_reply(history, partial: bool, keep_last, focus_topic, agg_ pv_msgs = [ {"role": m.get("role"), "content": m.get("content")} for m in history - if m.get("role") in {"user", "assistant"} and m.get("content") - ] + if m.get("role") in {"user", "assistant"} and m.get("content")] report = summarize_compress_preview( - pv_msgs, partial, keep_last, focus_topic, estimate_request_tokens_rough(pv_msgs) - ) + pv_msgs, partial, keep_last, focus_topic, estimate_request_tokens_rough(pv_msgs)) lines = [f"🗜️ {line}" for line in report["lines"]] if agg_note: lines.append(agg_note) @@ -136,31 +132,26 @@ class GatewaySessionCommandsMixin: try: await asyncio.wait_for( self._run_in_executor_with_context(self._cleanup_agent_resources, _old_agent), - timeout=_RESET_CLEANUP_TIMEOUT_S, - ) + timeout=_RESET_CLEANUP_TIMEOUT_S) except asyncio.TimeoutError: logger.warning( "Agent resource cleanup for session %s exceeded %ss during /new reset; proceeding with " "reset (the worker thread is left to finish on its own). (#35994)", - session_key, _RESET_CLEANUP_TIMEOUT_S, - ) + session_key, _RESET_CLEANUP_TIMEOUT_S) except Exception as cleanup_exc: logger.warning( "Agent resource cleanup for session %s failed during /new reset: %s (#35994)", - session_key, cleanup_exc, - ) + session_key, cleanup_exc) async def _fire_session_reset_hooks( - self, source: SessionSource, session_key: str, old_sid, new_sid - ) -> None: + self, source: SessionSource, session_key: str, old_sid, new_sid) -> None: """Session-boundary hooks: plugin finalize (off-loop + bounded — trace exports can block arbitrarily), then session:end and session:reset.""" platform_value = source.platform.value if source.platform else "" with contextlib.suppress(Exception): await self._finalize_session_off_loop( session_id=old_sid, platform=platform_value, reason="new_session", - old_session_id=old_sid, new_session_id=new_sid, - ) + old_session_id=old_sid, new_session_id=new_sid) hook_payload = {"platform": platform_value, "user_id": source.user_id, "session_key": session_key} await self.hooks.emit("session:end", dict(hook_payload)) await self.hooks.emit("session:reset", dict(hook_payload)) @@ -172,8 +163,7 @@ class GatewaySessionCommandsMixin: _invoke_hook( "on_session_reset", session_id=new_sid, platform=source.platform.value if source.platform else "", reason="new_session", - old_session_id=old_sid, new_session_id=new_sid, - ) + old_session_id=old_sid, new_session_id=new_sid) except Exception: pass @@ -200,8 +190,7 @@ class GatewaySessionCommandsMixin: interrupt_for_session( session_key=session_key, parent_session_id=str(getattr(old_entry, "session_id", "") or ""), - reason="session_reset", - ) + reason="session_reset") except Exception: pass _reset_process_scoped_tool_state() @@ -209,8 +198,7 @@ class GatewaySessionCommandsMixin: new_entry = await self.async_session_store.reset_session(session_key) _old_sid = old_entry.session_id if old_entry else None await self._fire_session_reset_hooks( - source, session_key, _old_sid, new_entry.session_id if new_entry else None - ) + source, session_key, _old_sid, new_entry.session_id if new_entry else None) # Scoped to the profile serving this source so a multiplexed /new banner reports the # profile's model, not the base config's. try: @@ -234,8 +222,7 @@ class GatewaySessionCommandsMixin: except Exception: logger.debug("Failed to rebind Telegram topic after /new", exc_info=True) self._invoke_session_reset_lifecycle_hook( - source, _old_sid, new_entry.session_id if new_entry else None - ) + source, _old_sid, new_entry.session_id if new_entry else None) try: from hermes_cli.tips import get_random_tip _tip_line = t("gateway.reset.tip", tip=get_random_tip()) @@ -277,8 +264,7 @@ class GatewaySessionCommandsMixin: entries = getattr(self.session_store, "_entries", {}) or {} return next( (getattr(e, "origin", None) for e in entries.values() if getattr(e, "session_id", None) == session_id), - None, - ) + None) @staticmethod def _same_matrix_room(current: SessionSource, origin: Optional[SessionSource]) -> bool: @@ -289,8 +275,7 @@ class GatewaySessionCommandsMixin: and origin.platform == Platform.MATRIX and current.platform == Platform.MATRIX and origin.chat_id == current.chat_id - and _sattr(current, "thread_id") == _sattr(origin, "thread_id") - ) + and _sattr(current, "thread_id") == _sattr(origin, "thread_id")) def _same_origin_chat(self, current: SessionSource, origin: Optional[SessionSource]) -> bool: """Platform-agnostic counterpart to ``_same_matrix_room``. @@ -327,8 +312,7 @@ class GatewaySessionCommandsMixin: build_session_key's isolation rules so the guards stay in lock-step with the key.""" return is_shared_multi_user_session( source, group_sessions_per_user=getattr(self.config, "group_sessions_per_user", True), - thread_sessions_per_user=getattr(self.config, "thread_sessions_per_user", False), - ) + thread_sessions_per_user=getattr(self.config, "thread_sessions_per_user", False)) def _resume_caller_is_admin(self, source: SessionSource) -> bool: """Whether *source* is an EXPLICITLY-configured admin (cross-origin /resume, /sessions). @@ -381,8 +365,7 @@ class GatewaySessionCommandsMixin: return bool(row_uid) and row_uid == caller_uid async def _resume_target_allowed( - self, source: SessionSource, target_id: str, allow_override: bool = False - ) -> bool: + self, source: SessionSource, target_id: str, allow_override: bool = False) -> bool: """Whether *source* may resume session *target_id* (IDOR guard for every adapter). The live origin decides when the target is active; otherwise the DB row must PROVE @@ -404,8 +387,7 @@ class GatewaySessionCommandsMixin: return self._persisted_row_proves_owner(source, row) async def _resume_row_visible( - self, source: SessionSource, row: dict, allow_all: bool - ) -> bool: + self, source: SessionSource, row: dict, allow_all: bool) -> bool: """Whether a listing *row* belongs to the caller's origin (blocks cross-origin enumeration of ids/previews); Matrix is room-scoped, ``--all`` needs a configured admin everywhere.""" if allow_all and self._resume_caller_is_admin(source): @@ -425,16 +407,14 @@ class GatewaySessionCommandsMixin: history_before_user_originated_turn, retryable_user_text, split_user_originated_turn, - user_originated_turn_view, - ) + user_originated_turn_view) source = event.source session_entry = await self.async_session_store.get_or_create_session(source) history = await self.async_session_store.load_transcript(session_entry.session_id) last_user_idx = next( (i for i in range(len(history) - 1, -1, -1) if user_originated_turn_view(history[i]) is not None), - None, - ) + None) if last_user_idx is None: return t("gateway.retry.no_previous") # Resolve text + scaffold-preserving prefix BEFORE any write; messaging retries cannot @@ -452,8 +432,7 @@ class GatewaySessionCommandsMixin: # on the same snapshot so a concurrent newer turn is never removed for stale text. try: rewind_result = await self.async_session_store.rewind_session( - session_entry.session_id, 1, require_retryable_composite=True, - ) + session_entry.session_id, 1, require_retryable_composite=True) except ValueError as exc: return f"Cannot retry that message safely: {exc}" if rewind_result is None: @@ -461,14 +440,12 @@ class GatewaySessionCommandsMixin: last_user_msg = rewind_result["target_text"] # active_only preserves the active=0/compacted=1 archive left by in-place compaction. elif not await self.async_session_store.rewrite_transcript( - session_entry.session_id, truncated, active_only=True, reject_active_turn_lease=True, - ): + session_entry.session_id, truncated, active_only=True, reject_active_turn_lease=True): return "Retry failed; transcript was not changed." session_entry.last_prompt_tokens = 0 # transcript was truncated retry_event = MessageEvent( text=last_user_msg, message_type=MessageType.TEXT, source=source, - raw_message=event.raw_message, channel_prompt=event.channel_prompt, - ) + raw_message=event.raw_message, channel_prompt=event.channel_prompt) return await self._handle_message(retry_event) async def _handle_undo_command(self, event: MessageEvent) -> str: @@ -519,8 +496,7 @@ class GatewaySessionCommandsMixin: return ( "🗜️ Nothing to compact: this session runs on the Codex app-server runtime, whose " "context lives in a Codex-owned thread that only exists while the agent is active. " - "Send a message first, then /compress — or /reset to start fresh." - ) + "Send a message first, then /compress — or /reset to start fresh.") compressor = getattr(agent, "context_compressor", None) count_before = getattr(compressor, "compression_count", 0) try: @@ -530,12 +506,10 @@ class GatewaySessionCommandsMixin: if getattr(compressor, "compression_count", 0) > count_before: return ( "🗜️ Codex app-server thread compacted (thread/compact). The transcript mirror is " - "unchanged by design — the app-server now carries the compacted context." - ) + "unchanged by design — the app-server now carries the compacted context.") return ( "⚠️ Codex app-server compaction did not complete — the thread is unchanged. Check the " - "app-server logs, retry /compress, or /reset for a clean session." - ) + "app-server logs, retry /compress, or /reset for a clean session.") async def _handle_compress_command_inner(self, event: MessageEvent) -> str: """Handle /compress -- manually compress conversation context; ``/compress `` tells @@ -562,15 +536,13 @@ class GatewaySessionCommandsMixin: return _compress_preview_reply(history, partial, keep_last, focus_topic, _agg_note) try: return await self._run_manual_compression( - source, session_entry, history, partial, keep_last, focus_topic - ) + source, session_entry, history, partial, keep_last, focus_topic) except Exception as e: logger.warning("Manual compress failed: %s", e) return t("gateway.compress.failed", error=e) async def _run_manual_compression( - self, source, session_entry, history: list, partial: bool, keep_last, focus_topic - ) -> str: + self, source, session_entry, history: list, partial: bool, keep_last, focus_topic) -> str: """Build a temporary agent, compress the transcript, persist, and describe the outcome.""" from agent.conversation_compression import finalize_context_engine_compression_notification from agent.manual_compression_feedback import summarize_manual_compression @@ -578,8 +550,7 @@ class GatewaySessionCommandsMixin: from gateway.run import _platform_config_key from hermes_cli.partial_compress import ( rejoin_compressed_head_and_tail, - split_history_for_partial_compress, - ) + split_history_for_partial_compress) session_key = self._session_key_for_source(source) # Platform + stable gateway session key bind this agent (for external context engines) to @@ -622,9 +593,7 @@ class GatewaySessionCommandsMixin: compressed, _ = await self._run_in_executor_with_context( lambda: tmp_agent._compress_context( head, "", approx_tokens=approx_tokens, focus_topic=focus_topic, force=True, - defer_context_engine_notification=True, - ) - ) + defer_context_engine_notification=True)) # A held compression lock returns unchanged; say so instead of the misleading no-op text. _lock_skipped = getattr(tmp_agent, "_compression_skipped_due_to_lock", None) if _lock_skipped is True or isinstance(_lock_skipped, str): @@ -636,8 +605,7 @@ class GatewaySessionCommandsMixin: finalize_context_engine_compression_notification(tmp_agent, committed=True) new_tokens = estimate_request_tokens_rough(compressed, system_prompt=_sys_prompt, tools=_tools) summary = summarize_manual_compression( - msgs, compressed, approx_tokens, new_tokens, compression_state=compressor, - ) + msgs, compressed, approx_tokens, new_tokens, compression_state=compressor) finally: finalize_context_engine_compression_notification(tmp_agent, committed=False) self._evict_cached_agent(session_key) # next turn rebuilds the prompt from current files @@ -663,20 +631,17 @@ class GatewaySessionCommandsMixin: logger.warning( "Manual compression 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_id, exc, exc_info=True, - ) + "providers.", session_id, exc, exc_info=True) # compression.checkpoint_required needs the memory provider loaded so _compress_context() # can write the pre-compression checkpoint; otherwise keep the fast path (no provider init). _checkpoint_required = _is_truthy( ((_load_cfg() or {}).get("compression") or {}).get("checkpoint_required"), - default=False, - ) + default=False) tmp_agent = AIAgent( **runtime_kwargs, model=model, max_iterations=4, quiet_mode=True, skip_memory=not _checkpoint_required, enabled_toolsets=["memory"], - session_id=session_id, session_db=getattr(self._session_db, "_db", self._session_db), - ) + session_id=session_id, session_db=getattr(self._session_db, "_db", self._session_db)) _seed_hygiene_system_prompt(tmp_agent, session_row) # Real platform during construction (context engines bind correctly); afterwards a prompt # rebuilt by compression is stamped as the provider-less fallback, stale for the next turn. @@ -698,18 +663,15 @@ class GatewaySessionCommandsMixin: if new_session_id != session_entry.session_id: if not await self.async_session_store.rewrite_transcript(new_session_id, compressed): raise RuntimeError( - f"failed to persist compressed transcript for session {new_session_id}" - ) + f"failed to persist compressed transcript for session {new_session_id}") session_entry.session_id = new_session_id await self.async_session_store._save() await asyncio.to_thread( - self._sync_telegram_topic_binding, source, session_entry, reason="compress-command", - ) + self._sync_telegram_topic_binding, source, session_entry, reason="compress-command") elif not getattr(tmp_agent, "_last_compaction_in_place", False): logger.warning( "Manual /compress: session rotation did not occur (session_id unchanged) and in-place " - "mode is off — preserving original transcript instead of overwriting it (#44794)." - ) + "mode is off — preserving original transcript instead of overwriting it (#44794).") await self.async_session_store.update_session(session_entry.session_key, last_prompt_tokens=0) # ------------------------------------------------------------------------ /topic @@ -756,8 +718,7 @@ class GatewaySessionCommandsMixin: await self._session_db.enable_telegram_topic_mode( chat_id=str(source.chat_id), user_id=str(source.user_id), profile_name=profile_name, has_topics_enabled=capabilities.get("has_topics_enabled"), - allows_users_to_create_topics=capabilities.get("allows_users_to_create_topics"), - ) + allows_users_to_create_topics=capabilities.get("allows_users_to_create_topics")) except Exception as exc: logger.exception("Failed to enable Telegram topic mode") return t("gateway.topic.enable_failed", error=exc) @@ -768,8 +729,7 @@ class GatewaySessionCommandsMixin: try: binding = await self._session_db.get_telegram_topic_binding( chat_id=str(source.chat_id), thread_id=str(source.thread_id), - profile_name=profile_name, - ) + profile_name=profile_name) except Exception: logger.debug("Failed to read Telegram topic binding", exc_info=True) binding = None @@ -782,8 +742,7 @@ class GatewaySessionCommandsMixin: title = None return t( "gateway.topic.bound_status", label=title or t("gateway.topic.untitled_session"), - session_id=session_id, - ) + session_id=session_id) # ------------------------------------------------------------------ /save, /title @@ -794,8 +753,7 @@ class GatewaySessionCommandsMixin: SAVE_USAGE, default_save_filename, normalize_save_format, - render_session_for_save, - ) + render_session_for_save) parts = event.get_command_args().split() redact = bool(parts) and parts[-1].lower() in ("redact", "--redact") @@ -837,8 +795,7 @@ class GatewaySessionCommandsMixin: return "Platform adapter not found to send the document." await adapter.send_document( chat_id=source.chat_id, file_path=temp_path, caption=f"Session export: {filename}", - file_name=filename, - ) + file_name=filename) return "Export complete." except Exception as e: logger.warning("Session /save failed: %s", e) @@ -865,8 +822,7 @@ class GatewaySessionCommandsMixin: session_id=session_id, source=source.platform.value if source.platform else "unknown", user_id=source.user_id, chat_id=source.chat_id, chat_type=source.chat_type, - thread_id=source.thread_id, - ) + thread_id=source.thread_id) title_arg = event.get_command_args().strip() if not title_arg: title = await self._session_db.get_session_title(session_id) @@ -899,8 +855,7 @@ class GatewaySessionCommandsMixin: widen = allow_all and self._resume_caller_is_admin(source) sessions = await self._session_db.list_sessions_rich( source=source.platform.value if source.platform else None, - session_key=None if widen else session_key, limit=10, - ) + session_key=None if widen else session_key, limit=10) titled = [s for s in sessions if s.get("title")][:10] return [s for s in titled if await self._resume_row_visible(source, s, allow_all)] @@ -942,8 +897,7 @@ class GatewaySessionCommandsMixin: return t("gateway.resume.matrix_blocked_no_origin", name=name) return t( "gateway.resume.matrix_blocked_other_room", - room=target_origin.chat_name or target_origin.chat_id, name=name, - ) + room=target_origin.chat_name or target_origin.chat_id, name=name) if await self._resume_target_allowed(source, target_id, allow_override=(allow_all or allow_cross_room)): return None return t("gateway.resume.blocked_not_owner", name=name) @@ -995,8 +949,7 @@ class GatewaySessionCommandsMixin: msg_part = f" ({msg_count} message{'s' if msg_count != 1 else ''})" if msg_count else "" return t( "gateway.resume.matrix_cross_room_success", title=title, - room=source.chat_name or source.chat_id, msg_part=msg_part, - ) + room=source.chat_name or source.chat_id, msg_part=msg_part) if not msg_count: return t("gateway.resume.resumed_no_count", title=title) if msg_count == 1: @@ -1036,13 +989,11 @@ class GatewaySessionCommandsMixin: from hermes_cli.session_listing import ( format_gateway_session_listing, parse_session_listing_args, - query_session_listing, - ) + query_session_listing) try: include_all, include_unnamed, target, search_query = parse_session_listing_args( - event.get_command_args().strip() - ) + event.get_command_args().strip()) except ValueError as exc: return t("gateway.resume.parse_error", error=exc) if search_query == "": @@ -1070,8 +1021,7 @@ class GatewaySessionCommandsMixin: search_query=search_query, # Search filters in SQL: over-fetch so origin-invisible matches don't consume the page. limit=50 if search_query else 10, - exclude_sources=["tool"], - ) + exclude_sources=["tool"]) if not cross_origin: rows = [row for row in rows if await self._resume_row_visible(source, row, allow_all=False)] rows = rows[:10] @@ -1124,8 +1074,7 @@ class GatewaySessionCommandsMixin: chat_type=source.chat_type, thread_id=source.thread_id, origin_json=_branch_origin_json, - display_name=current_entry.display_name, - ) + display_name=current_entry.display_name) except Exception as e: logger.error("Failed to create branch session: %s", e) return t("gateway.branch.create_failed", error=e) @@ -1133,8 +1082,7 @@ class GatewaySessionCommandsMixin: # Chunked transactions; best-effort — a failed copy still yields a usable (partial) branch. with contextlib.suppress(Exception): await self._session_db.append_messages_batch( - new_session_id, [_branch_row(msg) for msg in history], chunk_rows=500, - ) + new_session_id, [_branch_row(msg) for msg in history], chunk_rows=500) with contextlib.suppress(Exception): await self._session_db.set_session_title(new_session_id, branch_title) new_entry = await self.async_session_store.switch_session(session_key, new_session_id) diff --git a/gateway/stream_consumer.py b/gateway/stream_consumer.py index 46172c7de4..8918e0ad27 100644 --- a/gateway/stream_consumer.py +++ b/gateway/stream_consumer.py @@ -26,16 +26,13 @@ from gateway.platforms.base import _custom_unit_to_cp from gateway.config import ( DEFAULT_STREAMING_EDIT_INTERVAL as _DEFAULT_STREAMING_EDIT_INTERVAL, DEFAULT_STREAMING_BUFFER_THRESHOLD as _DEFAULT_STREAMING_BUFFER_THRESHOLD, - DEFAULT_STREAMING_CURSOR as _DEFAULT_STREAMING_CURSOR, -) + DEFAULT_STREAMING_CURSOR as _DEFAULT_STREAMING_CURSOR) from gateway.response_filters import ( is_intentional_silence_response as _is_intentional_silence_response, - is_partial_silence_marker as _is_partial_silence_marker, -) + is_partial_silence_marker as _is_partial_silence_marker) from gateway.stream_consumer_fences import ( # noqa: F401 (re-exported) ensure_closed_code_fences, - escape_code_fences_for_display, -) + escape_code_fences_for_display) from gateway.stream_consumer_transport import StreamTransportMixin from gateway.stream_consumer_fallback import StreamFallbackMixin from gateway.stream_consumer_think import StreamThinkFilterMixin @@ -120,8 +117,7 @@ class GatewayStreamConsumer(StreamTransportMixin, StreamFallbackMixin, StreamThi on_new_message: Optional[callable] = None, on_before_finalize: Optional[Callable[[], Any]] = None, initial_reply_to_id: Optional[str] = None, - run_still_current: Optional[Callable[[], bool]] = None, - ): + run_still_current: Optional[Callable[[], bool]] = None): self.adapter = adapter self.chat_id = chat_id self.cfg = config or StreamConsumerConfig() @@ -781,8 +777,7 @@ class GatewayStreamConsumer(StreamTransportMixin, StreamFallbackMixin, StreamThi # answer. tick.update_visible = await self._send_or_edit( display_text, finalize=tick.got_done or tick.got_segment_break, - is_turn_final=tick.got_done, - ) + is_turn_final=tick.got_done) self._last_edit_time = time.monotonic() # Lines stay in _tool_progress_lines for the next compose. self._tool_progress_active = False @@ -885,8 +880,7 @@ class GatewayStreamConsumer(StreamTransportMixin, StreamFallbackMixin, StreamThi if self._accumulated and self._message_id: with contextlib.suppress(Exception): best_effort_ok = bool(await self._send_or_edit( - self._accumulated, finalize=True, is_turn_final=False, - )) + self._accumulated, finalize=True, is_turn_final=False)) elif self._message_id is None: # Draft path keeps _message_id=None; seal in place (else the stream stays # visibly live and the adapter keeps armed interception state). diff --git a/gateway/stream_consumer_fallback.py b/gateway/stream_consumer_fallback.py index e0a2e4d4e9..1439277ef9 100644 --- a/gateway/stream_consumer_fallback.py +++ b/gateway/stream_consumer_fallback.py @@ -26,8 +26,7 @@ class StreamFallbackMixin: try: result = await self.adapter.send( chat_id=self.chat_id, content=text, reply_to=reply_to_id, - metadata=self._metadata_for_send(final=final, expect_edits=not final), - ) + metadata=self._metadata_for_send(final=final, expect_edits=not final)) if not (result.success and result.message_id): self._edit_supported = False return reply_to_id @@ -216,8 +215,7 @@ class StreamFallbackMixin: try: result = await self._send_with_flood_retry( content=final_text, reply_to=self._initial_reply_to_id, - retry_log="Flood control on empty fallback final send; retrying in %.1fs", - ) + retry_log="Flood control on empty fallback final send; retrying in %.1fs") except Exception as exc: logger.debug("Empty fallback final send failed: %s", exc) return "ambiguous" if self._send_failure_may_have_delivered(exc) else "failed" @@ -323,8 +321,7 @@ class StreamFallbackMixin: _needs_reply_anchor = _platform_name in ("buzz", "slack", "mattermost", "feishu") result = await self.adapter.send( chat_id=self.chat_id, content=text, - reply_to=self._initial_reply_to_id if _needs_reply_anchor else None, metadata=_md, - ) + reply_to=self._initial_reply_to_id if _needs_reply_anchor else None, metadata=_md) # Do NOT set _already_sent: commentary is interim, and the flag would # suppress the real final after multiple tool calls. if result.success: diff --git a/gateway/stream_consumer_transport.py b/gateway/stream_consumer_transport.py index 70121c75ed..2cfed0a686 100644 --- a/gateway/stream_consumer_transport.py +++ b/gateway/stream_consumer_transport.py @@ -31,8 +31,7 @@ class StreamTransportMixin: try: params = inspect.signature(self.adapter.edit_message).parameters if "metadata" in params or any( - param.kind is inspect.Parameter.VAR_KEYWORD for param in params.values() - ): + param.kind is inspect.Parameter.VAR_KEYWORD for param in params.values()): kwargs["metadata"] = self.metadata except (TypeError, ValueError): pass @@ -43,8 +42,7 @@ class StreamTransportMixin: bool; a raise logs ``fail_log`` at DEBUG (error formatted in, or the traceback when ``exc_info``) and reads as False.""" seed = self.adapter.send_stream_frame( - "", chat_id=self.chat_id, reply_to=self._initial_reply_to_id, turn_id=self._turn_id, - ) + "", chat_id=self.chat_id, reply_to=self._initial_reply_to_id, turn_id=self._turn_id) return await self._try_frame(seed, fail_log, exc_info=exc_info) @staticmethod @@ -63,8 +61,7 @@ class StreamTransportMixin: """One native-stream frame; every frame carries the same chat/reply/turn routing.""" return await self.adapter.send_stream_frame( text, finalize=finalize, chat_id=self.chat_id, reply_to=self._initial_reply_to_id, - turn_id=self._turn_id, - ) + turn_id=self._turn_id) def _close_native_state(self) -> None: """Mark the native stream closed (next content re-seeds or falls back).""" @@ -166,8 +163,7 @@ class StreamTransportMixin: try: result = await self.adapter.send_draft( chat_id=self.chat_id, draft_id=self._draft_id, content=text, - metadata=self._draft_metadata(), - ) + metadata=self._draft_metadata()) except Exception as e: logger.debug("send_draft raised, disabling draft transport for this run: %s", e) else: @@ -191,8 +187,7 @@ class StreamTransportMixin: try: await self.adapter.abandon_open_draft( self.chat_id, self._last_sent_text or self._clean_for_display(self._accumulated), - metadata=self._draft_metadata(), - ) + metadata=self._draft_metadata()) except Exception as e: logger.debug("abandon_open_draft failed (best-effort): %s", e) @@ -255,8 +250,7 @@ class StreamTransportMixin: stale_ids = self._stale_preview_ids() try: result = await self.adapter.send( - chat_id=self.chat_id, content=text, metadata=self._metadata_for_send(final=True), - ) + chat_id=self.chat_id, content=text, metadata=self._metadata_for_send(final=True)) except Exception as e: logger.debug("Fresh-final send failed, falling back to edit: %s", e) return False @@ -284,8 +278,7 @@ class StreamTransportMixin: self._message_created_ts = None async def _send_or_edit( - self, text: str, *, finalize: bool = False, is_turn_final: bool = True, - ) -> bool: + self, text: str, *, finalize: bool = False, is_turn_final: bool = True) -> bool: """Send or edit the streaming message; True if delivered. ``finalize`` marks the last edit. Transport order: native frame → draft frame → edit existing → first send; a transport returns None to fall through to the next.""" @@ -413,8 +406,7 @@ class StreamTransportMixin: """First send, threaded to the user's message (correct topic/thread).""" result = await self.adapter.send( chat_id=self.chat_id, content=text, reply_to=self._initial_reply_to_id, - metadata=self._metadata_for_send(final=finalize, expect_edits=not finalize), - ) + metadata=self._metadata_for_send(final=finalize, expect_edits=not finalize)) if not result.success: self._edit_supported = False return False @@ -443,8 +435,7 @@ class StreamTransportMixin: # CLASS (MagicMock auto-creates attrs) plus instance __dict__ (test doubles). has_prefers_hook = ( hasattr(type(self.adapter), "prefers_fresh_final_streaming") - or "prefers_fresh_final_streaming" in getattr(self.adapter, "__dict__", {}) - ) + or "prefers_fresh_final_streaming" in getattr(self.adapter, "__dict__", {})) prefers_fresh = self._adapter_prefers_fresh_final(text) # probed every edit (hook contract) if finalize and ( prefers_fresh or (not has_prefers_hook and self._should_send_fresh_final()) @@ -518,8 +509,7 @@ class StreamTransportMixin: logger.debug("Flood control on edit (strike %d/%d), backoff interval → %.1fs", self._flood_strikes, self._MAX_FLOOD_STRIKES, self._current_edit_interval) immediate_final_fallback = ( - turn_final and getattr(self.adapter, "FALLBACK_ON_FINAL_EDIT_FLOOD", False) is True - ) + turn_final and getattr(self.adapter, "FALLBACK_ON_FINAL_EDIT_FLOOD", False) is True) if self._flood_strikes < self._MAX_FLOOD_STRIKES and not immediate_final_fallback: self._last_edit_time = time.monotonic() # honor the new interval return False