refactor(cli): split HermesCLI into 10 cohesive mixins (cli.py 22,284 -> 9,150)
326 methods lifted by AST (bodies identical; ast.dump-verified) into
hermes_cli/cli_{tui,status_bar,voice,model_switch,session,stream,modal,
terminal,info,loops}_mixin.py. cli.py-internal symbols resolve via lazy
'from cli import ...' inside each method (no import cycle; patch('cli.X')
keeps working). The three 'global' writers (_skill_commands, _cli_wake_owner)
now write the cli module attribute explicitly so the origin's readers still
see them. Dropped imports left unused in cli.py; kept display_hermes_home /
build_welcome_banner as re-exports (mixins + tests resolve them via cli).
Repointed two AST change-detector tests to cli_tui_mixin.py; one test
fixture now keeps 'cli' in sys.modules across its patch.dict scope.
This commit is contained in:
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,725 @@
|
||||
"""Simple slash-command wrappers plus goal/heartbeat/loop manager hooks for the interactive CLI
|
||||
|
||||
Mixin split out of ``cli.py``; bound onto ``HermesCLI`` via the MRO. cli.py-internal
|
||||
symbols are imported LAZILY inside each method (``from cli import ...``) — the mixin
|
||||
never imports ``cli`` at module load time (import cycle).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import shutil
|
||||
import threading
|
||||
import time
|
||||
|
||||
from rich.markup import escape as _escape
|
||||
|
||||
|
||||
class CLILoopsMixin:
|
||||
"""Simple slash-command wrappers plus goal/heartbeat/loop manager hooks for the interactive CLI"""
|
||||
|
||||
def _cmd_exit(self, cmd_original: str):
|
||||
# /exit --delete also removes the session's transcripts + SQLite history.
|
||||
from cli import _DIM, _RST, _cprint, _slash_args
|
||||
_args = _slash_args(cmd_original).lower()
|
||||
if _args in {"--delete", "-d"}:
|
||||
self._delete_session_on_exit = True
|
||||
elif _args:
|
||||
_cprint(f" {_DIM}✗ Unknown argument: {_escape(_args)}. Use /exit --delete to also remove session history.{_RST}")
|
||||
return True
|
||||
return False
|
||||
|
||||
def _cmd_help(self, cmd_original: str):
|
||||
from cli import _slash_args
|
||||
self.show_help(_slash_args(cmd_original))
|
||||
|
||||
def _cmd_redraw(self, cmd_original: str):
|
||||
# Manual recovery for terminal buffer drift from multiplexer
|
||||
# tab switches, subshell ``clear``, SSH window restores, etc.
|
||||
# See issue #8688 (cmux). Ctrl+L is bound to the same helper.
|
||||
from cli import _DIM, _RST, _cprint
|
||||
self._force_full_redraw()
|
||||
_cprint(f" {_DIM}✓ UI redrawn{_RST}")
|
||||
|
||||
def _cmd_clear(self, cmd_original: str):
|
||||
from cli import (
|
||||
ChatConsole,
|
||||
_build_compact_banner,
|
||||
_clear_output_history,
|
||||
_cprint,
|
||||
build_welcome_banner,
|
||||
get_tool_definitions,
|
||||
)
|
||||
if self._confirm_destructive_slash(
|
||||
"clear",
|
||||
"This clears the screen and starts a new session.\n"
|
||||
"The current conversation history will be discarded.",
|
||||
cmd_original=cmd_original,
|
||||
) is None:
|
||||
return True # confirmation cancelled — command handled, keep REPL alive
|
||||
self.new_session(silent=True)
|
||||
_clear_output_history()
|
||||
# Clear terminal screen. Inside the TUI, Rich's console.clear()
|
||||
# goes through patch_stdout's StdoutProxy which swallows the
|
||||
# screen-clear escape sequences. Use prompt_toolkit's output
|
||||
# object directly to actually clear the terminal.
|
||||
if self._app:
|
||||
out = self._app.output
|
||||
out.erase_screen()
|
||||
out.cursor_goto(0, 0)
|
||||
out.flush()
|
||||
else:
|
||||
self.console.clear()
|
||||
# Show fresh banner. Inside the TUI we must route Rich output
|
||||
# through ChatConsole (which uses prompt_toolkit's native ANSI
|
||||
# renderer) instead of self.console (which writes raw to stdout
|
||||
# and gets mangled by patch_stdout).
|
||||
if self._app:
|
||||
cc = ChatConsole()
|
||||
term_w = shutil.get_terminal_size().columns
|
||||
if self.compact or term_w < 80:
|
||||
cc.print(_build_compact_banner())
|
||||
else:
|
||||
tools = get_tool_definitions(enabled_toolsets=self.enabled_toolsets, quiet_mode=True)
|
||||
cwd = os.getenv("TERMINAL_CWD", os.getcwd())
|
||||
ctx_len = None
|
||||
if hasattr(self, 'agent') and self.agent and hasattr(self.agent, 'context_compressor'):
|
||||
ctx_len = self.agent.context_compressor.context_length
|
||||
build_welcome_banner(
|
||||
console=cc,
|
||||
model=self.model,
|
||||
cwd=cwd,
|
||||
tools=tools,
|
||||
enabled_toolsets=self.enabled_toolsets,
|
||||
session_id=self.session_id,
|
||||
context_length=ctx_len,
|
||||
provider=self.provider,
|
||||
)
|
||||
_cprint(" ✨ (◕‿◕)✨ Fresh start! Screen cleared and conversation reset.\n")
|
||||
self._print_random_tip()
|
||||
else:
|
||||
self.show_banner()
|
||||
print(" ✨ (◕‿◕)✨ Fresh start! Screen cleared and conversation reset.\n")
|
||||
self._print_random_tip()
|
||||
|
||||
def _cmd_title(self, cmd_original: str):
|
||||
from cli import _cprint
|
||||
parts = cmd_original.split(maxsplit=1)
|
||||
if len(parts) > 1:
|
||||
raw_title = parts[1].strip()
|
||||
if raw_title:
|
||||
if self._session_db:
|
||||
# Sanitize the title early so feedback matches what gets stored
|
||||
try:
|
||||
from hermes_state import SessionDB
|
||||
new_title = SessionDB.sanitize_title(raw_title)
|
||||
except ValueError as e:
|
||||
# sanitize_title rejected the input (e.g. too long).
|
||||
# Print that one reason and stop — don't fall
|
||||
# through to the "empty after cleanup" branch and
|
||||
# print a second, contradictory error (SC-05).
|
||||
_cprint(f" {e}")
|
||||
return True
|
||||
if not new_title:
|
||||
_cprint(" Title is empty after cleanup. Please use printable characters.")
|
||||
elif self._session_db.get_session(self.session_id):
|
||||
# Session exists in DB — set title directly
|
||||
try:
|
||||
if self._session_db.set_session_title(self.session_id, new_title):
|
||||
self._status_bar_title_checked_at = 0.0
|
||||
_cprint(f" Session title set: {new_title}")
|
||||
else:
|
||||
_cprint(" Session not found in database.")
|
||||
except ValueError as e:
|
||||
_cprint(f" {e}")
|
||||
else:
|
||||
# Session not created yet — defer the title
|
||||
# Check uniqueness proactively with the sanitized title
|
||||
existing = self._session_db.get_session_by_title(new_title)
|
||||
if existing:
|
||||
_cprint(f" Title '{new_title}' is already in use by session {existing['id']}")
|
||||
else:
|
||||
self._pending_title = new_title
|
||||
_cprint(f" Session title queued: {new_title} (will be saved on first message)")
|
||||
else:
|
||||
from hermes_state import format_session_db_unavailable
|
||||
_cprint(f" {format_session_db_unavailable()}")
|
||||
else:
|
||||
_cprint(" Usage: /title <your session title>")
|
||||
# Show current title and session ID if no argument given
|
||||
elif self._session_db:
|
||||
_cprint(f" Session ID: {self.session_id}")
|
||||
session = self._session_db.get_session(self.session_id)
|
||||
if session and session.get("title"):
|
||||
_cprint(f" Title: {session['title']}")
|
||||
elif self._pending_title:
|
||||
_cprint(f" Title (pending): {self._pending_title}")
|
||||
else:
|
||||
_cprint(" No title set. Usage: /title <your session title>")
|
||||
else:
|
||||
from hermes_state import format_session_db_unavailable
|
||||
_cprint(f" {format_session_db_unavailable()}")
|
||||
|
||||
def _cmd_new(self, cmd_original: str):
|
||||
# Strip inline-skip tokens (now/--yes/-y) before deriving the title
|
||||
# so "/new now My Session" yields title="My Session" instead of
|
||||
# title="now My Session". See _split_destructive_skip.
|
||||
_new_args, _ = self._split_destructive_skip(cmd_original)
|
||||
title = _new_args.strip() or None
|
||||
if self._confirm_destructive_slash(
|
||||
"new",
|
||||
"This starts a fresh session.\n"
|
||||
"The current conversation history will be discarded.",
|
||||
cmd_original=cmd_original,
|
||||
) is None:
|
||||
return True # confirmation cancelled — command handled, keep REPL alive
|
||||
self.new_session(title=title)
|
||||
|
||||
def _cmd_retry(self, cmd_original: str):
|
||||
retry_msg = self.retry_last()
|
||||
if retry_msg and hasattr(self, '_pending_input'):
|
||||
# Re-queue the message so process_loop sends it to the agent
|
||||
self._pending_input.put(retry_msg)
|
||||
|
||||
def _cmd_undo(self, cmd_original: str):
|
||||
# Parse optional turn count: "/undo" → 1, "/undo 3" → 3.
|
||||
_undo_n = 1
|
||||
_undo_parts = cmd_original.split()
|
||||
if len(_undo_parts) > 1:
|
||||
try:
|
||||
_undo_n = int(_undo_parts[1])
|
||||
except ValueError:
|
||||
print(f"(._.) Invalid count {_undo_parts[1]!r} — use /undo or /undo N.")
|
||||
return True # bad arg — command handled, keep the REPL alive
|
||||
if _undo_n < 1:
|
||||
_undo_n = 1
|
||||
# Nothing to undo → say so immediately; don't pop a destructive
|
||||
# confirmation dialog for a guaranteed no-op (SC-06).
|
||||
if not self.conversation_history:
|
||||
print("(._.) No messages to undo.")
|
||||
return True
|
||||
_undo_desc = (
|
||||
"This removes the last user/assistant exchange from history."
|
||||
if _undo_n == 1
|
||||
else f"This removes the last {_undo_n} user turns from history."
|
||||
)
|
||||
if self._confirm_destructive_slash(
|
||||
"undo",
|
||||
_undo_desc,
|
||||
cmd_original=cmd_original,
|
||||
) is None:
|
||||
return True # confirmation cancelled — command handled, keep REPL alive
|
||||
self.undo_last(_undo_n)
|
||||
|
||||
def _cmd_skills(self, cmd_original: str):
|
||||
with self._busy_command(self._slow_command_status(cmd_original)):
|
||||
self._handle_skills_command(cmd_original)
|
||||
|
||||
def _cmd_egress(self, cmd_original: str):
|
||||
from hermes_cli.slash_exec import CommandContext, execute_command
|
||||
|
||||
self._console_print(
|
||||
execute_command("egress", CommandContext(surface="cli")).text,
|
||||
highlight=False, markup=False,
|
||||
)
|
||||
|
||||
def _cmd_statusbar(self, cmd_original: str):
|
||||
self._status_bar_visible = not self._status_bar_visible
|
||||
state = "visible" if self._status_bar_visible else "hidden"
|
||||
self._console_print(f" Status bar {state}")
|
||||
|
||||
def _cmd_update(self, cmd_original: str) -> bool:
|
||||
# A truthy result means the process is relaunching — leave the REPL.
|
||||
return not self._handle_update_command()
|
||||
|
||||
def _cmd_version(self, cmd_original: str):
|
||||
from hermes_cli.main import _print_version_info
|
||||
|
||||
_print_version_info(check_updates=True)
|
||||
|
||||
def _cmd_reload(self, cmd_original: str):
|
||||
from hermes_cli.config import reload_env
|
||||
count = reload_env()
|
||||
print(f" Reloaded .env ({count} var(s) updated)")
|
||||
|
||||
def _cmd_reload_skills(self, cmd_original: str):
|
||||
with self._busy_command(self._slow_command_status(cmd_original)):
|
||||
self._reload_skills()
|
||||
|
||||
def _cmd_plugins(self, cmd_original: str):
|
||||
from cli import display_hermes_home
|
||||
try:
|
||||
# Discover from disk (bundled + user), matching `hermes plugins
|
||||
# list` — so installed-but-not-enabled plugins are visible here
|
||||
# too. The plugin manager only knows about *loaded* plugins, so
|
||||
# using it alone made freshly-installed, not-yet-enabled plugins
|
||||
# look like "nothing installed".
|
||||
from hermes_cli.plugins_cmd import (
|
||||
_discover_all_plugins,
|
||||
_get_disabled_set,
|
||||
_get_enabled_set,
|
||||
_plugin_status,
|
||||
)
|
||||
|
||||
entries = _discover_all_plugins()
|
||||
enabled = _get_enabled_set()
|
||||
disabled = _get_disabled_set()
|
||||
|
||||
# `/plugins` is a quick glance — default to user-installed
|
||||
# plugins (what the user actually added). Bundled provider/
|
||||
# platform plugins are summarized on one line; the full
|
||||
# catalog lives behind `hermes plugins list`.
|
||||
user_entries = [e for e in entries if e[3] != "bundled"]
|
||||
bundled_count = len(entries) - len(user_entries)
|
||||
|
||||
if not user_entries:
|
||||
print("No user plugins installed.")
|
||||
print(" Install one: hermes plugins install owner/repo")
|
||||
print(f" Or drop a plugin directory into {display_hermes_home()}/plugins/")
|
||||
if bundled_count:
|
||||
print(f" ({bundled_count} bundled plugins available — see: hermes plugins list)")
|
||||
else:
|
||||
# Loaded-plugin details (tools/hooks/commands counts, errors)
|
||||
# keyed by name, when available.
|
||||
loaded: dict = {}
|
||||
try:
|
||||
from hermes_cli.plugins import get_plugin_manager
|
||||
for p in get_plugin_manager().list_plugins():
|
||||
loaded[p["name"]] = p
|
||||
except Exception:
|
||||
loaded = {}
|
||||
|
||||
print(f"User plugins ({len(user_entries)}):")
|
||||
for name, version, _desc, source, _dir, key in sorted(user_entries):
|
||||
state = _plugin_status(name, enabled, disabled, key=key)
|
||||
glyph = {"enabled": "✓", "disabled": "✗"}.get(state, "○")
|
||||
ver = f" v{version}" if version else ""
|
||||
info = loaded.get(name) or {}
|
||||
bits = []
|
||||
if info.get("tools"):
|
||||
bits.append(f"{info['tools']} tools")
|
||||
if info.get("hooks"):
|
||||
bits.append(f"{info['hooks']} hooks")
|
||||
if info.get("commands"):
|
||||
bits.append(f"{info['commands']} commands")
|
||||
detail = f" ({', '.join(bits)})" if bits else ""
|
||||
label = "" if state == "enabled" else f" [{state}]"
|
||||
error = f" — {info['error']}" if info.get("error") else ""
|
||||
print(f" {glyph} {name}{ver}{label}{detail}{error}")
|
||||
if bundled_count:
|
||||
print(f" (+{bundled_count} bundled — see: hermes plugins list)")
|
||||
print(" Enable/disable: hermes plugins enable/disable <name>")
|
||||
except Exception as e:
|
||||
print(f"Plugin system error: {e}")
|
||||
|
||||
def _cmd_queue(self, cmd_original: str):
|
||||
from cli import _cprint, _slash_args
|
||||
payload = self._expand_paste_references(_slash_args(cmd_original))
|
||||
if not payload:
|
||||
_cprint(" Usage: /queue <prompt>")
|
||||
else:
|
||||
self._pending_input.put(payload)
|
||||
if self._agent_running:
|
||||
_cprint(f" Queued for the next turn: {payload[:80]}{'...' if len(payload) > 80 else ''}")
|
||||
else:
|
||||
_cprint(f" Queued: {payload[:80]}{'...' if len(payload) > 80 else ''}")
|
||||
|
||||
def _cmd_steer(self, cmd_original: str):
|
||||
# Inject a message after the next tool call without interrupting.
|
||||
# If the agent is actively running, push the text into the agent's
|
||||
# pending_steer slot — the drain hook in _execute_tool_calls_*
|
||||
# will append it to the next tool result's content. If no agent
|
||||
# is running, fall back to queue semantics (same as /queue).
|
||||
from cli import _cprint, _slash_args
|
||||
payload = _slash_args(cmd_original)
|
||||
if not payload:
|
||||
_cprint(" Usage: /steer <prompt>")
|
||||
elif self._agent_running and self.agent is not None and hasattr(self.agent, "steer"):
|
||||
try:
|
||||
accepted = self.agent.steer(payload)
|
||||
except Exception as exc:
|
||||
_cprint(f" Steer failed: {exc}")
|
||||
else:
|
||||
if accepted:
|
||||
_cprint(f" ⏩ Steer queued — arrives after the next tool call: {payload[:80]}{'...' if len(payload) > 80 else ''}")
|
||||
else:
|
||||
_cprint(" Steer rejected (empty payload).")
|
||||
else:
|
||||
# No active run — treat as a normal next-turn message.
|
||||
self._pending_input.put(payload)
|
||||
_cprint(f" No agent running; queued as next turn: {payload[:80]}{'...' if len(payload) > 80 else ''}")
|
||||
|
||||
# ────────────────────────────────────────────────────────────────
|
||||
# /goal — persistent cross-turn goals (Ralph-style loop)
|
||||
# ────────────────────────────────────────────────────────────────
|
||||
def _get_goal_manager(self):
|
||||
"""Return the GoalManager bound to the current session_id.
|
||||
|
||||
Cached on ``self._goal_manager`` and rebound lazily when
|
||||
``session_id`` changes (e.g. after /new or a compression-driven
|
||||
session split).
|
||||
"""
|
||||
try:
|
||||
from hermes_cli.goals import GoalManager
|
||||
from hermes_cli.config import load_config
|
||||
except Exception as exc:
|
||||
logging.debug("goal manager unavailable: %s", exc)
|
||||
return None
|
||||
|
||||
sid = getattr(self, "session_id", None) or ""
|
||||
if not sid:
|
||||
return None
|
||||
|
||||
existing = getattr(self, "_goal_manager", None)
|
||||
if existing is not None and getattr(existing, "session_id", None) == sid:
|
||||
return existing
|
||||
|
||||
try:
|
||||
cfg = load_config() or {}
|
||||
goals_cfg = cfg.get("goals") or {}
|
||||
max_turns = int(goals_cfg.get("max_turns", 20) or 20)
|
||||
except Exception:
|
||||
max_turns = 20
|
||||
|
||||
mgr = GoalManager(session_id=sid, default_max_turns=max_turns)
|
||||
self._goal_manager = mgr
|
||||
return mgr
|
||||
|
||||
def _get_heartbeat_manager(self):
|
||||
"""Return the HeartbeatManager bound to the current session_id.
|
||||
|
||||
Cached on ``self._heartbeat_manager`` and rebound lazily when
|
||||
``session_id`` changes (mirrors ``_get_goal_manager``).
|
||||
"""
|
||||
try:
|
||||
from hermes_cli.heartbeat import HeartbeatManager
|
||||
except Exception as exc:
|
||||
logging.debug("heartbeat manager unavailable: %s", exc)
|
||||
return None
|
||||
|
||||
sid = getattr(self, "session_id", None) or ""
|
||||
if not sid:
|
||||
return None
|
||||
|
||||
existing = getattr(self, "_heartbeat_manager", None)
|
||||
if existing is not None and getattr(existing, "session_id", None) == sid:
|
||||
return existing
|
||||
|
||||
mgr = HeartbeatManager(session_id=sid)
|
||||
self._heartbeat_manager = mgr
|
||||
return mgr
|
||||
|
||||
def _start_heartbeat_watchdog(self):
|
||||
"""Start the idle-poll thread that fires due heartbeats.
|
||||
|
||||
Same pattern as the wake-word watchdog: a daemon thread polls a few
|
||||
times a minute; when the session is idle (no agent running, empty
|
||||
input queue) and the heartbeat is due, its prompt is injected into
|
||||
``_pending_input`` as a normal user turn. Missed ticks coalesce —
|
||||
the anchor resets on fire, so a busy hour yields ONE heartbeat turn,
|
||||
not a backlog. Idempotent; safe to call on every /heartbeat set.
|
||||
"""
|
||||
if getattr(self, "_heartbeat_watchdog_started", False):
|
||||
return
|
||||
self._heartbeat_watchdog_started = True
|
||||
|
||||
from hermes_cli.heartbeat import POLL_SECONDS
|
||||
|
||||
def _loop():
|
||||
try:
|
||||
while not getattr(self, "_should_exit", False):
|
||||
time.sleep(POLL_SECONDS)
|
||||
try:
|
||||
mgr = self._get_heartbeat_manager()
|
||||
if mgr is None or not mgr.is_active():
|
||||
continue
|
||||
busy = (
|
||||
self._agent_running
|
||||
or getattr(self, "_voice_recording", False)
|
||||
or getattr(self, "_voice_processing", False)
|
||||
or not self._pending_input.empty()
|
||||
)
|
||||
if busy:
|
||||
continue
|
||||
prompt = mgr.due_prompt()
|
||||
if prompt:
|
||||
self._pending_input.put(prompt)
|
||||
except Exception as exc:
|
||||
logging.debug("heartbeat watchdog tick failed: %s", exc)
|
||||
finally:
|
||||
self._heartbeat_watchdog_started = False
|
||||
|
||||
threading.Thread(target=_loop, daemon=True, name="heartbeat-watchdog").start()
|
||||
|
||||
# ────────────────────────────────────────────────────────────────
|
||||
# /loop — recurring in-session wakeups (Claude Code /loop parity)
|
||||
# ────────────────────────────────────────────────────────────────
|
||||
def _get_loop_manager(self):
|
||||
"""Return the LoopManager bound to the current session_id.
|
||||
|
||||
Cached on ``self._loop_manager`` and rebound lazily when
|
||||
``session_id`` changes (mirrors ``_get_goal_manager``).
|
||||
"""
|
||||
try:
|
||||
from hermes_cli.loops import LoopManager
|
||||
except Exception as exc:
|
||||
logging.debug("loop manager unavailable: %s", exc)
|
||||
return None
|
||||
|
||||
sid = getattr(self, "session_id", None) or ""
|
||||
if not sid:
|
||||
return None
|
||||
|
||||
existing = getattr(self, "_loop_manager", None)
|
||||
if existing is not None and getattr(existing, "session_id", None) == sid:
|
||||
return existing
|
||||
|
||||
mgr = LoopManager(session_id=sid)
|
||||
self._loop_manager = mgr
|
||||
return mgr
|
||||
|
||||
def _maybe_fire_loop_tick(self) -> None:
|
||||
"""Idle hook run from process_loop: fire a due /loop wakeup.
|
||||
|
||||
Only runs while the agent is idle and nothing is queued — a real
|
||||
user message always wins the idle boundary. An active (non-parked)
|
||||
/goal also wins: its judge-driven continuations own the idle
|
||||
boundary, so the loop defers to the next poll.
|
||||
"""
|
||||
from cli import _DIM, _RST, _cprint
|
||||
mgr = self._get_loop_manager()
|
||||
if mgr is None or not mgr.is_due():
|
||||
return
|
||||
# The idle poll runs at ~10 Hz; once a tick is due but deferred
|
||||
# (queued input / active goal), every poll would otherwise hit the
|
||||
# DB via goal_blocks_loop_tick. Throttle the deferred re-check.
|
||||
now = time.time()
|
||||
if now - getattr(self, "_last_loop_tick_check", 0.0) < 2.0:
|
||||
return
|
||||
self._last_loop_tick_check = now
|
||||
# Real user input (or anything else queued) takes priority; the
|
||||
# loop stays due and fires at the next idle poll.
|
||||
try:
|
||||
if not self._pending_input.empty():
|
||||
return
|
||||
except Exception:
|
||||
return
|
||||
try:
|
||||
from hermes_cli.loops import goal_blocks_loop_tick
|
||||
|
||||
if goal_blocks_loop_tick(mgr.session_id):
|
||||
return
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
wakeup = mgr.fire_tick()
|
||||
if not wakeup:
|
||||
return
|
||||
try:
|
||||
state = mgr.state
|
||||
tick_no = state.ticks_fired if state else "?"
|
||||
_cprint(f" {_DIM}↻ /loop wakeup #{tick_no} firing…{_RST}")
|
||||
self._pending_input.put(wakeup)
|
||||
except Exception as exc:
|
||||
logging.debug("loop tick injection failed: %s", exc)
|
||||
try:
|
||||
mgr.abandon_tick()
|
||||
except Exception:
|
||||
pass
|
||||
return
|
||||
# A slash-command loop (e.g. `/loop 10m /recap`) is dispatched via
|
||||
# process_command, which never reaches the post-turn chat() finally
|
||||
# block — so the tick would never complete and the loop would wedge
|
||||
# on awaiting_response. Slash ticks have no model reply to evaluate;
|
||||
# complete them immediately (caps and scheduling still apply).
|
||||
if wakeup.lstrip().startswith("/"):
|
||||
try:
|
||||
decision = mgr.complete_tick("")
|
||||
msg = decision.get("message") or ""
|
||||
if msg:
|
||||
_cprint(f" {msg}")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def _maybe_complete_loop_tick_after_turn(self) -> None:
|
||||
"""Post-turn hook: evaluate a finished /loop wakeup turn.
|
||||
|
||||
No-op unless the turn that just ended was a loop wakeup
|
||||
(``awaiting_response`` set by ``fire_tick``). Detects the
|
||||
LOOP_COMPLETE marker, judges --until, applies caps, and schedules
|
||||
the next tick. Mirrors _maybe_continue_goal_after_turn's shape.
|
||||
"""
|
||||
from cli import _DIM, _RST, _cprint
|
||||
mgr = self._get_loop_manager()
|
||||
if mgr is None:
|
||||
return
|
||||
state = mgr.state
|
||||
if state is None or not state.awaiting_response:
|
||||
return
|
||||
|
||||
# A user-interrupted wakeup turn pauses the loop (recoverable via
|
||||
# /loop resume) — same contract as the goal loop's Ctrl+C handling.
|
||||
if getattr(self, "_last_turn_interrupted", False):
|
||||
try:
|
||||
mgr.pause(reason="user-interrupted (Ctrl+C)")
|
||||
except Exception:
|
||||
pass
|
||||
_cprint(
|
||||
f" {_DIM}⏸ Loop paused — wakeup turn was interrupted. "
|
||||
f"Use /loop resume to continue, or /loop stop to end it.{_RST}"
|
||||
)
|
||||
return
|
||||
|
||||
last_response = ""
|
||||
try:
|
||||
hist = self.conversation_history or []
|
||||
for msg in reversed(hist):
|
||||
if msg.get("role") == "assistant":
|
||||
content = msg.get("content", "")
|
||||
if isinstance(content, list):
|
||||
parts = [
|
||||
p.get("text", "")
|
||||
for p in content
|
||||
if isinstance(p, dict) and p.get("type") in {"text", "output_text"}
|
||||
]
|
||||
last_response = "\n".join(t for t in parts if t)
|
||||
else:
|
||||
last_response = str(content or "")
|
||||
break
|
||||
except Exception:
|
||||
last_response = ""
|
||||
|
||||
decision = mgr.complete_tick(last_response)
|
||||
msg = decision.get("message") or ""
|
||||
if msg:
|
||||
_cprint(f" {msg}")
|
||||
elif decision.get("status") == "active" and mgr.state is not None:
|
||||
_cprint(f" {_DIM}↻ Loop: {mgr.state.remaining_label()}.{_RST}")
|
||||
|
||||
def _maybe_continue_goal_after_turn(self) -> None:
|
||||
"""Hook run after every CLI turn. Judges + maybe re-queues.
|
||||
|
||||
Safe to call when no goal is set — returns quickly.
|
||||
|
||||
Preemption is automatic: if a real user message is already in
|
||||
``_pending_input`` we skip judging (the user's new input takes
|
||||
priority and we'll re-judge after that turn). If judge says done,
|
||||
mark it done and tell the user. If judge says continue and we're
|
||||
under budget, push the continuation prompt onto the queue.
|
||||
|
||||
Interrupt handling: if the turn was user-cancelled (Ctrl+C), we
|
||||
AUTO-PAUSE the goal instead of judging + re-queuing. Otherwise
|
||||
Ctrl+C feels like it did nothing — the judge runs on whatever
|
||||
partial output landed, almost always says "continue", and the
|
||||
loop keeps going. Auto-pause keeps the goal recoverable via
|
||||
``/goal resume`` once the user has sorted out what they want.
|
||||
The empty-response skip mirrors the gateway guard at
|
||||
``_handle_message`` in ``gateway/run.py``.
|
||||
"""
|
||||
from cli import _DIM, _RST, _cprint, _looks_like_slash_command
|
||||
mgr = self._get_goal_manager()
|
||||
if mgr is None or not mgr.is_active():
|
||||
return
|
||||
|
||||
# If a real user message is already queued, don't inject a
|
||||
# continuation prompt on top — let the user's turn go first.
|
||||
# Slash commands don't count as "real user messages" for this
|
||||
# check: they're inspection/mutation (e.g. /subgoal added mid-
|
||||
# run) and the process_loop dispatches them via process_command,
|
||||
# not via chat(). If we treat a queued /subgoal as preempting,
|
||||
# the goal loop silently stalls — we'd return here, then the
|
||||
# slash command consumes its queue slot via process_command()
|
||||
# which never re-fires the goal hook. Peek at all queued entries
|
||||
# and only defer when there's a non-slash payload.
|
||||
try:
|
||||
pending = getattr(self, "_pending_input", None)
|
||||
if pending is not None and not pending.empty():
|
||||
has_real_message = False
|
||||
try:
|
||||
# Queue.queue is the underlying deque — direct peek
|
||||
# without disturbing FIFO order.
|
||||
for entry in list(pending.queue):
|
||||
# Bundled payloads are (text, images) tuples;
|
||||
# unpack for inspection.
|
||||
if isinstance(entry, tuple) and entry:
|
||||
entry = entry[0]
|
||||
if isinstance(entry, str) and _looks_like_slash_command(entry):
|
||||
continue
|
||||
has_real_message = True
|
||||
break
|
||||
except Exception:
|
||||
# Fallback: if we can't introspect the queue, behave
|
||||
# like the old check and defer to be safe.
|
||||
has_real_message = True
|
||||
if has_real_message:
|
||||
return
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# If the turn was user-interrupted (Ctrl+C), auto-pause the goal
|
||||
# and bail. The judge call would almost always return "continue"
|
||||
# on the partial output and immediately re-queue another turn,
|
||||
# which is exactly what the user cancelled. Pausing (rather than
|
||||
# silently skipping) is the observable, recoverable behavior.
|
||||
if getattr(self, "_last_turn_interrupted", False):
|
||||
try:
|
||||
mgr.pause(reason="user-interrupted (Ctrl+C)")
|
||||
except Exception as exc:
|
||||
logging.debug("goal pause-on-interrupt failed: %s", exc)
|
||||
_cprint(
|
||||
f" {_DIM}⏸ Goal paused — turn was interrupted. "
|
||||
f"Use /goal resume to continue, or /goal clear to stop.{_RST}"
|
||||
)
|
||||
return
|
||||
|
||||
# Extract the agent's final response for this turn.
|
||||
last_response = ""
|
||||
try:
|
||||
hist = self.conversation_history or []
|
||||
for msg in reversed(hist):
|
||||
if msg.get("role") == "assistant":
|
||||
content = msg.get("content", "")
|
||||
if isinstance(content, list):
|
||||
# Multimodal content — flatten text parts.
|
||||
parts = [
|
||||
p.get("text", "")
|
||||
for p in content
|
||||
if isinstance(p, dict) and p.get("type") in {"text", "output_text"}
|
||||
]
|
||||
last_response = "\n".join(t for t in parts if t)
|
||||
else:
|
||||
last_response = str(content or "")
|
||||
break
|
||||
except Exception:
|
||||
last_response = ""
|
||||
|
||||
# Skip judging on empty/whitespace-only responses. These are almost
|
||||
# always transient failures (API error, empty stream) where the
|
||||
# judge would say "continue" and trip the consecutive-parse-failures
|
||||
# backstop unnecessarily. Mirrors the gateway guard.
|
||||
if not last_response.strip():
|
||||
return
|
||||
|
||||
try:
|
||||
from hermes_cli.goals import gather_background_processes as _gather_bg
|
||||
_bg_procs = _gather_bg()
|
||||
except Exception:
|
||||
_bg_procs = None
|
||||
|
||||
decision = mgr.evaluate_after_turn(
|
||||
last_response,
|
||||
user_initiated=True,
|
||||
background_processes=_bg_procs,
|
||||
)
|
||||
msg = decision.get("message") or ""
|
||||
if msg:
|
||||
_cprint(f" {msg}")
|
||||
|
||||
if decision.get("should_continue"):
|
||||
prompt = decision.get("continuation_prompt")
|
||||
if prompt:
|
||||
try:
|
||||
self._pending_input.put(prompt)
|
||||
except Exception as exc:
|
||||
logging.debug("goal continuation enqueue failed: %s", exc)
|
||||
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,986 @@
|
||||
"""Streaming output, reasoning preview, tool progress callbacks, and busy-command spinner for the interactive CLI
|
||||
|
||||
Mixin split out of ``cli.py``; bound onto ``HermesCLI`` via the MRO. cli.py-internal
|
||||
symbols are imported LAZILY inside each method (``from cli import ...``) — the mixin
|
||||
never imports ``cli`` at module load time (import cycle).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import re
|
||||
import shutil
|
||||
import textwrap
|
||||
import time
|
||||
|
||||
from contextlib import contextmanager
|
||||
from pathlib import Path
|
||||
from rich.markup import escape as _escape
|
||||
|
||||
|
||||
class CLIStreamMixin:
|
||||
"""Streaming output, reasoning preview, tool progress callbacks, and busy-command spinner for the interactive CLI"""
|
||||
|
||||
def _on_thinking(self, text: str) -> None:
|
||||
"""Called by agent when thinking starts/stops. Updates TUI spinner."""
|
||||
if not text:
|
||||
self._flush_reasoning_preview(force=True)
|
||||
self._spinner_text = text or ""
|
||||
self._tool_start_time = 0.0 # clear tool timer when switching to thinking
|
||||
self._invalidate()
|
||||
|
||||
def _on_notice(self, notice) -> None:
|
||||
"""Queue an out-of-band AgentNotice for rendering at the next clean boundary.
|
||||
|
||||
Notices fire from inside the agent turn (cold-start seed during _init_agent,
|
||||
per-turn _capture_credits after the API call) — printing immediately races the
|
||||
streaming response and the line gets buried behind the prompt (see _cprint's
|
||||
bg-thread caveat). So we QUEUE here and flush in _flush_credit_notices(), called
|
||||
right after run_conversation returns. Fail-soft: never break the turn.
|
||||
"""
|
||||
try:
|
||||
text = getattr(notice, "text", "") or ""
|
||||
if not text:
|
||||
return
|
||||
level = getattr(notice, "level", "info") or "info"
|
||||
if not hasattr(self, "_pending_credit_notices"):
|
||||
self._pending_credit_notices = []
|
||||
self._pending_credit_notices.append((level, text))
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def _flush_credit_notices(self) -> None:
|
||||
"""Print any queued credit notices as level-colored lines. Called at turn end
|
||||
(after run_conversation) where _cprint paints cleanly above the prompt."""
|
||||
from cli import _DIM, _RST, _cprint
|
||||
try:
|
||||
pending = getattr(self, "_pending_credit_notices", None)
|
||||
if not pending:
|
||||
return
|
||||
self._pending_credit_notices = []
|
||||
for level, text in pending:
|
||||
color = {
|
||||
"error": "\033[31m",
|
||||
"warn": "\033[33m",
|
||||
"success": "\033[32m",
|
||||
"info": _DIM,
|
||||
}.get(level, _DIM)
|
||||
_cprint(f" {color}{text}{_RST}")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def _on_notice_clear(self, key: str) -> None:
|
||||
"""Notice cleared. The REPL prints lines (no persistent slot to wipe), so
|
||||
this drops any still-queued notice with that key is not tracked by key here;
|
||||
it's a no-op for rendering — kept so the agent's clear callback is bound
|
||||
symmetrically with the show callback (and so future REPL UIs can hook it)."""
|
||||
return
|
||||
|
||||
def _current_reasoning_callback(self):
|
||||
"""Return the active reasoning display callback for the current mode."""
|
||||
if self.show_reasoning and self.streaming_enabled:
|
||||
return self._stream_reasoning_delta
|
||||
if self.verbose and not self.show_reasoning:
|
||||
return self._on_reasoning
|
||||
return None
|
||||
|
||||
def _emit_reasoning_preview(self, reasoning_text: str) -> None:
|
||||
"""Render a buffered reasoning preview as a single [thinking] block."""
|
||||
from cli import _DIM, _RST, _cprint
|
||||
preview_text = reasoning_text.strip()
|
||||
if not preview_text:
|
||||
return
|
||||
|
||||
try:
|
||||
term_width = shutil.get_terminal_size().columns
|
||||
except Exception:
|
||||
term_width = 80
|
||||
prefix = " [thinking] "
|
||||
wrap_width = max(30, term_width - len(prefix) - 2)
|
||||
|
||||
paragraphs = []
|
||||
raw_paragraphs = re.split(r"\n\s*\n+", preview_text.replace("\r\n", "\n"))
|
||||
for paragraph in raw_paragraphs:
|
||||
compact = " ".join(line.strip() for line in paragraph.splitlines() if line.strip())
|
||||
if compact:
|
||||
paragraphs.append(textwrap.fill(compact, width=wrap_width))
|
||||
preview_text = "\n".join(paragraphs)
|
||||
if not preview_text:
|
||||
return
|
||||
|
||||
if self.verbose:
|
||||
_cprint(f" {_DIM}[thinking] {preview_text}{_RST}")
|
||||
return
|
||||
|
||||
lines = preview_text.splitlines()
|
||||
if len(lines) > 5:
|
||||
preview = "\n".join(lines[:5])
|
||||
preview += f"\n ... ({len(lines) - 5} more lines)"
|
||||
else:
|
||||
preview = preview_text
|
||||
_cprint(f" {_DIM}[thinking] {preview}{_RST}")
|
||||
|
||||
def _flush_reasoning_preview(self, *, force: bool = False) -> None:
|
||||
"""Flush buffered reasoning text at natural boundaries.
|
||||
|
||||
Some providers stream reasoning in tiny word or punctuation chunks.
|
||||
Buffer them here so the preview path does not print one `[thinking]`
|
||||
line per token.
|
||||
"""
|
||||
buf = getattr(self, "_reasoning_preview_buf", "")
|
||||
if not buf:
|
||||
return
|
||||
|
||||
try:
|
||||
term_width = shutil.get_terminal_size().columns
|
||||
except Exception:
|
||||
term_width = 80
|
||||
target_width = max(40, term_width - len(" [thinking] ") - 4)
|
||||
|
||||
flush_text = ""
|
||||
|
||||
if force:
|
||||
flush_text = buf
|
||||
buf = ""
|
||||
else:
|
||||
line_break = buf.rfind("\n")
|
||||
min_newline_flush = max(16, target_width // 3)
|
||||
if line_break != -1 and (
|
||||
line_break >= min_newline_flush
|
||||
or buf.endswith("\n\n")
|
||||
or buf.endswith(".\n")
|
||||
or buf.endswith("!\n")
|
||||
or buf.endswith("?\n")
|
||||
or buf.endswith(":\n")
|
||||
):
|
||||
flush_text = buf[: line_break + 1]
|
||||
buf = buf[line_break + 1 :]
|
||||
elif len(buf) >= target_width:
|
||||
search_start = max(20, target_width // 2)
|
||||
search_end = min(len(buf), max(target_width + (target_width // 3), target_width + 8))
|
||||
cut = -1
|
||||
for boundary in (" ", "\t", ".", "!", "?", ",", ";", ":"):
|
||||
cut = max(cut, buf.rfind(boundary, search_start, search_end))
|
||||
if cut != -1:
|
||||
flush_text = buf[: cut + 1]
|
||||
buf = buf[cut + 1 :]
|
||||
|
||||
self._reasoning_preview_buf = buf.lstrip() if flush_text else buf
|
||||
if flush_text:
|
||||
self._emit_reasoning_preview(flush_text)
|
||||
|
||||
def _format_submitted_user_message_preview(self, user_input: str) -> str:
|
||||
"""Format the submitted user-message scrollback preview."""
|
||||
from cli import _accent_hex, datetime
|
||||
ts_suffix = (
|
||||
f" [dim]{datetime.now().strftime(getattr(self, 'timestamp_format', '%H:%M'))}[/]"
|
||||
if getattr(self, "show_timestamps", False) else ""
|
||||
)
|
||||
lines = user_input.split("\n")
|
||||
if len(lines) <= 1:
|
||||
return f"[bold {_accent_hex()}]●[/] [bold]{_escape(user_input)}[/]{ts_suffix}"
|
||||
|
||||
first_lines = int(getattr(self, "user_message_preview_first_lines", 2))
|
||||
last_lines = int(getattr(self, "user_message_preview_last_lines", 2))
|
||||
first_lines = max(1, first_lines)
|
||||
last_lines = max(0, last_lines)
|
||||
head = lines[:first_lines]
|
||||
remaining_after_head = max(0, len(lines) - len(head))
|
||||
tail_count = min(last_lines, remaining_after_head)
|
||||
tail = lines[-tail_count:] if tail_count else []
|
||||
|
||||
hidden_middle_count = len(lines) - len(head) - len(tail)
|
||||
if hidden_middle_count < 0:
|
||||
hidden_middle_count = 0
|
||||
tail = []
|
||||
|
||||
preview_lines = [
|
||||
f"[bold {_accent_hex()}]●[/] [bold]{_escape(head[0])}[/]{ts_suffix}"
|
||||
]
|
||||
preview_lines.extend(f"[bold]{_escape(line)}[/]" for line in head[1:])
|
||||
|
||||
if hidden_middle_count > 0:
|
||||
noun = "line" if hidden_middle_count == 1 else "lines"
|
||||
preview_lines.append(f"[dim]... (+{hidden_middle_count} more {noun})[/]")
|
||||
|
||||
preview_lines.extend(f"[bold]{_escape(line)}[/]" for line in tail)
|
||||
return "\n".join(preview_lines)
|
||||
|
||||
def _expand_paste_references(self, text: str | None) -> str:
|
||||
"""Expand [Pasted text #N -> file] placeholders into file contents."""
|
||||
from cli import logger
|
||||
if not isinstance(text, str) or "[Pasted text #" not in text:
|
||||
return text or ""
|
||||
paste_ref_re = re.compile(r'\[Pasted text #\d+: \d+ lines \u2192 (.+?)\]')
|
||||
|
||||
def _expand_ref(match):
|
||||
path = Path(match.group(1))
|
||||
# Use try/except instead of path.exists() to avoid TOCTOU race:
|
||||
# the paste file may be deleted between check and read, causing
|
||||
# the input to be silently dropped (#17666).
|
||||
try:
|
||||
return path.read_text(encoding="utf-8")
|
||||
except (OSError, IOError):
|
||||
logger.warning("Paste file gone or unreadable, returning placeholder: %s", path)
|
||||
return match.group(0)
|
||||
|
||||
return paste_ref_re.sub(_expand_ref, text)
|
||||
|
||||
def _print_user_message_preview(self, user_input: str) -> None:
|
||||
"""Render a user message using the normal chat scrollback style."""
|
||||
from cli import ChatConsole, _accent_hex
|
||||
ChatConsole().print(f"[{_accent_hex()}]{'─' * 40}[/]")
|
||||
text = str(user_input or "")
|
||||
if "\n" in text:
|
||||
ChatConsole().print(self._format_submitted_user_message_preview(text))
|
||||
else:
|
||||
ChatConsole().print(f"[bold {_accent_hex()}]●[/] [bold]{_escape(text)}[/]")
|
||||
|
||||
def _stream_reasoning_delta(self, text: str) -> None:
|
||||
"""Stream reasoning/thinking tokens into a dim box above the response.
|
||||
|
||||
Opens a dim reasoning box on first token, streams line-by-line.
|
||||
The box is closed automatically when content tokens start arriving
|
||||
(via _stream_delta → _emit_stream_text).
|
||||
|
||||
Once the response box is open, suppress any further reasoning
|
||||
rendering — a late thinking block (e.g. after an interrupt) would
|
||||
otherwise draw a reasoning box inside the response box.
|
||||
"""
|
||||
from cli import _DIM, _RST, _cprint
|
||||
if not text:
|
||||
return
|
||||
self._reasoning_shown_this_turn = True
|
||||
if getattr(self, "_stream_box_opened", False):
|
||||
return
|
||||
|
||||
# Open reasoning box on first reasoning token
|
||||
if not getattr(self, "_reasoning_box_opened", False):
|
||||
self._reasoning_box_opened = True
|
||||
w = self._scrollback_box_width()
|
||||
r_label = " Reasoning "
|
||||
r_fill = w - 2 - len(r_label)
|
||||
_cprint(f"\n{_DIM}┌─{r_label}{'─' * max(r_fill - 1, 0)}┐{_RST}")
|
||||
|
||||
self._reasoning_buf = getattr(self, "_reasoning_buf", "") + text
|
||||
|
||||
# Emit complete lines, and force-flush long partial lines so
|
||||
# reasoning is visible in real-time even without newlines.
|
||||
while "\n" in self._reasoning_buf:
|
||||
line, self._reasoning_buf = self._reasoning_buf.split("\n", 1)
|
||||
_cprint(f"{_DIM}{line}{_RST}")
|
||||
if len(self._reasoning_buf) > 80:
|
||||
_cprint(f"{_DIM}{self._reasoning_buf}{_RST}")
|
||||
self._reasoning_buf = ""
|
||||
|
||||
def _close_reasoning_box(self) -> None:
|
||||
"""Close the live reasoning box if it's open."""
|
||||
from cli import _DIM, _RST, _cprint
|
||||
if getattr(self, "_reasoning_box_opened", False):
|
||||
# Flush remaining reasoning buffer
|
||||
buf = getattr(self, "_reasoning_buf", "")
|
||||
if buf:
|
||||
_cprint(f"{_DIM}{buf}{_RST}")
|
||||
self._reasoning_buf = ""
|
||||
w = self._scrollback_box_width()
|
||||
_cprint(f"{_DIM}└{'─' * (w - 2)}┘{_RST}")
|
||||
self._reasoning_box_opened = False
|
||||
|
||||
# Flush any content that was deferred while reasoning was rendering.
|
||||
deferred = getattr(self, "_deferred_content", "")
|
||||
if deferred:
|
||||
self._deferred_content = ""
|
||||
self._emit_stream_text(deferred)
|
||||
|
||||
def _stream_delta(self, text) -> None:
|
||||
"""Line-buffered streaming callback for real-time token rendering.
|
||||
|
||||
Receives text deltas from the agent as tokens arrive. Buffers
|
||||
partial lines and emits complete lines via _cprint to work
|
||||
reliably with prompt_toolkit's patch_stdout.
|
||||
|
||||
Reasoning/thinking blocks (<REASONING_SCRATCHPAD>, <think>, etc.)
|
||||
are suppressed during streaming since they'd display raw XML tags.
|
||||
The agent strips them from the final response anyway.
|
||||
|
||||
A ``None`` value signals an intermediate turn boundary (tools are
|
||||
about to execute). Flushes any open boxes and resets state so
|
||||
tool feed lines render cleanly between turns.
|
||||
"""
|
||||
if text is None:
|
||||
self._flush_stream()
|
||||
self._reset_stream_state()
|
||||
return
|
||||
if not text:
|
||||
return
|
||||
|
||||
self._stream_started = True
|
||||
|
||||
# ── Tag-based reasoning suppression ──
|
||||
# Track whether we're inside a reasoning/thinking block.
|
||||
# These tags are model-generated (system prompt tells the model
|
||||
# to use them) and get stripped from final_response. We must
|
||||
# suppress them during streaming too — unless show_reasoning is
|
||||
# enabled, in which case we route the inner content to the
|
||||
# reasoning display box instead of discarding it.
|
||||
_OPEN_TAGS = ("<REASONING_SCRATCHPAD>", "<think>", "<reasoning>", "<THINKING>", "<thinking>", "<thought>")
|
||||
_CLOSE_TAGS = ("</REASONING_SCRATCHPAD>", "</think>", "</reasoning>", "</THINKING>", "</thinking>", "</thought>")
|
||||
|
||||
# Append to a pre-filter buffer first
|
||||
self._stream_prefilt = getattr(self, "_stream_prefilt", "") + text
|
||||
|
||||
# Check if we're entering a reasoning block.
|
||||
# Only match tags that appear at a "block boundary": start of the
|
||||
# stream, after a newline (with optional whitespace), or when nothing
|
||||
# but whitespace has been emitted on the current line.
|
||||
# This prevents false positives when models *mention* tags in prose
|
||||
# like "(/think not producing <think> tags)".
|
||||
#
|
||||
# _stream_last_was_newline tracks whether the last character emitted
|
||||
# (or the start of the stream) is a line boundary. It's True at
|
||||
# stream start and set True whenever emitted text ends with '\n'.
|
||||
if not hasattr(self, "_stream_last_was_newline"):
|
||||
self._stream_last_was_newline = True # start of stream = boundary
|
||||
|
||||
if not getattr(self, "_in_reasoning_block", False):
|
||||
# Case-insensitive matching against a lowercased view so
|
||||
# mixed-case tag variants (<Think>, <THINKING>, …) are caught.
|
||||
prefilt_lower = self._stream_prefilt.lower()
|
||||
for tag in _OPEN_TAGS:
|
||||
tag_lower = tag.lower()
|
||||
search_start = 0
|
||||
while True:
|
||||
idx = prefilt_lower.find(tag_lower, search_start)
|
||||
if idx == -1:
|
||||
break
|
||||
# Check if this is a block boundary position
|
||||
preceding = self._stream_prefilt[:idx]
|
||||
if idx == 0:
|
||||
# At buffer start — only a boundary if we're at
|
||||
# a line start (stream start or last emit ended
|
||||
# with newline)
|
||||
is_block_boundary = getattr(self, "_stream_last_was_newline", True)
|
||||
else:
|
||||
# Find last newline in the buffer before the tag
|
||||
last_nl = preceding.rfind("\n")
|
||||
if last_nl == -1:
|
||||
# No newline in buffer — boundary only if
|
||||
# last emit was a newline AND only whitespace
|
||||
# has accumulated before the tag
|
||||
is_block_boundary = (
|
||||
getattr(self, "_stream_last_was_newline", True)
|
||||
and preceding.strip() == ""
|
||||
)
|
||||
else:
|
||||
# Text between last newline and tag must be
|
||||
# whitespace-only
|
||||
is_block_boundary = preceding[last_nl + 1:].strip() == ""
|
||||
if is_block_boundary:
|
||||
# Emit everything before the tag
|
||||
if preceding:
|
||||
self._emit_stream_text(preceding)
|
||||
self._stream_last_was_newline = preceding.endswith("\n")
|
||||
self._in_reasoning_block = True
|
||||
self._stream_prefilt = self._stream_prefilt[idx + len(tag):]
|
||||
break
|
||||
# Not a block boundary — keep searching after this occurrence
|
||||
search_start = idx + 1
|
||||
if getattr(self, "_in_reasoning_block", False):
|
||||
break
|
||||
|
||||
# Could also be a partial open tag at the end — hold it back
|
||||
if not getattr(self, "_in_reasoning_block", False):
|
||||
# Check for partial tag match at the end (case-insensitive)
|
||||
safe = self._stream_prefilt
|
||||
for tag in _OPEN_TAGS:
|
||||
tag_lower = tag.lower()
|
||||
for i in range(1, len(tag)):
|
||||
if prefilt_lower.endswith(tag_lower[:i]):
|
||||
safe = self._stream_prefilt[:-i]
|
||||
break
|
||||
if safe:
|
||||
self._emit_stream_text(safe)
|
||||
self._stream_last_was_newline = safe.endswith("\n")
|
||||
self._stream_prefilt = self._stream_prefilt[len(safe):]
|
||||
return
|
||||
|
||||
# Inside a reasoning block — look for close tag.
|
||||
# Keep accumulating _stream_prefilt because close tags can arrive
|
||||
# split across multiple tokens (e.g. "</REASONING_SCRATCH" + "PAD>...").
|
||||
if getattr(self, "_in_reasoning_block", False):
|
||||
prefilt_lower = self._stream_prefilt.lower()
|
||||
for tag in _CLOSE_TAGS:
|
||||
idx = prefilt_lower.find(tag.lower())
|
||||
if idx != -1:
|
||||
self._in_reasoning_block = False
|
||||
# When show_reasoning is on, route inner content to
|
||||
# the reasoning display box instead of discarding.
|
||||
if self.show_reasoning:
|
||||
inner = self._stream_prefilt[:idx]
|
||||
if inner:
|
||||
self._stream_reasoning_delta(inner)
|
||||
after = self._stream_prefilt[idx + len(tag):]
|
||||
self._stream_prefilt = ""
|
||||
# Process remaining text after close tag through full
|
||||
# filtering (it could contain another open tag)
|
||||
if after:
|
||||
self._stream_delta(after)
|
||||
return
|
||||
# When show_reasoning is on, stream reasoning content live
|
||||
# instead of silently accumulating. Keep only the tail that
|
||||
# could be a partial close tag prefix.
|
||||
max_tag_len = max(len(t) for t in _CLOSE_TAGS)
|
||||
if len(self._stream_prefilt) > max_tag_len:
|
||||
if self.show_reasoning:
|
||||
# Route the safe prefix to reasoning display
|
||||
safe_reasoning = self._stream_prefilt[:-max_tag_len]
|
||||
self._stream_reasoning_delta(safe_reasoning)
|
||||
self._stream_prefilt = self._stream_prefilt[-max_tag_len:]
|
||||
return
|
||||
|
||||
def _emit_stream_text(self, text: str) -> None:
|
||||
"""Emit filtered text to the streaming display."""
|
||||
from cli import (
|
||||
HermesCLI,
|
||||
_ACCENT,
|
||||
_RST,
|
||||
_STREAM_PAD,
|
||||
_STREAM_PARTIAL_PREVIEW_LEN,
|
||||
_cprint,
|
||||
_strip_markdown_syntax,
|
||||
_terminal_width_for_streaming,
|
||||
datetime,
|
||||
is_table_divider,
|
||||
looks_like_table_row,
|
||||
realign_markdown_tables,
|
||||
)
|
||||
if not text:
|
||||
return
|
||||
|
||||
# When show_reasoning is on and reasoning is still rendering,
|
||||
# defer content until the reasoning box closes. This ensures the
|
||||
# reasoning block always appears BEFORE the response in the terminal.
|
||||
if self.show_reasoning and getattr(self, "_reasoning_box_opened", False):
|
||||
self._deferred_content = getattr(self, "_deferred_content", "") + text
|
||||
return
|
||||
|
||||
# Close the live reasoning box before opening the response box
|
||||
self._close_reasoning_box()
|
||||
|
||||
# Open the response box header on the very first visible text
|
||||
if not self._stream_box_opened:
|
||||
# Strip leading whitespace/newlines before first visible content
|
||||
text = text.lstrip("\n")
|
||||
if not text:
|
||||
return
|
||||
self._stream_box_opened = True
|
||||
try:
|
||||
from hermes_cli.skin_engine import get_active_skin
|
||||
_skin = get_active_skin()
|
||||
label = _skin.get_branding("response_label", "⚕ Hermes")
|
||||
_text_hex = _skin.get_color("banner_text", "#FFF8DC")
|
||||
except Exception:
|
||||
label = "⚕ Hermes"
|
||||
_text_hex = "#FFF8DC"
|
||||
# Build a true-color ANSI escape for the response text color
|
||||
# so streamed content matches the Rich Panel appearance.
|
||||
try:
|
||||
_r = int(_text_hex[1:3], 16)
|
||||
_g = int(_text_hex[3:5], 16)
|
||||
_b = int(_text_hex[5:7], 16)
|
||||
self._stream_text_ansi = f"\033[38;2;{_r};{_g};{_b}m"
|
||||
except (ValueError, IndexError):
|
||||
self._stream_text_ansi = ""
|
||||
if self.show_timestamps:
|
||||
label = f"{label} {datetime.now().strftime(getattr(self, 'timestamp_format', '%H:%M'))}"
|
||||
w = self._scrollback_box_width()
|
||||
fill = w - 2 - HermesCLI._status_bar_display_width(label)
|
||||
_cprint(f"\n{_ACCENT}╭─{label}{'─' * max(fill - 1, 0)}╮{_RST}")
|
||||
|
||||
self._stream_buf += text
|
||||
|
||||
# Emit complete lines, keep partial remainder in buffer
|
||||
_tc = getattr(self, "_stream_text_ansi", "")
|
||||
|
||||
def _emit_one(printed_line: str) -> None:
|
||||
_cprint(f"{_STREAM_PAD}{_tc}{printed_line}{_RST}" if _tc else f"{_STREAM_PAD}{printed_line}")
|
||||
|
||||
def _flush_table_buf() -> None:
|
||||
buf = self._stream_table_buf
|
||||
self._stream_table_buf = []
|
||||
self._in_stream_table = False
|
||||
if not buf:
|
||||
return
|
||||
# Strip cell-level markdown (`code`, **bold**, ~~strike~~) FIRST
|
||||
# so the realigner pads to the final visible cell width, not
|
||||
# the marker-decorated source width. Otherwise a body row
|
||||
# like `` | Bold | `**bold**` | `` lands narrower than its
|
||||
# header column once the markers are removed.
|
||||
joined = "\n".join(buf)
|
||||
if self.final_response_markdown == "strip":
|
||||
joined = _strip_markdown_syntax(joined)
|
||||
block = realign_markdown_tables(joined, _terminal_width_for_streaming())
|
||||
for ln in block.split("\n"):
|
||||
_emit_one(ln)
|
||||
|
||||
while "\n" in self._stream_buf:
|
||||
line, self._stream_buf = self._stream_buf.split("\n", 1)
|
||||
|
||||
# Hold table-shaped lines in a side-buffer so we can re-pad
|
||||
# the whole block once it ends. Streaming line-by-line, we
|
||||
# cannot re-align mid-table without reflowing already-printed
|
||||
# rows; the cost is that the user sees the table appear in a
|
||||
# single batch when the block closes instead of row-by-row.
|
||||
if self._in_stream_table:
|
||||
if looks_like_table_row(line) or is_table_divider(line):
|
||||
self._stream_table_buf.append(line)
|
||||
continue
|
||||
# Block ended — flush the realigned table, then fall
|
||||
# through to print the current (non-table) line.
|
||||
_flush_table_buf()
|
||||
elif looks_like_table_row(line):
|
||||
self._stream_table_buf.append(line)
|
||||
self._in_stream_table = True
|
||||
continue
|
||||
|
||||
if self.final_response_markdown == "strip":
|
||||
line = _strip_markdown_syntax(line)
|
||||
_emit_one(line)
|
||||
|
||||
# Long partial lines are emitted ONLY at real newlines — we no
|
||||
# longer hard-wrap paragraphs at terminal width ourselves. Each
|
||||
# logical line lands in scrollback as one line; the TERMINAL
|
||||
# soft-wraps it visually, and emulators (iTerm2/kitty/VTE/
|
||||
# xterm.js/Windows Terminal) rejoin soft-wrapped rows on copy,
|
||||
# so highlight-copy yields the original unwrapped text — same
|
||||
# outcome as the TUI's selection copy. (The pre-July-2026 chunk
|
||||
# emitter baked real '\n's into every long paragraph, which is
|
||||
# exactly what polluted copy/paste.)
|
||||
#
|
||||
# TTFT perception: while a long opening paragraph accumulates
|
||||
# without a newline, mirror its tail into the status-bar spinner
|
||||
# line so the user sees tokens arriving instead of a blank box.
|
||||
if (
|
||||
self._stream_buf
|
||||
and not self._in_stream_table
|
||||
and not self._stream_buf.lstrip().startswith("|")
|
||||
and len(self._stream_buf) >= 80
|
||||
):
|
||||
preview = self._stream_buf[-int(_STREAM_PARTIAL_PREVIEW_LEN):]
|
||||
cut = preview.find(" ")
|
||||
if 0 < cut < len(preview) - 1:
|
||||
preview = preview[cut + 1:]
|
||||
try:
|
||||
self._spinner_text = f"… {preview}"
|
||||
self._invalidate()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def _flush_stream(self) -> None:
|
||||
"""Emit any remaining partial line from the stream buffer and close the box."""
|
||||
from cli import (
|
||||
_ACCENT,
|
||||
_RST,
|
||||
_STREAM_PAD,
|
||||
_cprint,
|
||||
_strip_markdown_syntax,
|
||||
_terminal_width_for_streaming,
|
||||
is_table_divider,
|
||||
looks_like_table_row,
|
||||
realign_markdown_tables,
|
||||
)
|
||||
# If we're still inside a "reasoning block" at end-of-stream, it was
|
||||
# a false positive — the model mentioned a tag like <think> in prose
|
||||
# but never closed it. Recover the buffered content as regular text.
|
||||
if getattr(self, "_in_reasoning_block", False) and getattr(self, "_stream_prefilt", ""):
|
||||
self._in_reasoning_block = False
|
||||
self._emit_stream_text(self._stream_prefilt)
|
||||
self._stream_prefilt = ""
|
||||
|
||||
# Close reasoning box if still open (in case no content tokens arrived)
|
||||
self._close_reasoning_box()
|
||||
|
||||
_tc = getattr(self, "_stream_text_ansi", "")
|
||||
|
||||
# If the stream buffer has a trailing partial line that looks like
|
||||
# a table row, fold it into the table buffer so the whole block
|
||||
# gets re-aligned together. Otherwise the final row prints raw
|
||||
# (with the model's original under-padded spacing) while the rows
|
||||
# above it are aligned.
|
||||
if (
|
||||
self._stream_buf
|
||||
and getattr(self, "_in_stream_table", False)
|
||||
and (looks_like_table_row(self._stream_buf) or is_table_divider(self._stream_buf))
|
||||
):
|
||||
self._stream_table_buf.append(self._stream_buf)
|
||||
self._stream_buf = ""
|
||||
|
||||
# Flush any buffered table rows first so their padding is
|
||||
# finalised before the stream remainder lands.
|
||||
if getattr(self, "_stream_table_buf", None):
|
||||
joined = "\n".join(self._stream_table_buf)
|
||||
self._stream_table_buf = []
|
||||
self._in_stream_table = False
|
||||
if self.final_response_markdown == "strip":
|
||||
joined = _strip_markdown_syntax(joined)
|
||||
block = realign_markdown_tables(joined, _terminal_width_for_streaming())
|
||||
for ln in block.split("\n"):
|
||||
_cprint(f"{_STREAM_PAD}{_tc}{ln}{_RST}" if _tc else f"{_STREAM_PAD}{ln}")
|
||||
|
||||
if self._stream_buf:
|
||||
line = _strip_markdown_syntax(self._stream_buf) if self.final_response_markdown == "strip" else self._stream_buf
|
||||
_cprint(f"{_STREAM_PAD}{_tc}{line}{_RST}" if _tc else f"{_STREAM_PAD}{line}")
|
||||
self._stream_buf = ""
|
||||
|
||||
# Close the response box
|
||||
if self._stream_box_opened:
|
||||
w = self._scrollback_box_width()
|
||||
_cprint(f"{_ACCENT}╰{'─' * (w - 2)}╯{_RST}")
|
||||
|
||||
def _reset_stream_state(self) -> None:
|
||||
"""Reset streaming state before each agent invocation."""
|
||||
self._stream_buf = ""
|
||||
self._stream_started = False
|
||||
self._stream_box_opened = False
|
||||
self._stream_text_ansi = ""
|
||||
self._stream_prefilt = ""
|
||||
self._in_reasoning_block = False
|
||||
self._stream_last_was_newline = True
|
||||
self._reasoning_box_opened = False
|
||||
self._reasoning_buf = ""
|
||||
self._reasoning_preview_buf = ""
|
||||
self._deferred_content = ""
|
||||
self._stream_table_buf = []
|
||||
self._in_stream_table = False
|
||||
|
||||
def _slow_command_status(self, command: str) -> str:
|
||||
"""Return a user-facing status message for slower slash commands."""
|
||||
cmd_lower = command.lower().strip()
|
||||
if cmd_lower.startswith("/skills search"):
|
||||
return "Searching skills..."
|
||||
if cmd_lower.startswith("/skills browse"):
|
||||
return "Loading skills..."
|
||||
if cmd_lower.startswith("/skills inspect"):
|
||||
return "Inspecting skill..."
|
||||
if cmd_lower.startswith("/skills install"):
|
||||
return "Installing skill..."
|
||||
if cmd_lower.startswith("/skills"):
|
||||
return "Processing skills command..."
|
||||
if cmd_lower == "/reload-mcp":
|
||||
return "Reloading MCP servers..."
|
||||
if cmd_lower == "/reload-skills" or cmd_lower == "/reload_skills":
|
||||
return "Reloading skills..."
|
||||
if cmd_lower.startswith("/browser"):
|
||||
return "Configuring browser..."
|
||||
return "Processing command..."
|
||||
|
||||
def _command_spinner_frame(self) -> str:
|
||||
"""Return the current spinner frame for slow slash commands."""
|
||||
from cli import _COMMAND_SPINNER_FRAMES
|
||||
frame_idx = int(time.monotonic() * 10) % len(_COMMAND_SPINNER_FRAMES)
|
||||
return _COMMAND_SPINNER_FRAMES[frame_idx]
|
||||
|
||||
@contextmanager
|
||||
def _busy_command(self, status: str, *, blocks_input: bool = True):
|
||||
"""Expose a temporary busy state in the TUI while a slash command runs.
|
||||
|
||||
Most synchronous slash commands must reserve the composer because their
|
||||
completion changes the active session state. Manual compression is safe
|
||||
to draft through: the queued input is processed against the compacted
|
||||
history after the command completes.
|
||||
"""
|
||||
previous_blocks_input = getattr(self, "_command_blocks_input", False)
|
||||
self._command_running = True
|
||||
self._command_blocks_input = blocks_input
|
||||
self._command_status = status
|
||||
self._invalidate(min_interval=0.0)
|
||||
try:
|
||||
print(f"⏳ {status}")
|
||||
yield
|
||||
finally:
|
||||
self._command_running = False
|
||||
self._command_blocks_input = previous_blocks_input
|
||||
self._command_status = ""
|
||||
self._invalidate(min_interval=0.0)
|
||||
|
||||
def _preprocess_images_with_vision(self, text: str, images: list, *, announce: bool = True) -> str:
|
||||
"""Analyze attached images via the vision tool and return enriched text.
|
||||
|
||||
Instead of embedding raw base64 ``image_url`` content parts in the
|
||||
conversation (which only works with vision-capable models), this
|
||||
pre-processes each image through the auxiliary vision model (Gemini
|
||||
Flash) and prepends the descriptions to the user's message — the
|
||||
same approach the messaging gateway uses.
|
||||
|
||||
The local file path is included so the agent can re-examine the
|
||||
image later with ``vision_analyze`` if needed.
|
||||
"""
|
||||
from cli import _DIM, _RST, _cprint
|
||||
import asyncio as _asyncio
|
||||
from tools.vision_tools import vision_analyze_tool
|
||||
|
||||
analysis_prompt = (
|
||||
"Describe everything visible in this image in thorough detail. "
|
||||
"Include any text, code, data, objects, people, layout, colors, "
|
||||
"and any other notable visual information."
|
||||
)
|
||||
|
||||
enriched_parts = []
|
||||
for img_path in images:
|
||||
if not img_path.exists():
|
||||
continue
|
||||
size_kb = img_path.stat().st_size // 1024
|
||||
if announce:
|
||||
_cprint(f" {_DIM}👁️ analyzing {img_path.name} ({size_kb}KB)...{_RST}")
|
||||
try:
|
||||
result_json = _asyncio.run(
|
||||
vision_analyze_tool(image_url=str(img_path), user_prompt=analysis_prompt)
|
||||
)
|
||||
result = json.loads(result_json)
|
||||
if result.get("success"):
|
||||
description = result.get("analysis", "")
|
||||
enriched_parts.append(
|
||||
f"[The user attached an image. Here's what it contains:\n{description}]\n"
|
||||
f"[If you need a closer look, use vision_analyze with "
|
||||
f"image_url: {img_path}]"
|
||||
)
|
||||
if announce:
|
||||
_cprint(f" {_DIM}✓ image analyzed{_RST}")
|
||||
else:
|
||||
enriched_parts.append(
|
||||
f"[The user attached an image but it couldn't be analyzed. "
|
||||
f"You can try examining it with vision_analyze using "
|
||||
f"image_url: {img_path}]"
|
||||
)
|
||||
if announce:
|
||||
_cprint(f" {_DIM}⚠ vision analysis failed — path included for retry{_RST}")
|
||||
except Exception as e:
|
||||
enriched_parts.append(
|
||||
f"[The user attached an image but analysis failed ({e}). "
|
||||
f"You can try examining it with vision_analyze using "
|
||||
f"image_url: {img_path}]"
|
||||
)
|
||||
if announce:
|
||||
_cprint(f" {_DIM}⚠ vision analysis error — path included for retry{_RST}")
|
||||
|
||||
# Combine: vision descriptions first, then the user's original text
|
||||
user_text = text if isinstance(text, str) and text else ""
|
||||
if enriched_parts:
|
||||
prefix = "\n\n".join(enriched_parts)
|
||||
return f"{prefix}\n\n{user_text}" if user_text else prefix
|
||||
return user_text or "What do you see in this image?"
|
||||
|
||||
def _output_console(self):
|
||||
"""Use prompt_toolkit-safe Rich rendering once the TUI is live."""
|
||||
from cli import ChatConsole
|
||||
if getattr(self, "_app", None):
|
||||
return ChatConsole()
|
||||
return self.console
|
||||
|
||||
def _console_print(self, *args, **kwargs):
|
||||
"""Print through the active command-safe console."""
|
||||
self._output_console().print(*args, **kwargs)
|
||||
|
||||
def _on_tool_gen_start(self, tool_name: str) -> None:
|
||||
"""Called when the model begins generating tool-call arguments.
|
||||
|
||||
Closes any open streaming boxes (reasoning / response) exactly once,
|
||||
then prints a short status line so the user sees activity instead of
|
||||
a frozen screen while a large payload (e.g. 45 KB write_file) streams.
|
||||
"""
|
||||
from cli import _cprint
|
||||
if getattr(self, "_stream_box_opened", False):
|
||||
self._flush_stream()
|
||||
self._stream_box_opened = False
|
||||
self._close_reasoning_box()
|
||||
|
||||
from agent.display import get_tool_emoji
|
||||
emoji = get_tool_emoji(tool_name, default="⚡")
|
||||
_cprint(f" ┊ {emoji} preparing {tool_name}…")
|
||||
|
||||
def _on_tool_progress(self, event_type: str, function_name: str = None, preview: str = None, function_args: dict = None, **kwargs):
|
||||
"""Called on tool lifecycle events (tool.started, tool.completed, reasoning.available, etc.).
|
||||
|
||||
Updates the TUI spinner widget so the user can see what the agent
|
||||
is doing during tool execution (fills the gap between thinking
|
||||
spinner and next response).
|
||||
|
||||
On tool.started, records a monotonic timestamp so get_spinner_text()
|
||||
can show a live elapsed timer (the TUI poll loop already invalidates
|
||||
every ~0.15s, so the counter updates automatically).
|
||||
|
||||
When tool_progress_mode is "all" or "new", also prints a persistent
|
||||
stacked line to scrollback on tool.completed so users can see the
|
||||
full history of tool calls (not just the current one in the spinner).
|
||||
"""
|
||||
from cli import CLI_CONFIG, _DIM, _RST, _cprint, _hermes_home
|
||||
# MoA reference-model outputs: render each reference's answer as a
|
||||
# labelled thinking-style block BEFORE the aggregator acts, so the user
|
||||
# sees the mixture-of-agents process instead of a silent pause. These
|
||||
# are display-only events emitted by the MoA facade (agent_init relay);
|
||||
# they never enter message history.
|
||||
if event_type == "moa.reference":
|
||||
label = function_name or "reference"
|
||||
text = preview or ""
|
||||
idx = kwargs.get("moa_index")
|
||||
count = kwargs.get("moa_count")
|
||||
header = f"Reference {idx}/{count} — {label}" if idx and count else f"Reference — {label}"
|
||||
try:
|
||||
self._flush_reasoning_preview(force=True)
|
||||
except Exception:
|
||||
pass
|
||||
_cprint(f" {_DIM}┊ ◇ {header}{_RST}")
|
||||
try:
|
||||
self._emit_reasoning_preview(text)
|
||||
except Exception:
|
||||
# Fallback: print the raw text dimmed if the preview helper fails.
|
||||
if text.strip():
|
||||
_cprint(f" {_DIM}{text.strip()}{_RST}")
|
||||
self._invalidate()
|
||||
return
|
||||
if event_type == "moa.aggregating":
|
||||
agg = function_name or ""
|
||||
self._spinner_text = f"◆ aggregating ({agg})" if agg else "◆ aggregating"
|
||||
self._invalidate()
|
||||
return
|
||||
|
||||
# Feed the pet: tools mean "running" (not reasoning); a failed tool
|
||||
# latches the turn so it ends on a sulk.
|
||||
if event_type == "tool.started":
|
||||
self._pet_reasoning = False
|
||||
elif event_type == "tool.completed" and kwargs.get("is_error"):
|
||||
self._pet_turn_error = True
|
||||
elif event_type and event_type.startswith("reasoning"):
|
||||
self._pet_reasoning = True
|
||||
|
||||
if event_type == "tool.completed":
|
||||
self._tool_start_time = 0.0
|
||||
# Per-turn accounting: this feed already sees every tool call with
|
||||
# its result, so the summary line needs no agent-loop state.
|
||||
self._turn_summary_record(
|
||||
function_name, kwargs.get("result"), kwargs.get("is_error", False)
|
||||
)
|
||||
# Focus view: count the scrollback line we are NOT printing, so the
|
||||
# post-turn recovery line can report how much was hidden. Counted
|
||||
# against the pre-focus tool-progress mode, so a user who already
|
||||
# had /verbose off is never told focus hid something it didn't.
|
||||
if getattr(self, "_focus_view_enabled", False):
|
||||
try:
|
||||
self._note_focus_hidden_line(function_name or "")
|
||||
except Exception:
|
||||
pass
|
||||
# Print stacked scrollback line for "new" / "all" / "verbose" modes.
|
||||
# "verbose" was previously omitted here, so non-streaming model
|
||||
# calls (MoA aggregator, copilot-acp) rendered each tool only into
|
||||
# the transient spinner line — which overwrites itself, so no
|
||||
# scrollable tool history accumulated. Streaming models hid the bug
|
||||
# because _on_tool_gen_start commits a "preparing" line per tool;
|
||||
# non-streaming calls never emit that, leaving verbose mode with no
|
||||
# committed line at all. "verbose" is strictly more than "all", so
|
||||
# it must commit at least the same line.
|
||||
if function_name and self.tool_progress_mode in {"new", "all", "verbose"}:
|
||||
duration = kwargs.get("duration", 0.0)
|
||||
# Pop stored args from tool.started for this function
|
||||
stored = self._pending_tool_info.get(function_name)
|
||||
stored_args = stored.pop(0) if stored else {}
|
||||
if stored is not None and not stored:
|
||||
del self._pending_tool_info[function_name]
|
||||
# "new" mode: skip consecutive repeats of the same tool
|
||||
if self.tool_progress_mode == "new" and function_name == self._last_scrollback_tool:
|
||||
self._invalidate()
|
||||
return
|
||||
self._last_scrollback_tool = function_name
|
||||
try:
|
||||
from agent.display import get_cute_tool_message
|
||||
line = get_cute_tool_message(function_name, stored_args, duration, result=kwargs.get("result"))
|
||||
_cprint(f" {line}")
|
||||
except Exception:
|
||||
pass
|
||||
# First-touch onboarding: on the first tool in this process
|
||||
# that takes longer than the threshold while we're in the
|
||||
# noisiest progress mode, print a one-time hint about
|
||||
# /verbose. Latched on self so it fires at most once per
|
||||
# process; persisted to config.yaml so it never fires again
|
||||
# across processes either.
|
||||
try:
|
||||
if (
|
||||
not getattr(self, "_long_tool_hint_fired", False)
|
||||
and self.tool_progress_mode == "all"
|
||||
and duration >= 30.0
|
||||
):
|
||||
from agent.onboarding import (
|
||||
TOOL_PROGRESS_FLAG,
|
||||
is_seen,
|
||||
mark_seen,
|
||||
tool_progress_hint_cli,
|
||||
)
|
||||
if not is_seen(CLI_CONFIG, TOOL_PROGRESS_FLAG):
|
||||
self._long_tool_hint_fired = True
|
||||
_cprint(f" {_DIM}{tool_progress_hint_cli()}{_RST}")
|
||||
mark_seen(_hermes_home / "config.yaml", TOOL_PROGRESS_FLAG)
|
||||
CLI_CONFIG.setdefault("onboarding", {}).setdefault("seen", {})[TOOL_PROGRESS_FLAG] = True
|
||||
except Exception:
|
||||
pass
|
||||
self._invalidate()
|
||||
return
|
||||
if event_type != "tool.started":
|
||||
return
|
||||
if function_name and not function_name.startswith("_"):
|
||||
from agent.display import get_tool_emoji
|
||||
emoji = get_tool_emoji(function_name)
|
||||
label = preview or function_name
|
||||
from agent.display import get_tool_preview_max_len
|
||||
_pl = get_tool_preview_max_len()
|
||||
if _pl > 0 and len(label) > _pl:
|
||||
label = label[:_pl - 3] + "..."
|
||||
self._spinner_text = f"{emoji} {label}"
|
||||
self._tool_start_time = time.monotonic()
|
||||
# Store args for stacked scrollback line on completion
|
||||
self._pending_tool_info.setdefault(function_name, []).append(
|
||||
function_args if function_args is not None else {}
|
||||
)
|
||||
self._invalidate()
|
||||
|
||||
def _on_tool_start(self, tool_call_id: str, function_name: str, function_args: dict):
|
||||
"""Capture local before-state for write-capable tools."""
|
||||
from cli import logger
|
||||
try:
|
||||
from agent.display import capture_local_edit_snapshot
|
||||
|
||||
snapshot = capture_local_edit_snapshot(function_name, function_args)
|
||||
if snapshot is not None:
|
||||
self._pending_edit_snapshots[tool_call_id] = snapshot
|
||||
except Exception:
|
||||
logger.debug("Edit snapshot capture failed for %s", function_name, exc_info=True)
|
||||
|
||||
def _on_tool_complete(self, tool_call_id: str, function_name: str, function_args: dict, function_result: str):
|
||||
"""Render file edits with inline diff after write-capable tools complete."""
|
||||
from cli import _cprint, logger
|
||||
# A top-level delegate_task dispatches in the background and re-enters as
|
||||
# a fresh turn when done. Say so once — no spinner, nothing to poll — so
|
||||
# the idle prompt doesn't read as "nothing happened" (⛓ tracks the work).
|
||||
if function_name == "delegate_task":
|
||||
try:
|
||||
parsed = json.loads(function_result) if isinstance(function_result, str) else (function_result or {})
|
||||
except Exception:
|
||||
parsed = {}
|
||||
if isinstance(parsed, dict) and parsed.get("status") == "dispatched" and parsed.get("mode") == "background":
|
||||
n = parsed.get("count") or 1
|
||||
noun, tail = ("task", "it finishes") if n == 1 else (f"{n} tasks", "they finish")
|
||||
try:
|
||||
_cprint(f"\033[2m\u21a9 Background {noun} running — I'll resume when {tail}. Keep chatting.\033[0m")
|
||||
except Exception:
|
||||
pass
|
||||
snapshot = self._pending_edit_snapshots.pop(tool_call_id, None)
|
||||
try:
|
||||
from agent.display import render_edit_diff_with_delta
|
||||
|
||||
render_edit_diff_with_delta(
|
||||
function_name,
|
||||
function_result,
|
||||
function_args=function_args,
|
||||
snapshot=snapshot,
|
||||
print_fn=_cprint,
|
||||
)
|
||||
except Exception:
|
||||
logger.debug("Edit diff preview failed for %s", function_name, exc_info=True)
|
||||
@@ -0,0 +1,624 @@
|
||||
"""Terminal repaint/resize recovery, input-mode healing, and clipboard helpers for the interactive CLI
|
||||
|
||||
Mixin split out of ``cli.py``; bound onto ``HermesCLI`` via the MRO. cli.py-internal
|
||||
symbols are imported LAZILY inside each method (``from cli import ...``) — the mixin
|
||||
never imports ``cli`` at module load time (import cycle).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import errno
|
||||
import os
|
||||
import shutil
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
|
||||
from hermes_constants import get_hermes_home
|
||||
|
||||
|
||||
class CLITerminalMixin:
|
||||
"""Terminal repaint/resize recovery, input-mode healing, and clipboard helpers for the interactive CLI"""
|
||||
|
||||
def _mark_terminal_io_broken(self, reason: str = "") -> None:
|
||||
"""Stop UI paints after the PTY/stdout becomes unusable (#81521)."""
|
||||
from cli import logger
|
||||
if getattr(self, "_terminal_io_broken", False):
|
||||
return
|
||||
self._terminal_io_broken = True
|
||||
try:
|
||||
self._pet_stop_anim()
|
||||
except Exception:
|
||||
pass
|
||||
logger.warning(
|
||||
"Terminal I/O broken%s — freezing UI paints to avoid redraw storm (#81521)",
|
||||
f" ({reason})" if reason else "",
|
||||
)
|
||||
|
||||
def _invalidate(self, min_interval: float = 0.25) -> None:
|
||||
"""Throttled UI repaint for high-frequency background updates.
|
||||
|
||||
Use this for spinner frames, streaming token flushes, and other
|
||||
repaints that can fire many times per second — the throttle prevents
|
||||
terminal blinking on slow/SSH connections, and the resize-recovery
|
||||
guard avoids stamping footer/status-bar chrome into scrollback while a
|
||||
SIGWINCH reflow is in flight.
|
||||
|
||||
Do NOT use this for user-blocking modal prompts (approval / clarify /
|
||||
sudo). Those are rare, one-shot, user-blocking events that must paint
|
||||
immediately; route them through ``self._app.invalidate()`` directly, the
|
||||
same way the modal key-binding handlers already do. Sending a modal's
|
||||
entry paint through this throttle lets an unrelated background repaint
|
||||
within the 250ms window — or an in-flight resize — silently drop it, so
|
||||
the prompt never renders and times out unseen (#41098).
|
||||
"""
|
||||
if getattr(self, "_terminal_io_broken", False):
|
||||
return
|
||||
if getattr(self, "_resize_recovery_pending", False):
|
||||
return
|
||||
now = time.monotonic()
|
||||
if hasattr(self, "_app") and self._app and (now - getattr(self, "_last_invalidate", 0.0)) >= min_interval:
|
||||
self._last_invalidate = now
|
||||
try:
|
||||
self._app.invalidate()
|
||||
except OSError as exc:
|
||||
if getattr(exc, "errno", None) == errno.EIO:
|
||||
self._mark_terminal_io_broken("invalidate")
|
||||
return
|
||||
raise
|
||||
|
||||
def _paint_now(self) -> None:
|
||||
"""Immediate, unthrottled repaint for user-blocking modal prompts.
|
||||
|
||||
Background-thread callbacks (approval / clarify / sudo) set their modal
|
||||
state then call this to make the panel visible at once. It deliberately
|
||||
bypasses the ``_invalidate`` throttle and resize-recovery guard — a
|
||||
modal the user is actively waiting on must never be dropped — mirroring
|
||||
the direct ``event.app.invalidate()`` the modal key-binding handlers
|
||||
already use. See ``_invalidate`` for why the throttle must not gate
|
||||
these paints (#41098).
|
||||
"""
|
||||
if getattr(self, "_terminal_io_broken", False):
|
||||
return
|
||||
app = getattr(self, "_app", None)
|
||||
if app is not None:
|
||||
try:
|
||||
app.invalidate()
|
||||
except OSError as exc:
|
||||
if getattr(exc, "errno", None) == errno.EIO:
|
||||
self._mark_terminal_io_broken("paint_now")
|
||||
return
|
||||
raise
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def _force_full_redraw(self) -> None:
|
||||
"""Force a clean full-screen repaint of the prompt_toolkit UI.
|
||||
|
||||
Used to recover from terminal buffer drift caused by external
|
||||
redraws we can't detect — e.g. macOS cmux / tmux tab switches,
|
||||
``clear`` issued from a subshell, or SSH window restores. These
|
||||
wipe or repaint the terminal without firing SIGWINCH, so
|
||||
prompt_toolkit's tracked ``_cursor_pos`` no longer matches reality
|
||||
and the next incremental redraw stacks on top of stale content
|
||||
(ghost status bars, duplicated prompts).
|
||||
|
||||
Bound to Ctrl+L and exposed as the ``/redraw`` slash command,
|
||||
matching the standard terminal-UX convention (bash, zsh, fish,
|
||||
vim, htop).
|
||||
"""
|
||||
from cli import _replay_output_history
|
||||
if getattr(self, "_terminal_io_broken", False):
|
||||
return
|
||||
app = getattr(self, "_app", None)
|
||||
if not app:
|
||||
return
|
||||
self._clear_prompt_toolkit_screen(
|
||||
app,
|
||||
rebuild_scrollback=self._redraw_rebuilds_scrollback(),
|
||||
)
|
||||
if getattr(self, "_terminal_io_broken", False):
|
||||
return
|
||||
_replay_output_history()
|
||||
self._pet_queue_kitty_frame()
|
||||
try:
|
||||
app.invalidate()
|
||||
except OSError as exc:
|
||||
if getattr(exc, "errno", None) == errno.EIO:
|
||||
self._mark_terminal_io_broken("force_full_redraw")
|
||||
return
|
||||
raise
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def _schedule_focus_regain_redraw(self, min_interval: float = 1.0) -> None:
|
||||
"""Repaint after a terminal focus-in report (``CSI I``), rate-limited.
|
||||
|
||||
Terminals with focus tracking active (Ghostty, iTerm2, xterm builds,
|
||||
multiplexers that toggle DECSET 1004 upstream) emit ``\\x1b[I`` when
|
||||
the Hermes tab/window becomes visible again. Emulators can coalesce
|
||||
or drop hidden-tab output and repaint the surface while we're
|
||||
invisible, so on regain prompt_toolkit's incremental diff stacks on
|
||||
stale content — a second copy of the composer/prompt chrome next to
|
||||
the ghost of the old one (#60920 focus-regain variant, #25337).
|
||||
|
||||
The stock handling maps ``CSI I``/``CSI O`` to ``Keys.Ignore`` so the
|
||||
bytes never pollute the input buffer; this hook additionally routes
|
||||
focus-in through the same recovery as Ctrl+L / ``/redraw``. It is
|
||||
self-gating: terminals that never enable focus tracking never emit
|
||||
the sequence, so nothing changes for them. Rate-limited so a burst of
|
||||
focus reports (rapid Alt+Tab, mux pane hops) repaints at most once
|
||||
per ``min_interval`` seconds.
|
||||
"""
|
||||
now = time.monotonic()
|
||||
last = getattr(self, "_last_focus_regain_redraw", 0.0)
|
||||
if now - last < min_interval:
|
||||
return
|
||||
self._last_focus_regain_redraw = now
|
||||
self._force_full_redraw()
|
||||
|
||||
@staticmethod
|
||||
def _redraw_rebuilds_scrollback() -> bool:
|
||||
"""Return whether CLI redraw/resize recovery should clear scrollback.
|
||||
|
||||
Some terminal/tmux stacks move prompt_toolkit's non-fullscreen bottom
|
||||
chrome into scrollback when the window is maximized/restored. A normal
|
||||
CSI 2J viewport clear cannot remove those stale prompt/input-rule rows,
|
||||
so users who hit that class of bug need CSI 3J as well, followed by the
|
||||
existing bounded output-history replay.
|
||||
"""
|
||||
from cli import CLI_CONFIG
|
||||
display_config = CLI_CONFIG.get("display") if isinstance(CLI_CONFIG, dict) else {}
|
||||
if not isinstance(display_config, dict):
|
||||
display_config = {}
|
||||
raw = display_config.get("cli_rebuild_scrollback_on_redraw", False)
|
||||
if isinstance(raw, str):
|
||||
return raw.strip().lower() in {"1", "true", "yes", "on", "always"}
|
||||
return bool(raw)
|
||||
|
||||
def _recover_terminal_after_interrupt(self) -> None:
|
||||
"""Recover the terminal after an interrupted agent turn (#33271).
|
||||
|
||||
When the user interrupts a running turn by typing a new message,
|
||||
prompt_toolkit may have an in-flight ``CSI 6n`` cursor-position query
|
||||
whose reply (``ESC[<row>;<col>R``) arrives on stdin after the input
|
||||
parser has torn down. The reply then leaks as literal text
|
||||
(``^[[19;1R``) and the VT100 parser can stall in a partial-escape
|
||||
state, accepting no further keystrokes — the terminal appears frozen.
|
||||
|
||||
Two steps recover a sane state:
|
||||
1. ``flush_stdin()`` drains stray escape bytes from the OS input
|
||||
buffer (``termios.tcflush(TCIFLUSH)``; no-op on non-TTY).
|
||||
2. ``_force_full_redraw()`` drops prompt_toolkit's cached
|
||||
screen/cursor state and forces a clean repaint.
|
||||
|
||||
Both steps are independently safe and self-guard, so a failure of one
|
||||
never prevents the other. If the PTY is already dead (EIO), skip the
|
||||
redraw entirely — painting a broken fd is the #81521 redraw storm.
|
||||
"""
|
||||
if getattr(self, "_terminal_io_broken", False):
|
||||
return
|
||||
try:
|
||||
from hermes_cli.curses_ui import flush_stdin
|
||||
flush_stdin()
|
||||
except Exception:
|
||||
pass
|
||||
# #60920: The interruption marker is now printed with
|
||||
# _suspend_output_history in chat(), so _OUTPUT_HISTORY only
|
||||
# contains the normal response text (no marker text). Do NOT
|
||||
# clear history here — _force_full_redraw → _replay_output_history
|
||||
# replays the response correctly without duplicating the marker.
|
||||
# The /redraw + Ctrl+L paths also preserve replay for scrollback
|
||||
# recovery as intended.
|
||||
self._force_full_redraw()
|
||||
|
||||
def _clear_prompt_toolkit_screen(self, app, *, rebuild_scrollback: bool = False) -> None:
|
||||
"""Clear the terminal and reset prompt_toolkit renderer state."""
|
||||
if getattr(self, "_terminal_io_broken", False):
|
||||
return
|
||||
try:
|
||||
renderer = app.renderer
|
||||
out = renderer.output
|
||||
out.reset_attributes()
|
||||
out.erase_screen()
|
||||
if rebuild_scrollback:
|
||||
try:
|
||||
out.write_raw("\x1b[3J")
|
||||
except Exception:
|
||||
pass
|
||||
out.cursor_goto(0, 0)
|
||||
out.flush()
|
||||
# Drop prompt_toolkit's cached screen + cursor state so the
|
||||
# next _redraw() starts from a known (0, 0) origin and
|
||||
# re-renders every cell rather than diffing against stale.
|
||||
renderer.reset(leave_alternate_screen=False)
|
||||
except OSError as exc:
|
||||
if getattr(exc, "errno", None) == errno.EIO:
|
||||
self._mark_terminal_io_broken("clear_screen")
|
||||
return
|
||||
pass
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def _recover_after_resize(self, app, original_on_resize) -> None:
|
||||
"""Recover a resized classic CLI without desynchronizing cursor state.
|
||||
|
||||
Unlike _force_full_redraw, we do NOT clear the physical screen or
|
||||
scrollback here. The startup banner and tool summary are printed
|
||||
before prompt_toolkit owns the live chrome, so they live in normal
|
||||
terminal scrollback. Erasing the screen on SIGWINCH removes that
|
||||
startup UI and ``_replay_output_history`` cannot reconstruct it
|
||||
(the banner was never added to ``_OUTPUT_HISTORY``).
|
||||
|
||||
Let prompt_toolkit's own resize path run with its renderer cursor
|
||||
cache intact. Its Application._on_resize() starts with
|
||||
renderer.erase(leave_alternate_screen=False), which needs the cached
|
||||
cursor position to move back to the live prompt origin before
|
||||
erase_down(). Resetting the renderer before that erase loses the
|
||||
origin and can leave stale prompt glyphs after a narrow resize.
|
||||
|
||||
We also flag ``_status_bar_suppressed_after_resize`` so the dynamic
|
||||
status bar and input separator rules stay hidden while the terminal
|
||||
reflow settles. On column shrink the terminal reflows already-rendered
|
||||
status bar rows into scrollback before prompt_toolkit can erase them;
|
||||
drawing a fresh full-width bar immediately makes the old and new
|
||||
versions look duplicated (#19280, #22976).
|
||||
|
||||
Suppression alone is not enough on a WIDTH change. prompt_toolkit's
|
||||
``renderer.erase()`` does ``cursor_up(_cursor_pos.y)`` + ``erase_down()``
|
||||
using the ``_cursor_pos.y`` cached from the LAST render at the OLD
|
||||
width (renderer.py). When the column count shrinks, the terminal
|
||||
reflows each already-painted full-width chrome row into 2+ physical
|
||||
rows, so the cached ``y`` undershoots: ``cursor_up`` does not climb
|
||||
past the reflowed rows and ``erase_down`` leaves the stale bar stranded
|
||||
ABOVE the live origin. The next paint then stacks a fresh bar below it
|
||||
— the duplicated-status-bar report (two bars, two elapsed readings).
|
||||
Suppression hides the *new* bar but never erases the already-reflowed
|
||||
*old* one, so the ghost survives the whole suppression window.
|
||||
|
||||
Fix: on a width change, wipe the visible viewport with ``erase_screen``
|
||||
(CSI 2J) BEFORE delegating to prompt_toolkit's resize, then let its
|
||||
repaint redraw from a clean origin. This is banner-safe: 2J clears
|
||||
only the visible screen, NOT scrollback history (that is CSI 3J, which
|
||||
we do not send here — ``rebuild_scrollback=False``), so the startup
|
||||
banner that scrolled into history is preserved and
|
||||
``_replay_output_history`` is not needed. Row-count-only changes skip
|
||||
the clear (no reflow, so no ghost) to avoid an unnecessary repaint.
|
||||
|
||||
The suppression is transient: a short follow-up timer clears it and
|
||||
repaints once the reflow has settled, so the bar returns on its own
|
||||
during idle. Previously the flag was only cleared on the next
|
||||
*submitted* user input, so a resize/reflow (tmux pane change, SSH
|
||||
window restore, font zoom) followed by idle left the status bar hidden
|
||||
indefinitely even while the refresh clock kept ticking (the dynamic
|
||||
chrome rendered at height 0 on every repaint). The next-submit clear
|
||||
at the input loop remains as a fast path.
|
||||
"""
|
||||
from cli import _replay_output_history
|
||||
self._status_bar_suppressed_after_resize = True
|
||||
# On a WIDTH change the terminal has already reflowed the old full-width
|
||||
# chrome into extra physical rows that prompt_toolkit's stale-cursor
|
||||
# erase (cursor_up(_cursor_pos.y) cached at the OLD width) will not
|
||||
# reach, leaving a duplicated status bar stranded above the live origin.
|
||||
# Ctrl+L / /redraw clears it cleanly, so route the resize path through
|
||||
# the SAME recovery: wipe the visible viewport (banner-safe — CSI 2J
|
||||
# by default; CSI 3J only when display.cli_rebuild_scrollback_on_redraw
|
||||
# is enabled) and replay the transcript so nothing is lost.
|
||||
# Same-width SIGWINCH (tmux attach, benign focus/tab signals) is left
|
||||
# untouched — no clear, no replay — because a 2J without replay erases
|
||||
# the visible transcript and a replay against preserved scrollback
|
||||
# duplicates it (#65293). The stale-previous_screen crash tmux attach
|
||||
# used to trigger is handled by _hermes_call_output_screen_diff's
|
||||
# retry-with-first-paint instead (#83874).
|
||||
try:
|
||||
new_width = self._get_tui_terminal_width()
|
||||
except Exception:
|
||||
new_width = None
|
||||
prev_width = getattr(self, "_last_resize_width", None)
|
||||
# Replay only on an OBSERVED width change. The first signal of a
|
||||
# session must not count as one (#65293): GNOME Terminal and friends
|
||||
# deliver benign SIGWINCHes (tab bar appearing, monitor-scale change,
|
||||
# focus events), and a 2J+replay against preserved scrollback
|
||||
# duplicates everything ``_OUTPUT_HISTORY`` holds — after a resume
|
||||
# that is the entire "Previous Conversation" recap plus the first
|
||||
# live exchange. ``_install_resize_recovery`` seeds the baseline at
|
||||
# startup, so an initial maximize/restore still differs from it and
|
||||
# is still recovered; with no baseline (width probe failed) this
|
||||
# signal just records one for the next comparison.
|
||||
width_changed = (
|
||||
new_width is not None
|
||||
and prev_width is not None
|
||||
and new_width != prev_width
|
||||
)
|
||||
if width_changed:
|
||||
try:
|
||||
self._clear_prompt_toolkit_screen(
|
||||
app,
|
||||
rebuild_scrollback=self._redraw_rebuilds_scrollback(),
|
||||
)
|
||||
_replay_output_history()
|
||||
except Exception:
|
||||
pass
|
||||
if new_width is not None:
|
||||
self._last_resize_width = new_width
|
||||
if width_changed:
|
||||
self._pet_queue_kitty_frame()
|
||||
original_on_resize()
|
||||
self._schedule_status_bar_unsuppress(app)
|
||||
|
||||
def _schedule_status_bar_unsuppress(self, app, delay: float = 0.35) -> None:
|
||||
"""Clear the post-resize status-bar suppression after the reflow settles.
|
||||
|
||||
Debounced: a fresh resize cancels the pending unsuppress and restarts
|
||||
the timer, so a resize storm only repaints the bar once it stops.
|
||||
"""
|
||||
try:
|
||||
old_timer = getattr(self, "_status_bar_unsuppress_timer", None)
|
||||
if old_timer is not None:
|
||||
try:
|
||||
old_timer.cancel()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def _clear():
|
||||
self._status_bar_suppressed_after_resize = False
|
||||
try:
|
||||
app.invalidate()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def _fire():
|
||||
try:
|
||||
loop = getattr(app, "loop", None)
|
||||
except Exception:
|
||||
loop = None
|
||||
if loop is not None:
|
||||
try:
|
||||
loop.call_soon_threadsafe(_clear)
|
||||
return
|
||||
except Exception:
|
||||
pass
|
||||
_clear()
|
||||
|
||||
timer = threading.Timer(delay, _fire)
|
||||
timer.daemon = True
|
||||
self._status_bar_unsuppress_timer = timer
|
||||
timer.start()
|
||||
except Exception:
|
||||
# Fail open: never leave the bar stuck hidden.
|
||||
self._status_bar_suppressed_after_resize = False
|
||||
|
||||
def _schedule_resize_recovery(self, app, original_on_resize, delay: float = 0.12) -> None:
|
||||
"""Debounce resize redraws so footer chrome is not stamped into scrollback."""
|
||||
try:
|
||||
old_timer = getattr(self, "_resize_recovery_timer", None)
|
||||
lock = getattr(self, "_resize_recovery_lock", None)
|
||||
if lock is None:
|
||||
lock = threading.Lock()
|
||||
self._resize_recovery_lock = lock
|
||||
|
||||
def _timer_fired(timer_ref):
|
||||
def _run_recovery():
|
||||
with lock:
|
||||
if getattr(self, "_resize_recovery_timer", None) is not timer_ref:
|
||||
return
|
||||
self._resize_recovery_timer = None
|
||||
self._resize_recovery_pending = False
|
||||
self._recover_after_resize(app, original_on_resize)
|
||||
|
||||
try:
|
||||
loop = app.loop # type: ignore[attr-defined]
|
||||
except Exception:
|
||||
loop = None
|
||||
if loop is not None:
|
||||
try:
|
||||
loop.call_soon_threadsafe(_run_recovery)
|
||||
return
|
||||
except Exception:
|
||||
pass
|
||||
_run_recovery()
|
||||
|
||||
with lock:
|
||||
if old_timer is not None:
|
||||
try:
|
||||
old_timer.cancel()
|
||||
except Exception:
|
||||
pass
|
||||
self._resize_recovery_pending = True
|
||||
timer = threading.Timer(delay, lambda: _timer_fired(timer))
|
||||
timer.daemon = True
|
||||
self._resize_recovery_timer = timer
|
||||
timer.start()
|
||||
except Exception:
|
||||
self._resize_recovery_pending = False
|
||||
self._recover_after_resize(app, original_on_resize)
|
||||
|
||||
def _install_resize_recovery(self, app) -> None:
|
||||
"""Route prompt_toolkit's ``_on_resize`` through the debounced
|
||||
ghost-clearing recovery (#5474/#49120) and record the current terminal
|
||||
width as the baseline for width-change detection.
|
||||
|
||||
Seeding the baseline here is what keeps the session's FIRST SIGWINCH
|
||||
honest (#65293): ``_recover_after_resize`` replays the transcript only
|
||||
on an observed width change, and without a startup baseline it could
|
||||
not tell a benign signal (GNOME Terminal tab bar, monitor-scale
|
||||
change) from a real one. An initial maximize/restore still differs
|
||||
from the seeded width, so it is still recovered.
|
||||
|
||||
The probe reads ``app.output`` directly — NOT
|
||||
``_get_tui_terminal_width`` — because this runs before ``app.run()``,
|
||||
when ``get_app()`` still returns prompt_toolkit's DummyApplication
|
||||
whose DummyOutput reports a hardcoded 80 columns; seeding that fake
|
||||
width would make the first real signal look like a width change and
|
||||
resurrect the duplicate-replay bug this exists to fix.
|
||||
``app.output`` is the same object the running app's resize handler
|
||||
measures, so install-time and signal-time widths are comparable.
|
||||
"""
|
||||
width = None
|
||||
try:
|
||||
width = app.output.get_size().columns
|
||||
except Exception:
|
||||
width = None
|
||||
if not width or width <= 0:
|
||||
try:
|
||||
width = shutil.get_terminal_size((80, 24)).columns
|
||||
except Exception:
|
||||
width = None
|
||||
self._last_resize_width = width
|
||||
original_on_resize = app._on_resize
|
||||
|
||||
def _resize_clear_ghosts():
|
||||
self._schedule_resize_recovery(app, original_on_resize)
|
||||
|
||||
app._on_resize = _resize_clear_ghosts
|
||||
|
||||
def _try_attach_clipboard_image(self) -> bool:
|
||||
"""Check clipboard for an image and attach it if found.
|
||||
|
||||
Saves the image to ~/.hermes/images/ and appends the path to
|
||||
``_attached_images``. Returns True if an image was attached.
|
||||
"""
|
||||
from cli import datetime
|
||||
from hermes_cli.clipboard import save_clipboard_image
|
||||
|
||||
img_dir = get_hermes_home() / "images"
|
||||
self._image_counter += 1
|
||||
ts = datetime.now().strftime("%Y%m%d_%H%M%S")
|
||||
img_path = img_dir / f"clip_{ts}_{self._image_counter}.png"
|
||||
|
||||
if save_clipboard_image(img_path):
|
||||
self._attached_images.append(img_path)
|
||||
return True
|
||||
self._image_counter -= 1
|
||||
return False
|
||||
|
||||
def _write_osc52_clipboard(self, text: str) -> None:
|
||||
"""Copy *text* to terminal clipboard via OSC 52.
|
||||
|
||||
Wrapped for tmux/screen passthrough (mirrors the TUI's
|
||||
wrapForMultiplexer in ui-tui/src/lib/osc52.ts) — without the DCS
|
||||
wrapper the multiplexer consumes the sequence and the copy is
|
||||
silently lost.
|
||||
"""
|
||||
payload = base64.b64encode(text.encode("utf-8")).decode("ascii")
|
||||
seq = f"\x1b]52;c;{payload}\x07"
|
||||
if os.environ.get("TMUX"):
|
||||
seq = "\x1bPtmux;" + seq.replace("\x1b", "\x1b\x1b") + "\x1b\\"
|
||||
elif os.environ.get("STY"):
|
||||
seq = "\x1bP" + seq + "\x1b\\"
|
||||
out = getattr(self, "_app", None)
|
||||
output = getattr(out, "output", None) if out else None
|
||||
if output and hasattr(output, "write_raw"):
|
||||
output.write_raw(seq)
|
||||
output.flush()
|
||||
return
|
||||
if output and hasattr(output, "write"):
|
||||
output.write(seq)
|
||||
output.flush()
|
||||
return
|
||||
sys.stdout.write(seq)
|
||||
sys.stdout.flush()
|
||||
|
||||
def _recover_terminal_input_modes(self, *, reason: str) -> None:
|
||||
"""Best-effort reset when leaked mouse reports indicate mode drift."""
|
||||
from cli import (
|
||||
CLI_CONFIG,
|
||||
_DIM,
|
||||
_RST,
|
||||
_TERMINAL_INPUT_MODE_RESET_SEQ,
|
||||
_cli_multiline_shortcuts_enabled,
|
||||
_cprint,
|
||||
_enable_extended_enter_keys,
|
||||
logger,
|
||||
)
|
||||
now = time.monotonic()
|
||||
# Rate-limit to avoid thrashing if a terminal floods reports.
|
||||
if now - self._last_input_mode_recovery < 0.5:
|
||||
return
|
||||
self._last_input_mode_recovery = now
|
||||
|
||||
out = getattr(self, "_app", None)
|
||||
output = getattr(out, "output", None) if out else None
|
||||
try:
|
||||
if output and hasattr(output, "write_raw"):
|
||||
output.write_raw(_TERMINAL_INPUT_MODE_RESET_SEQ)
|
||||
output.flush()
|
||||
elif output and hasattr(output, "write"):
|
||||
output.write(_TERMINAL_INPUT_MODE_RESET_SEQ)
|
||||
output.flush()
|
||||
else:
|
||||
sys.stdout.write(_TERMINAL_INPUT_MODE_RESET_SEQ)
|
||||
sys.stdout.flush()
|
||||
except Exception:
|
||||
return
|
||||
|
||||
# The reset sequence above pops kitty keyboard mode and resets
|
||||
# modifyOtherKeys too — re-request extended keys so Shift+Enter /
|
||||
# modified-key reporting isn't silently dead for the rest of the
|
||||
# session after a recovery (sibling of the startup push).
|
||||
try:
|
||||
if _cli_multiline_shortcuts_enabled(self.config or CLI_CONFIG):
|
||||
_enable_extended_enter_keys(output)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
logger.warning("Recovered terminal input modes after leak: %s", reason)
|
||||
if not self._input_mode_recovery_notice_shown:
|
||||
self._input_mode_recovery_notice_shown = True
|
||||
_cprint(
|
||||
f" {_DIM}Recovered terminal input modes after leaked mouse reports. "
|
||||
f"If this repeats, run /new or restart this tab.{_RST}"
|
||||
)
|
||||
|
||||
def _check_termios_drift(self) -> None:
|
||||
"""Watchdog: heal the tty if it drifted back to cooked mode.
|
||||
|
||||
See ``_heal_cooked_mode_drift`` for the failure class (a lost
|
||||
``run_in_terminal`` cooked→raw restore leaves the terminal
|
||||
line-buffering keystrokes while the prompt_toolkit app believes it
|
||||
owns raw mode — the CLI looks dead but the process is healthy).
|
||||
|
||||
Called from ``process_loop``'s idle branch, so a drifted terminal
|
||||
self-heals within ~a second of the agent going idle instead of
|
||||
requiring an external ``stty`` rescue. Skipped while a
|
||||
``run_in_terminal`` window is legitimately holding cooked mode
|
||||
(``app._running_in_terminal``), while the agent is running (approval
|
||||
prompts and sudo prompts legitimately manipulate the tty), and on
|
||||
Windows (no termios).
|
||||
"""
|
||||
from cli import _DIM, _RST, _cprint, _heal_cooked_mode_drift, logger
|
||||
if os.name == "nt":
|
||||
return
|
||||
app = getattr(self, "_app", None)
|
||||
if app is None or not getattr(app, "_is_running", False):
|
||||
return
|
||||
# A run_in_terminal window is *supposed* to be cooked — don't fight it.
|
||||
if getattr(app, "_running_in_terminal", False):
|
||||
return
|
||||
now = time.monotonic()
|
||||
if now - self._last_termios_drift_check < 1.0:
|
||||
return
|
||||
self._last_termios_drift_check = now
|
||||
try:
|
||||
if not sys.stdin.isatty():
|
||||
return
|
||||
fd = sys.stdin.fileno()
|
||||
except Exception:
|
||||
return
|
||||
if _heal_cooked_mode_drift(fd):
|
||||
logger.warning(
|
||||
"Healed cooked-mode termios drift on stdin — a "
|
||||
"run_in_terminal cooked→raw restore was lost."
|
||||
)
|
||||
# Redraw so the prompt is visibly alive again.
|
||||
try:
|
||||
self._invalidate()
|
||||
except Exception:
|
||||
pass
|
||||
if not self._termios_drift_notice_shown:
|
||||
self._termios_drift_notice_shown = True
|
||||
_cprint(
|
||||
f" {_DIM}Recovered terminal from cooked-mode drift "
|
||||
f"(input should respond normally again).{_RST}"
|
||||
)
|
||||
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
@@ -128,7 +128,7 @@ def test_handle_enter_never_gates_on_command_running():
|
||||
``_command_blocks_input`` check inside ``handle_enter`` would let a
|
||||
future edit silently drop type-ahead submissions during /compress.
|
||||
"""
|
||||
cli_path = Path(__file__).resolve().parents[2] / "cli.py"
|
||||
cli_path = Path(__file__).resolve().parents[2] / "hermes_cli" / "cli_tui_mixin.py"
|
||||
tree = ast.parse(cli_path.read_text(encoding="utf-8"))
|
||||
|
||||
target = None
|
||||
@@ -136,7 +136,7 @@ def test_handle_enter_never_gates_on_command_running():
|
||||
if isinstance(node, ast.FunctionDef) and node.name == "_tui_handle_enter":
|
||||
target = node
|
||||
break
|
||||
assert target is not None, "handle_enter closure not found in cli.py"
|
||||
assert target is not None, "_tui_handle_enter not found in cli_tui_mixin.py"
|
||||
|
||||
offenders = [
|
||||
node.attr
|
||||
|
||||
@@ -27,7 +27,7 @@ from pathlib import Path
|
||||
|
||||
def _load_handle_enter_node() -> ast.FunctionDef:
|
||||
"""Extract the ``handle_enter`` nested function node from cli.py."""
|
||||
cli_path = Path(__file__).resolve().parents[2] / "cli.py"
|
||||
cli_path = Path(__file__).resolve().parents[2] / "hermes_cli" / "cli_tui_mixin.py"
|
||||
tree = ast.parse(cli_path.read_text(encoding="utf-8"))
|
||||
|
||||
target = None
|
||||
@@ -35,7 +35,7 @@ def _load_handle_enter_node() -> ast.FunctionDef:
|
||||
if isinstance(node, ast.FunctionDef) and node.name == "_tui_handle_enter":
|
||||
target = node
|
||||
break
|
||||
assert target is not None, "handle_enter closure not found in cli.py"
|
||||
assert target is not None, "_tui_handle_enter not found in cli_tui_mixin.py"
|
||||
return target
|
||||
|
||||
|
||||
|
||||
@@ -54,8 +54,15 @@ def _make_cli(tool_progress="all", verbose=_UNSET):
|
||||
with patch.object(mod, "get_tool_definitions", return_value=[]), \
|
||||
patch.dict(mod.__dict__, {"CLI_CONFIG": _clean_config}):
|
||||
if verbose is _UNSET:
|
||||
return mod.HermesCLI()
|
||||
return mod.HermesCLI(verbose=verbose)
|
||||
inst = mod.HermesCLI()
|
||||
else:
|
||||
inst = mod.HermesCLI(verbose=verbose)
|
||||
# patch.dict(sys.modules) above restores the pre-import state on exit, which
|
||||
# DROPS ``cli`` when this was the first import. The mixin handlers resolve
|
||||
# ``_cprint`` via a lazy ``from cli import ...``, so ``cli`` must stay
|
||||
# registered for ``patch.object(_cli_mod, "_cprint")`` to intercept.
|
||||
sys.modules["cli"] = mod
|
||||
return inst
|
||||
|
||||
|
||||
class TestToolProgressScrollback:
|
||||
|
||||
Reference in New Issue
Block a user