refactor(cli): move cli.py module-level helper clusters into topical siblings; cli.py lands under 2,000 lines (#116911)
Mechanical, behaviour-neutral extraction along the existing hermes_cli/cli_*.py
pattern. Six clusters of module-level helpers leave cli.py:
- cli_config_load.py prefill messages, reasoning/service-tier parsing, terminal
env mirroring, CLI defaults + YAML merge, logging bootstrap
- cli_render.py reasoning-tag stripping, ANSI/skin colours, light-mode
detection, markdown/final rendering, output-history
recording, _cprint, ChatConsole, compact banner, panel wrap
- cli_terminal_input.py file drops/attachments, bracketed-paste patch, extended
Enter keys, CPR guards, TUI input height, query images
- cli_shutdown.py process session-id sync, exit watchdog, cleanup steps,
session-finalize notifications, one-shot finalize
- cli_single_query.py kanban goal loops, exit-code mapping, quiet -q runner,
image routing, signal handlers, single-query mode
- cli_auto_maintenance.py state-db / checkpoint startup maintenance
Every moved body is AST-identical to the base copy modulo two mechanical
edits: cli-level names are late-bound with a call-time `from cli import ...`
(the cli_init_mixin pattern) and mutable cli state is read through `_cli().NAME`
(the gateway_service_unit `_gw()` pattern), so every monkeypatch seam on the
`cli` facade still intercepts moved code and no read snapshots stale state.
Mutable module state and every `global`-writing function stay in cli.py; the
facade re-exports the moved names in one import block.
test_bracketed_paste_timeout AST-loads the paste helper from its new home.
cli.py 3959 -> 1820 lines.
This commit is contained in:
@@ -8,7 +8,14 @@ Applies on top of the root `AGENTS.md`. Long-form: `website/docs/developer-guide
|
||||
`hermes_cli/cli_commands_mixin.py`, `cli_stream_mixin.py`, `cli_status_bar_mixin.py`,
|
||||
`cli_billing_mixin.py`, `cli_tui_mixin.py` (widgets, keybindings, panels), `cli_tui_runtime_mixin.py`
|
||||
(run-loop phases: input dispatch, startup, signals, shutdown), `cli_init_mixin.py` (the `__init__`
|
||||
phases), ... **Rich** renders banner/panels; **prompt_toolkit**
|
||||
phases), ... Module-level helpers live in topical siblings that `cli.py` re-exports:
|
||||
`cli_config_load.py` (defaults + YAML merge, env mirroring), `cli_render.py` (ANSI/skin colours,
|
||||
light mode, markdown, `_cprint`, panel wrap), `cli_terminal_input.py` (file drops, paste/Enter-key
|
||||
sequences, CPR guards), `cli_shutdown.py` (exit watchdog, cleanup steps, one-shot finalize),
|
||||
`cli_single_query.py` (`-q` runner, exit codes, kanban loops), `cli_auto_maintenance.py` (state-db/checkpoint maintenance).
|
||||
Moved bodies late-bind cli-level names via `from cli import ...` at call time, so patch seams on the
|
||||
`cli` facade still intercept them; mutable module state (`_cleanup_done`, `_OUTPUT_HISTORY`,
|
||||
`_LIGHT_MODE_CACHE`, ...) and every `global`-writing function stay in `cli.py`. **Rich** renders banner/panels; **prompt_toolkit**
|
||||
handles input + autocomplete; `KawaiiSpinner` (`agent/display.py`) animates API calls and prints
|
||||
the `┊` activity feed. `load_cli_config()` in `cli.py` merges CLI defaults + user YAML.
|
||||
`process_command()` resolves the canonical name via `resolve_command()` then dispatches through
|
||||
|
||||
74
hermes_cli/cli_auto_maintenance.py
Normal file
74
hermes_cli/cli_auto_maintenance.py
Normal file
@@ -0,0 +1,74 @@
|
||||
"""Startup auto-maintenance for the classic CLI: state-db pack/vacuum and checkpoint pruning run off the prompt thread.
|
||||
|
||||
Split out of ``cli.py``; ``cli`` re-exports every public name and moved bodies late-bind
|
||||
cli-level names through ``from cli import ...`` at call time so facade monkeypatch seams hold.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import threading
|
||||
|
||||
# Log-record parity with the origin module.
|
||||
logger = logging.getLogger("cli")
|
||||
|
||||
|
||||
def _run_state_db_auto_maintenance(session_db) -> None:
|
||||
"""One-time repairs + auto-archive/prune/vacuum per the ``sessions:`` config. Never raises."""
|
||||
if session_db is None:
|
||||
return
|
||||
try:
|
||||
from hermes_cli.config import load_config as _load_full_config
|
||||
from hermes_constants import get_hermes_home as _get_hermes_home # lazy: tests patch it
|
||||
_hermes_home_maint = _get_hermes_home()
|
||||
|
||||
# One-time repairs, each latched in state_meta once it has run.
|
||||
for meta_key, repair, done_msg, skip_msg in (
|
||||
(
|
||||
"ghost_session_prune_v1",
|
||||
lambda: session_db.prune_empty_ghost_sessions(sessions_dir=_hermes_home_maint / "sessions"),
|
||||
"Pruned %d empty TUI ghost sessions", "Ghost session prune skipped: %s",
|
||||
),
|
||||
(
|
||||
"orphaned_compression_finalize_v1",
|
||||
session_db.finalize_orphaned_compression_sessions,
|
||||
"Finalized %d orphaned compression sessions", "Orphan compression finalize skipped: %s",
|
||||
),
|
||||
):
|
||||
try:
|
||||
if not session_db.get_meta(meta_key):
|
||||
count = repair()
|
||||
session_db.set_meta(meta_key, "1")
|
||||
if count:
|
||||
logger.info(done_msg, count)
|
||||
except Exception as _exc:
|
||||
logger.debug(skip_msg, _exc)
|
||||
|
||||
cfg = (_load_full_config().get("sessions") or {})
|
||||
|
||||
# Auto-archive is independent of auto_prune: run it before prune's early return.
|
||||
if cfg.get("auto_archive", False):
|
||||
session_db.maybe_auto_archive(
|
||||
idle_days=float(cfg.get("auto_archive_days", 3)),
|
||||
min_interval_hours=int(cfg.get("min_interval_hours", 24)),
|
||||
)
|
||||
|
||||
if not cfg.get("auto_prune", False):
|
||||
return
|
||||
session_db.maybe_auto_prune_and_vacuum(
|
||||
retention_days=int(cfg.get("retention_days", 90)),
|
||||
min_interval_hours=int(cfg.get("min_interval_hours", 24)),
|
||||
min_vacuum_interval_days=int(cfg.get("min_vacuum_interval_days", 30)),
|
||||
vacuum=bool(cfg.get("vacuum_after_prune", True)),
|
||||
sessions_dir=_hermes_home_maint / "sessions",
|
||||
)
|
||||
except Exception as exc:
|
||||
logger.debug("state.db auto-maintenance skipped: %s", exc)
|
||||
|
||||
|
||||
def _run_checkpoint_auto_maintenance() -> None:
|
||||
"""Checkpoint store retention on a daemon thread: its ``git gc`` can block for tens of seconds
|
||||
on a large store, which used to stall the prompt once a day. ``auto_prune_from_config`` owns the
|
||||
config gate and the 24h marker and never raises."""
|
||||
from tools.checkpoint_manager import auto_prune_from_config
|
||||
threading.Thread(target=auto_prune_from_config, name="checkpoint-auto-prune", daemon=True).start()
|
||||
318
hermes_cli/cli_config_load.py
Normal file
318
hermes_cli/cli_config_load.py
Normal file
@@ -0,0 +1,318 @@
|
||||
"""Classic-CLI config loading: prefill messages, reasoning/service-tier parsing, terminal env mirroring, CLI defaults + user YAML merge and the logging/display bootstrap.
|
||||
|
||||
Split out of ``cli.py``; ``cli`` re-exports every public name and moved bodies late-bind
|
||||
cli-level names through ``from cli import ...`` at call time so facade monkeypatch seams hold.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
from pathlib import Path
|
||||
from typing import List, Dict, Any
|
||||
from utils import fast_safe_load
|
||||
|
||||
# Log-record parity with the origin module.
|
||||
logger = logging.getLogger("cli")
|
||||
|
||||
|
||||
def _cli():
|
||||
"""Late import of the ``cli`` facade: mutable CLI module state (and its test seams) lives there."""
|
||||
import cli
|
||||
|
||||
return cli
|
||||
|
||||
|
||||
def _load_prefill_messages(file_path: str) -> List[Dict[str, Any]]:
|
||||
"""Load prefill messages (JSON array) from *file_path*; relative to ~/.hermes/; missing/empty -> []."""
|
||||
from cli import _hermes_home
|
||||
if not file_path:
|
||||
return []
|
||||
path = Path(file_path).expanduser()
|
||||
if not path.is_absolute():
|
||||
path = _hermes_home / path
|
||||
if not path.exists():
|
||||
logger.warning("Prefill messages file not found: %s", path)
|
||||
return []
|
||||
try:
|
||||
with open(path, "r", encoding="utf-8") as f:
|
||||
data = json.load(f)
|
||||
if not isinstance(data, list):
|
||||
logger.warning("Prefill messages file must contain a JSON array: %s", path)
|
||||
return []
|
||||
return data
|
||||
except Exception as e:
|
||||
logger.warning("Failed to load prefill messages from %s: %s", path, e)
|
||||
return []
|
||||
|
||||
|
||||
def _resolve_prefill_messages_file(config: Dict[str, Any]) -> str:
|
||||
"""Prefill file path: env, then top-level ``prefill_messages_file``, then legacy ``agent.*``."""
|
||||
agent_cfg = config.get("agent", {})
|
||||
return (
|
||||
os.getenv("HERMES_PREFILL_MESSAGES_FILE", "").strip()
|
||||
or str(config.get("prefill_messages_file", "") or "").strip()
|
||||
or (str(agent_cfg.get("prefill_messages_file", "") or "").strip() if isinstance(agent_cfg, dict) else "")
|
||||
)
|
||||
|
||||
|
||||
def _parse_reasoning_config(effort) -> dict | None:
|
||||
"""Parse a reasoning effort level (string or YAML bool; ``false``/``off`` = disabled)."""
|
||||
from hermes_constants import parse_reasoning_effort
|
||||
result = parse_reasoning_effort(effort)
|
||||
if effort and str(effort).strip() and result is None:
|
||||
logger.warning("Unknown reasoning_effort '%s', using default (medium)", effort)
|
||||
return result
|
||||
|
||||
|
||||
def _parse_service_tier_config(raw: str) -> str | None:
|
||||
"""Parse a persisted fast-mode preference: None, "priority", "auto", or "cold"."""
|
||||
value = str(raw or "").strip().lower()
|
||||
if not value or value in {"normal", "default", "standard", "off", "none"}:
|
||||
return None
|
||||
if value in {"fast", "priority", "on"}:
|
||||
return "priority"
|
||||
if value in {"auto", "cold"}:
|
||||
return value
|
||||
logger.warning("Unknown service_tier '%s', ignoring", raw)
|
||||
return None
|
||||
|
||||
|
||||
# terminal.<key> -> TERMINAL_<KEY> env var. Container-resource keys apply to docker,
|
||||
# singularity, modal, daytona and vercel_sandbox only (ignored for local/ssh).
|
||||
_TERMINAL_ENV_MAPPINGS = {
|
||||
key: f"TERMINAL_{key.upper()}"
|
||||
for key in (
|
||||
"degraded_mode", "cwd", "timeout", "home_mode", "lifetime_seconds", "docker_image",
|
||||
"docker_forward_env", "singularity_image", "modal_image", "daytona_image", "vercel_runtime",
|
||||
"ssh_host", "ssh_user", "ssh_port", "ssh_key", "container_cpu", "container_memory",
|
||||
"container_disk", "container_persistent", "docker_volumes", "docker_env", "docker_extra_args",
|
||||
"docker_shm_size", "docker_mount_cwd_to_workspace", "docker_network", "docker_run_as_host_user",
|
||||
"docker_snap_compat",
|
||||
"docker_persist_across_processes", "docker_shared_container_key", "docker_orphan_reaper",
|
||||
"sandbox_dir", "persistent_shell",
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
_TERMINAL_ENV_MAPPINGS = {"env_type": "TERMINAL_ENV", **_TERMINAL_ENV_MAPPINGS, "sudo_password": "SUDO_PASSWORD"}
|
||||
|
||||
|
||||
# Per-task auxiliary endpoint tuples (config key -> env var).
|
||||
_AUXILIARY_TASK_ENV = {
|
||||
"vision": {
|
||||
"provider": "AUXILIARY_VISION_PROVIDER",
|
||||
"model": "AUXILIARY_VISION_MODEL",
|
||||
"base_url": "AUXILIARY_VISION_BASE_URL",
|
||||
"api_key": "AUXILIARY_VISION_API_KEY",
|
||||
},
|
||||
"approval": {
|
||||
"provider": "AUXILIARY_APPROVAL_PROVIDER",
|
||||
"model": "AUXILIARY_APPROVAL_MODEL",
|
||||
"base_url": "AUXILIARY_APPROVAL_BASE_URL",
|
||||
"api_key": "AUXILIARY_APPROVAL_API_KEY",
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
_CWD_PLACEHOLDERS = (".", "auto", "cwd")
|
||||
|
||||
|
||||
def _mirror_config_to_env(defaults, _file_has_terminal_config):
|
||||
"""Project config.yaml values into the env vars the tool modules read (terminal/browser/auxiliary/security/sessions). Env always wins when already set."""
|
||||
from cli import _AUXILIARY_TASK_ENV, _CWD_PLACEHOLDERS, _TERMINAL_ENV_MAPPINGS
|
||||
terminal_config = defaults.get("terminal", {})
|
||||
|
||||
# "backend" (documented) and legacy "env_type" are both accepted; "backend" wins.
|
||||
if "backend" in terminal_config:
|
||||
terminal_config["env_type"] = terminal_config["backend"]
|
||||
|
||||
# Local backend: cwd is always os.getcwd(). Non-local: a placeholder is popped so
|
||||
# terminal_tool uses its per-backend default; an explicit path is kept.
|
||||
effective_backend = terminal_config.get("env_type", "local")
|
||||
if effective_backend == "local":
|
||||
terminal_config["cwd"] = os.getcwd()
|
||||
defaults["terminal"]["cwd"] = terminal_config["cwd"]
|
||||
elif terminal_config.get("cwd") in _CWD_PLACEHOLDERS:
|
||||
terminal_config.pop("cwd", None)
|
||||
|
||||
# TERMINAL_CWD is force-exported (beats stale .env) except inside a gateway process,
|
||||
# whose config bridge already set it.
|
||||
_is_gateway = os.environ.get("_HERMES_GATEWAY") == "1"
|
||||
for config_key, env_var in _TERMINAL_ENV_MAPPINGS.items():
|
||||
if config_key not in terminal_config:
|
||||
continue
|
||||
val = terminal_config[config_key]
|
||||
if env_var == "TERMINAL_CWD":
|
||||
if not _is_gateway:
|
||||
os.environ[env_var] = str(val)
|
||||
elif _file_has_terminal_config or env_var not in os.environ:
|
||||
os.environ[env_var] = json.dumps(val) if isinstance(val, (list, dict)) else str(val)
|
||||
|
||||
browser_config = defaults.get("browser", {})
|
||||
if "inactivity_timeout" in browser_config:
|
||||
os.environ["BROWSER_INACTIVITY_TIMEOUT"] = str(browser_config["inactivity_timeout"])
|
||||
|
||||
# Only non-empty / non-"auto" auxiliary values are bridged so auto-detection still works.
|
||||
auxiliary_config = defaults.get("auxiliary", {})
|
||||
for task_key, env_map in _AUXILIARY_TASK_ENV.items():
|
||||
task_cfg = auxiliary_config.get(task_key, {})
|
||||
if not isinstance(task_cfg, dict):
|
||||
continue
|
||||
for field, env_var in env_map.items():
|
||||
val = str(task_cfg.get(field, "")).strip()
|
||||
if val and not (field == "provider" and val == "auto"):
|
||||
os.environ[env_var] = val
|
||||
|
||||
security_config = defaults.get("security", {})
|
||||
if isinstance(security_config, dict):
|
||||
redact = security_config.get("redact_secrets")
|
||||
if redact is not None:
|
||||
os.environ["HERMES_REDACT_SECRETS"] = str(redact).lower()
|
||||
|
||||
# Session-search index knobs (hermes_state reads the env carriers).
|
||||
sessions_config = defaults.get("sessions", {})
|
||||
if isinstance(sessions_config, dict):
|
||||
if "cjk_fts" in sessions_config:
|
||||
os.environ["HERMES_CJK_FTS"] = str(sessions_config["cjk_fts"])
|
||||
if "search_slow_ms" in sessions_config:
|
||||
os.environ["HERMES_SEARCH_SLOW_MS"] = str(sessions_config["search_slow_ms"])
|
||||
|
||||
|
||||
def _cli_config_defaults():
|
||||
"""Built-in defaults for every config key the CLI reads (the file overlays these)."""
|
||||
img = "nikolaik/python-nodejs:python3.11-nodejs20"
|
||||
return {
|
||||
"model": {"default": "", "base_url": "", "provider": "auto"},
|
||||
"terminal": {
|
||||
"env_type": "local", "cwd": ".", "home_mode": "auto", "lifetime_seconds": 300, # cwd "." -> os.getcwd()
|
||||
"docker_image": img, "docker_forward_env": [], "singularity_image": f"docker://{img}",
|
||||
"modal_image": img, "daytona_image": img, "docker_volumes": [],
|
||||
"docker_mount_cwd_to_workspace": False, # opt-in only: sandbox isolation
|
||||
"docker_shared_container_key": "",
|
||||
},
|
||||
"browser": {
|
||||
"inactivity_timeout": 120, "record_sessions": False, "engine": "auto", # auto (Chrome) | lightpanda | chrome
|
||||
"camofox": {"rewrite_loopback_urls": False, "loopback_host_alias": "host.docker.internal"},
|
||||
},
|
||||
# threshold: fraction of the model's context limit; min_tail: real user messages kept in the tail
|
||||
"compression": {"enabled": True, "threshold": 0.50, "min_tail_user_messages": 1},
|
||||
"agent": {
|
||||
"max_turns": 500, "verbose": False, "system_prompt": "", "prefill_messages_file": "", # max_turns shared with subagents
|
||||
"reasoning_effort": "", "service_tier": "",
|
||||
"personalities": {}, # user overrides merged by name over hermes_cli.personality builtins
|
||||
},
|
||||
"display": {
|
||||
"compact": False,
|
||||
# /resume recap tuning and show_reasoning: keep in sync with hermes_cli/config.py DEFAULT_CONFIG
|
||||
"resume_display": "full", "resume_exchanges": 10, "resume_max_user_chars": 300,
|
||||
"resume_max_assistant_chars": 200, "resume_max_assistant_lines": 3, "resume_skip_tool_only": True,
|
||||
"show_reasoning": True, "reasoning_full": False, "streaming": True, "busy_input_mode": "interrupt",
|
||||
"persistent_output": True, "persistent_output_max_lines": 200,
|
||||
# Also clear scrollback on redraw/resize recovery; off because users prefer history.
|
||||
"cli_rebuild_scrollback_on_redraw": False,
|
||||
"persist_prompts": True, # one-line summary of resolved modal prompts into scrollback
|
||||
"skin": "default",
|
||||
},
|
||||
"code_execution": {"timeout": 300, "max_tool_calls": 50},
|
||||
"auxiliary": {"vision": {"provider": "auto", "model": "", "base_url": "", "api_key": ""}},
|
||||
# delegation: empty model/provider = inherit parent; api_key falls back to OPENAI_API_KEY
|
||||
"delegation": {"max_iterations": 45, "model": "", "provider": "", "base_url": "", "api_key": ""},
|
||||
"onboarding": {"seen": {}}, # first-touch hint flags (agent/onboarding.py), latched once shown
|
||||
}
|
||||
|
||||
|
||||
def _merge_file_config(defaults: Dict[str, Any], file_config: Dict[str, Any]) -> None:
|
||||
"""Overlay a parsed config file onto *defaults* in place (model normalization, deep merge, legacy keys)."""
|
||||
# model: string (new format) or dict (old format with default/base_url)
|
||||
if "model" in file_config:
|
||||
if isinstance(file_config["model"], str):
|
||||
defaults["model"]["default"] = file_config["model"]
|
||||
elif isinstance(file_config["model"], dict):
|
||||
defaults["model"].update(file_config["model"])
|
||||
# Promote model.model -> model.default (HermesCLI checks "default" first).
|
||||
if "model" in file_config["model"] and "default" not in file_config["model"]:
|
||||
defaults["model"]["default"] = file_config["model"]["model"]
|
||||
|
||||
# Deep-merge dict sections, overwrite scalars; a None section keeps the defaults;
|
||||
# unknown keys (platform_toolsets, memory, ...) are carried over.
|
||||
for key, value in file_config.items():
|
||||
if key == "model":
|
||||
continue
|
||||
if isinstance(defaults.get(key), dict):
|
||||
if isinstance(value, dict):
|
||||
defaults[key].update(value)
|
||||
elif value is not None:
|
||||
defaults[key] = value
|
||||
else:
|
||||
defaults[key] = value
|
||||
|
||||
# Legacy root-level max_turns -> agent.max_turns whenever the nested key is missing.
|
||||
agent_file_config = file_config.get("agent")
|
||||
if "max_turns" in file_config and not (
|
||||
isinstance(agent_file_config, dict) and agent_file_config.get("max_turns") is not None
|
||||
):
|
||||
defaults["agent"]["max_turns"] = file_config["max_turns"]
|
||||
|
||||
|
||||
def load_cli_config() -> Dict[str, Any]:
|
||||
"""~/.hermes/config.yaml (else ./cli-config.yaml) over built-in defaults; env vars win.
|
||||
|
||||
``HERMES_IGNORE_USER_CONFIG=1`` skips the user config entirely (``.env`` still loads).
|
||||
"""
|
||||
from cli import _cli_config_defaults, _hermes_home, _merge_file_config, _mirror_config_to_env
|
||||
config_path = _hermes_home / 'config.yaml'
|
||||
if not config_path.exists() or os.environ.get("HERMES_IGNORE_USER_CONFIG") == "1":
|
||||
config_path = Path(__file__).parent / 'cli-config.yaml'
|
||||
|
||||
defaults = _cli_config_defaults()
|
||||
|
||||
# Only a file's terminal section may overwrite terminal env vars already set by .env.
|
||||
_file_has_terminal_config = False
|
||||
|
||||
if config_path.exists():
|
||||
try:
|
||||
with open(config_path, "r", encoding="utf-8") as f:
|
||||
from hermes_cli.config import _normalize_root_model_keys
|
||||
|
||||
file_config = _normalize_root_model_keys(fast_safe_load(f) or {})
|
||||
|
||||
_file_has_terminal_config = "terminal" in file_config
|
||||
_merge_file_config(defaults, file_config)
|
||||
except Exception as e:
|
||||
logger.warning("Failed to load cli-config.yaml: %s", e)
|
||||
|
||||
# Expand ${ENV_VAR} references before bridging to env vars.
|
||||
from hermes_cli.config import _expand_env_vars
|
||||
defaults = _expand_env_vars(defaults)
|
||||
|
||||
# Administrator-pinned (managed scope) values overlay LAST; cli.py builds its config
|
||||
# independently of hermes_cli.config, so this keeps parity with `hermes config`. Fail-open.
|
||||
from hermes_cli import managed_scope
|
||||
|
||||
defaults = managed_scope.apply_managed_overlay(defaults)
|
||||
|
||||
_mirror_config_to_env(defaults, _file_has_terminal_config)
|
||||
|
||||
return defaults
|
||||
|
||||
|
||||
def _init_logging_and_display_from_config() -> None:
|
||||
"""Best-effort startup side effects: logging, config warnings, skin, display knobs."""
|
||||
from importlib import import_module as _im
|
||||
|
||||
def _display(key, default):
|
||||
return _cli().CLI_CONFIG.get("display", {}).get(key, default)
|
||||
|
||||
for step in (
|
||||
lambda: _im("hermes_logging").setup_logging(mode="cli"),
|
||||
lambda: _im("hermes_cli.config").print_config_warnings(),
|
||||
lambda: _im("hermes_cli.skin_engine").init_skin_from_config(_cli().CLI_CONFIG),
|
||||
lambda: _im("agent.display").set_tool_preview_max_len(int(_display("tool_preview_length", 0) or 0)),
|
||||
lambda: _im("agent.display").set_friendly_tool_labels(bool(_display("friendly_tool_labels", True))),
|
||||
):
|
||||
try:
|
||||
step()
|
||||
except Exception:
|
||||
pass
|
||||
704
hermes_cli/cli_render.py
Normal file
704
hermes_cli/cli_render.py
Normal file
@@ -0,0 +1,704 @@
|
||||
"""Classic-CLI rendering helpers: reasoning-tag stripping, ANSI/skin colours, light-mode detection, markdown/final-content rendering, output-history replay, ``_cprint`` and the panel box/wrap helpers.
|
||||
|
||||
Split out of ``cli.py``; ``cli`` re-exports every public name and moved bodies late-bind
|
||||
cli-level names through ``from cli import ...`` at call time so facade monkeypatch seams hold.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import functools
|
||||
import os
|
||||
import re
|
||||
import shutil
|
||||
import sys
|
||||
import textwrap
|
||||
import time
|
||||
from contextlib import contextmanager, suppress
|
||||
from hermes_cli.banner import format_banner_version_label
|
||||
from rich.console import Console
|
||||
from rich.text import Text as _RichText
|
||||
from typing import Any
|
||||
|
||||
def _cli():
|
||||
"""Late import of the ``cli`` facade: mutable CLI module state (and its test seams) lives there."""
|
||||
import cli
|
||||
|
||||
return cli
|
||||
|
||||
|
||||
_REASONING_TAGS = ("REASONING_SCRATCHPAD", "think", "thinking", "reasoning", "thought")
|
||||
|
||||
|
||||
_TOOL_CALL_TAGS = ("tool_call", "tool_calls", "tool_result", "function_call", "function_calls")
|
||||
|
||||
|
||||
def _strip_reasoning_tags(text: str) -> str:
|
||||
"""Strip reasoning blocks (closed, unterminated, orphan-close) and leaked tool-call XML from display text.
|
||||
|
||||
Keep in sync with ``agent.agent_runtime_helpers.strip_think_blocks`` and the stream consumer's think-tag sets.
|
||||
|
||||
Also strips tool-call XML blocks some open models leak into visible content (``<tool_call>``,
|
||||
``<function_calls>``, Gemma-style ``<function name="…">…</function>``). Ported from
|
||||
openclaw/openclaw#67318.
|
||||
"""
|
||||
from cli import _REASONING_TAGS, _TOOL_CALL_TAGS
|
||||
cleaned = text
|
||||
for tag in _REASONING_TAGS:
|
||||
cleaned = re.sub(rf"<{tag}>.*?</{tag}>\s*", "", cleaned, flags=re.DOTALL | re.IGNORECASE)
|
||||
cleaned = re.sub(rf"<{tag}>.*$", "", cleaned, flags=re.DOTALL | re.IGNORECASE)
|
||||
cleaned = re.sub(rf"</{tag}>\s*", "", cleaned, flags=re.IGNORECASE)
|
||||
for tc_tag in _TOOL_CALL_TAGS:
|
||||
cleaned = re.sub(
|
||||
rf"<(?:[\w.-]+:)?{tc_tag}\b[^>]*>.*?</(?:[\w.-]+:)?{tc_tag}>\s*",
|
||||
"", cleaned, flags=re.DOTALL | re.IGNORECASE,
|
||||
)
|
||||
# <function name="..."> — boundary + attribute gated to avoid prose false positives.
|
||||
cleaned = re.sub(
|
||||
r'(?:(?<=^)|(?<=[\n\r.!?:]))[ \t]*<function\b[^>]*\bname\s*=[^>]*>(?:(?:(?!</function>).)*)</function>\s*',
|
||||
'', cleaned, flags=re.DOTALL | re.IGNORECASE,
|
||||
)
|
||||
cleaned = re.sub(
|
||||
r'</(?:(?:[\w.-]+:)?(?:tool_call|tool_calls|tool_result|function_call|function_calls|function))>\s*', '', cleaned,
|
||||
flags=re.IGNORECASE,
|
||||
)
|
||||
# Unterminated opener / stray <arg_key>/<arg_value> markup = stream cut
|
||||
# mid tool-call serialization (#101899); strip to end of text.
|
||||
cleaned = re.sub(
|
||||
r'(?:^|\n)[ \t]*<(?:[\w.-]+:)?(?:tool_call|tool_calls|tool_result|function_call|function_calls)\b[^>]*>.*$'
|
||||
r'|(?:^|\n)[^\n<]*</?arg_(?:key|value)\b.*$',
|
||||
'',
|
||||
cleaned,
|
||||
flags=re.DOTALL | re.IGNORECASE,
|
||||
)
|
||||
return cleaned.strip()
|
||||
|
||||
|
||||
def _assistant_content_as_text(content: Any) -> str:
|
||||
if content is None:
|
||||
return ""
|
||||
if isinstance(content, str):
|
||||
return content
|
||||
if isinstance(content, list):
|
||||
parts = [str(part.get("text", "")) for part in content if isinstance(part, dict) and part.get("type") == "text"]
|
||||
return "\n".join(p for p in parts if p)
|
||||
return str(content)
|
||||
|
||||
|
||||
def _assistant_copy_text(content: Any) -> str:
|
||||
from cli import _assistant_content_as_text, _strip_reasoning_tags
|
||||
return _strip_reasoning_tags(_assistant_content_as_text(content))
|
||||
|
||||
|
||||
_ACCENT_ANSI_DEFAULT = "\033[1;38;2;255;215;0m" # #FFD700 bold fallback
|
||||
|
||||
|
||||
_BOLD = "\033[1m"
|
||||
|
||||
|
||||
_RST = "\033[0m"
|
||||
|
||||
|
||||
_STREAM_PAD = "" # no indent: leading whitespace pollutes copy/paste
|
||||
|
||||
|
||||
_STREAM_PARTIAL_PREVIEW_LEN = 60 # tail of an unfinished line mirrored into the spinner
|
||||
|
||||
|
||||
def _hex_to_ansi(hex_color: str, *, bold: bool = False) -> str:
|
||||
"""Convert '#RRGGBB' to a true-color ANSI escape, remapping dark-tuned colors in light mode."""
|
||||
from cli import _ACCENT_ANSI_DEFAULT, _maybe_remap_for_light_mode
|
||||
hex_color = _maybe_remap_for_light_mode(hex_color)
|
||||
try:
|
||||
r, g, b = (int(hex_color[i:i + 2], 16) for i in (1, 3, 5))
|
||||
return f"\033[{'1;' if bold else ''}38;2;{r};{g};{b}m"
|
||||
except (ValueError, IndexError):
|
||||
return _ACCENT_ANSI_DEFAULT if bold else "\033[38;2;184;134;11m"
|
||||
|
||||
|
||||
_TRUE_RE = re.compile(r"^(1|true|on|yes|y)$")
|
||||
|
||||
|
||||
_FALSE_RE = re.compile(r"^(0|false|off|no|n)$")
|
||||
|
||||
|
||||
_LIGHT_DEFAULT_TERM_PROGRAMS = frozenset() # Apple_Terminal isn't reliable; require explicit config
|
||||
|
||||
|
||||
def _luminance_from_hex(hex_str: str) -> float | None:
|
||||
"""Rec.709 luma in [0, 1] for '#RGB'/'#RRGGBB', or None when malformed."""
|
||||
s = (hex_str or "").strip().lstrip("#")
|
||||
if len(s) == 3:
|
||||
s = "".join(c * 2 for c in s)
|
||||
if len(s) != 6 or not all(c in "0123456789abcdefABCDEF" for c in s):
|
||||
return None
|
||||
try:
|
||||
r, g, b = int(s[0:2], 16), int(s[2:4], 16), int(s[4:6], 16)
|
||||
except ValueError:
|
||||
return None
|
||||
return (0.2126 * r + 0.7152 * g + 0.0722 * b) / 255.0
|
||||
|
||||
|
||||
_DA1_REPLY_RE = re.compile(rb"\x1b\[\?[0-9;]*c")
|
||||
|
||||
|
||||
def _query_osc11_background() -> str | None:
|
||||
"""Terminal background via OSC 11 as "#RRGGBB", or None.
|
||||
|
||||
Fenced with a DA1 sentinel (``ESC[c``): terminals answer in order and virtually all
|
||||
answer DA1, so its reply proves our OSC 11 was processed — otherwise a late reply
|
||||
leaks into prompt_toolkit's stdin as typed text. Skipped over SSH (round-trip too
|
||||
slow; a late BEL reads as Ctrl+G). A 50 ms drain after TCSAFLUSH catches stragglers.
|
||||
|
||||
After the main read + TCSAFLUSH, a short drain window (50 ms) catches late-arriving bytes that slipped
|
||||
past the flush — a race observed on VPS and container terminals under load (#40250).
|
||||
"""
|
||||
from cli import _DA1_REPLY_RE
|
||||
if not sys.stdin.isatty() or not sys.stdout.isatty():
|
||||
return None
|
||||
if any(os.environ.get(v) for v in ("SSH_CONNECTION", "SSH_CLIENT", "SSH_TTY")):
|
||||
return None
|
||||
try:
|
||||
import select
|
||||
import termios
|
||||
import tty
|
||||
fd = sys.stdin.fileno()
|
||||
old = termios.tcgetattr(fd)
|
||||
except Exception:
|
||||
return None
|
||||
try:
|
||||
try:
|
||||
tty.setcbreak(fd)
|
||||
except Exception:
|
||||
return None
|
||||
try:
|
||||
# One write so the OSC 11 query and DA1 fence cannot reorder.
|
||||
sys.stdout.write("\x1b]11;?\x1b\\\x1b[c")
|
||||
sys.stdout.flush()
|
||||
except Exception:
|
||||
return None
|
||||
# Read until the DA1 fence closes; the 1s deadline only covers terminals ignoring DA1.
|
||||
deadline = time.monotonic() + 1.0
|
||||
buf = b""
|
||||
while time.monotonic() < deadline:
|
||||
r, _, _ = select.select([fd], [], [], deadline - time.monotonic())
|
||||
if not r:
|
||||
continue
|
||||
try:
|
||||
chunk = os.read(fd, 64)
|
||||
except OSError:
|
||||
break
|
||||
if not chunk:
|
||||
break
|
||||
buf += chunk
|
||||
if _DA1_REPLY_RE.search(buf):
|
||||
break
|
||||
# Reply: \x1b]11;rgb:RRRR/GGGG/BBBB\x1b\\ — components are 1-4 hex digits.
|
||||
m = re.search(rb"rgb:([0-9a-fA-F]+)/([0-9a-fA-F]+)/([0-9a-fA-F]+)", buf)
|
||||
if not m:
|
||||
return None
|
||||
|
||||
def norm(h: bytes) -> int:
|
||||
v = int(h, 16)
|
||||
bits = len(h) * 4
|
||||
return (v * 255) // ((1 << bits) - 1) if bits else 0
|
||||
r, g, b = norm(m.group(1)), norm(m.group(2)), norm(m.group(3))
|
||||
return f"#{r:02X}{g:02X}{b:02X}"
|
||||
finally:
|
||||
# TCSAFLUSH discards unread input, scrubbing a partial reply before prompt_toolkit reads it.
|
||||
with suppress(Exception):
|
||||
termios.tcsetattr(fd, termios.TCSAFLUSH, old)
|
||||
try:
|
||||
drain_deadline = time.monotonic() + 0.05
|
||||
while time.monotonic() < drain_deadline:
|
||||
r, _, _ = select.select([fd], [], [], drain_deadline - time.monotonic())
|
||||
if not r or not os.read(fd, 64):
|
||||
break
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def _heal_cooked_mode_drift(fd: int) -> bool:
|
||||
"""Re-apply raw mode on *fd* when termios drifted back to cooked (POSIX only).
|
||||
|
||||
A lost ``run_in_terminal`` cooked_mode() restore makes the kernel line-buffer every
|
||||
keystroke and the CLI looks dead. Mirrors prompt_toolkit's raw_mode flag surgery in
|
||||
place. Returns True when healed; False when already raw or not inspectable.
|
||||
"""
|
||||
try:
|
||||
import termios
|
||||
attrs = termios.tcgetattr(fd)
|
||||
except Exception:
|
||||
return False
|
||||
lflag = attrs[3]
|
||||
if not (lflag & (termios.ICANON | termios.ECHO)):
|
||||
return False # still raw — nothing to do
|
||||
attrs[3] = lflag & ~(termios.ECHO | termios.ICANON | termios.IEXTEN | termios.ISIG)
|
||||
attrs[0] = attrs[0] & ~(termios.IXON | termios.IXOFF | termios.ICRNL | termios.INLCR | termios.IGNCR)
|
||||
attrs[6][termios.VMIN] = 1
|
||||
try:
|
||||
termios.tcsetattr(fd, termios.TCSANOW, attrs)
|
||||
except Exception:
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
def _detect_light_mode_uncached() -> bool:
|
||||
"""The detection ladder documented above; may raise (caller maps errors to dark)."""
|
||||
from cli import _FALSE_RE, _LIGHT_DEFAULT_TERM_PROGRAMS, _TRUE_RE, _luminance_from_hex, _query_osc11_background
|
||||
for var in ("HERMES_LIGHT", "HERMES_TUI_LIGHT"):
|
||||
v = (os.environ.get(var) or "").strip().lower()
|
||||
if _TRUE_RE.match(v):
|
||||
return True
|
||||
if _FALSE_RE.match(v):
|
||||
return False
|
||||
theme = (os.environ.get("HERMES_TUI_THEME") or "").strip().lower()
|
||||
if theme == "light":
|
||||
return True
|
||||
if theme == "dark":
|
||||
return False
|
||||
bg_lum = _luminance_from_hex(os.environ.get("HERMES_TUI_BACKGROUND") or "")
|
||||
if bg_lum is not None:
|
||||
return bg_lum >= 0.5
|
||||
last = (os.environ.get("COLORFGBG") or "").strip().split(";")[-1]
|
||||
if last.isdigit() and 0 <= int(last) < 16:
|
||||
return int(last) in {7, 15}
|
||||
bg_color = _query_osc11_background()
|
||||
if bg_color:
|
||||
lum = _luminance_from_hex(bg_color)
|
||||
if lum is not None:
|
||||
return lum >= 0.5
|
||||
return (os.environ.get("TERM_PROGRAM") or "").strip() in _LIGHT_DEFAULT_TERM_PROGRAMS
|
||||
|
||||
|
||||
# Light-mode equivalents of skin colors unreadable on cream backgrounds. Only colors used
|
||||
# as STANDALONE foregrounds: ones paired with a dark bg (status bar text on #1a1a2e) would
|
||||
# become invisible the other direction, hence #C0C0C0/#888888/#555555/#8B8682 are skipped.
|
||||
_LIGHT_MODE_REMAP: dict[str, str] = {
|
||||
"#FFF8DC": "#1A1A1A", "#FFD700": "#9A6B00", "#FFBF00": "#8A5A00", "#B8860B": "#5C4500",
|
||||
"#DAA520": "#6B4F00", "#F1E6CF": "#1A1A1A", "#c9d1d9": "#24292F", "#EAF7FF": "#0F1B26",
|
||||
"#F5F5F5": "#1A1A1A", "#FFF0D4": "#1A1A1A", "#CD7F32": "#8A4F1A", "#FFEFB5": "#3A2A00",
|
||||
}
|
||||
|
||||
|
||||
_LIGHT_MODE_REMAP_UPPER = {k.upper(): v for k, v in _LIGHT_MODE_REMAP.items()}
|
||||
|
||||
|
||||
def _maybe_remap_for_light_mode(hex_color: str) -> str:
|
||||
"""In light mode, remap a dark-tuned color to its higher-contrast equivalent."""
|
||||
from cli import _LIGHT_MODE_REMAP_UPPER, _detect_light_mode
|
||||
if not _detect_light_mode():
|
||||
return hex_color
|
||||
if not hex_color or not hex_color.startswith("#"):
|
||||
return hex_color
|
||||
return _LIGHT_MODE_REMAP_UPPER.get(hex_color.upper(), hex_color)
|
||||
|
||||
|
||||
def _install_skin_light_mode_hook() -> None:
|
||||
"""Wrap SkinConfig.get_color so EVERY skin color read goes through the light-mode remap. Idempotent."""
|
||||
from cli import _maybe_remap_for_light_mode
|
||||
try:
|
||||
from hermes_cli.skin_engine import SkinConfig # type: ignore[import]
|
||||
except Exception:
|
||||
return
|
||||
if getattr(SkinConfig, "_hermes_light_mode_hook_installed", False):
|
||||
return
|
||||
_orig_get_color = SkinConfig.get_color
|
||||
|
||||
def _wrapped_get_color(self, key, fallback=""):
|
||||
value = _orig_get_color(self, key, fallback)
|
||||
try:
|
||||
return _maybe_remap_for_light_mode(value)
|
||||
except Exception:
|
||||
return value
|
||||
|
||||
SkinConfig.get_color = _wrapped_get_color # type: ignore[method-assign]
|
||||
SkinConfig._hermes_light_mode_hook_installed = True # type: ignore[attr-defined]
|
||||
|
||||
|
||||
class _SkinAwareAnsi:
|
||||
"""Lazy ANSI escape resolved from the skin on first use; ``.reset()`` after a ``/skin`` switch."""
|
||||
|
||||
def __init__(self, skin_key: str, fallback_hex: str = "#FFD700", *, bold: bool = False):
|
||||
self._skin_key = skin_key
|
||||
self._fallback_hex = fallback_hex
|
||||
self._bold = bold
|
||||
self._cached: str | None = None
|
||||
|
||||
def __str__(self) -> str:
|
||||
from cli import _hex_to_ansi
|
||||
if self._cached is None:
|
||||
try:
|
||||
from hermes_cli.skin_engine import get_active_skin
|
||||
self._cached = _hex_to_ansi(
|
||||
get_active_skin().get_color(self._skin_key, self._fallback_hex),
|
||||
bold=self._bold,
|
||||
)
|
||||
except Exception:
|
||||
self._cached = _hex_to_ansi(self._fallback_hex, bold=self._bold)
|
||||
return self._cached
|
||||
|
||||
def __add__(self, other: str) -> str:
|
||||
return str(self) + other
|
||||
|
||||
def __radd__(self, other: str) -> str:
|
||||
return other + str(self)
|
||||
|
||||
def reset(self) -> None:
|
||||
"""Clear cache so the next access re-reads the skin."""
|
||||
self._cached = None
|
||||
|
||||
|
||||
_ACCENT = _SkinAwareAnsi("response_border", "#FFD700", bold=True)
|
||||
|
||||
|
||||
# dim+italic attributes (not a hex) so dim text inherits the terminal foreground in both modes.
|
||||
_DIM = "\x1b[2;3m"
|
||||
|
||||
|
||||
def _tty_wrap(s: str, sgr: str) -> str:
|
||||
"""Wrap *s* in an SGR attribute when stdout is a real TTY; plain text otherwise."""
|
||||
try:
|
||||
return f"{sgr}{s}\x1b[0m" if sys.stdout.isatty() else str(s)
|
||||
except Exception:
|
||||
return str(s)
|
||||
|
||||
|
||||
_b = functools.partial(_tty_wrap, sgr="\x1b[1m") # bold when stdout is a real TTY
|
||||
|
||||
|
||||
_d = functools.partial(_tty_wrap, sgr="\x1b[2;3m") # dim-italic when stdout is a real TTY
|
||||
|
||||
|
||||
def _accent_hex() -> str:
|
||||
"""Return the active skin accent color for legacy CLI output lines."""
|
||||
try:
|
||||
from hermes_cli.skin_engine import get_active_skin
|
||||
return get_active_skin().get_color("ui_accent", "#FFBF00")
|
||||
except Exception:
|
||||
return "#FFBF00"
|
||||
|
||||
|
||||
def _rich_text_from_ansi(text: str) -> _RichText:
|
||||
"""Rich Text from ANSI output; literal ``[brackets]`` are not treated as markup."""
|
||||
return _RichText.from_ansi(text or "")
|
||||
|
||||
|
||||
def _strip_markdown_syntax(text: str) -> str:
|
||||
"""Best-effort markdown marker removal for plain-text display."""
|
||||
from cli import _rich_text_from_ansi
|
||||
plain = _rich_text_from_ansi(text or "").plain
|
||||
# HR markers: "-"/"_" runs of 3+, but "*" only when exactly 3 (cron schedules "* * * * *").
|
||||
plain = re.sub(r"^\s{0,3}(?:[-_]\s*){3,}$", "", plain, flags=re.MULTILINE)
|
||||
plain = re.sub(r"^\s{0,3}(?:\*\s*){3}\s*$", "", plain, flags=re.MULTILINE)
|
||||
plain = re.sub(r"^\s{0,3}#{1,6}\s+", "", plain, flags=re.MULTILINE)
|
||||
# Blockquotes, lists, and checkboxes are preserved because they carry structure.
|
||||
plain = re.sub(r"(```+|~~~+)", "", plain)
|
||||
plain = re.sub(r"`([^`]*)`", r"\1", plain)
|
||||
plain = re.sub(r"!\[([^\]]*)\]\([^\)]*\)", r"\1", plain)
|
||||
plain = re.sub(r"\[([^\]]+)\]\([^\)]*\)", r"\1", plain)
|
||||
plain = re.sub(r"\*\*\*([^*]+)\*\*\*", r"\1", plain)
|
||||
plain = re.sub(r"(?<!\w)___([^_]+)___(?!\w)", r"\1", plain)
|
||||
plain = re.sub(r"\*\*([^*]+)\*\*", r"\1", plain)
|
||||
plain = re.sub(r"(?<!\w)__([^_]+)__(?!\w)", r"\1", plain)
|
||||
# `*emphasis*` only when the inner text is non-whitespace (cron expressions again).
|
||||
plain = re.sub(r"\*([^\s*][^*]*?[^\s*])\*", r"\1", plain)
|
||||
plain = re.sub(r"(?<!\w)_([^_]+)_(?!\w)", r"\1", plain)
|
||||
plain = re.sub(r"~~([^~]+)~~", r"\1", plain)
|
||||
plain = re.sub(r"\n{3,}", "\n\n", plain)
|
||||
return plain.strip("\n")
|
||||
|
||||
|
||||
_WINDOWS_PATH_WITH_DOT_SEGMENT_RE = re.compile(r"(?i)(?:\b[a-z]:\\|\\\\)[^\s`]*\\\.[^\s`]*")
|
||||
|
||||
|
||||
def _preserve_windows_dot_segments_for_markdown(text: str) -> str:
|
||||
r"""Double the ``\`` before hidden dirs in Windows paths: CommonMark reads ``\.`` as an escaped dot."""
|
||||
from cli import _WINDOWS_PATH_WITH_DOT_SEGMENT_RE
|
||||
if "\\." not in text:
|
||||
return text
|
||||
|
||||
def _protect(match: re.Match[str]) -> str:
|
||||
return re.sub(r"(?<!\\)\\(?=\.)", r"\\\\", match.group(0))
|
||||
|
||||
return _WINDOWS_PATH_WITH_DOT_SEGMENT_RE.sub(_protect, text)
|
||||
|
||||
|
||||
def _terminal_columns() -> int:
|
||||
try:
|
||||
return shutil.get_terminal_size((80, 24)).columns
|
||||
except Exception:
|
||||
return 80
|
||||
|
||||
|
||||
def _terminal_width_for_streaming() -> int:
|
||||
"""Display cells inside the streamed response box (small margin for resize races)."""
|
||||
from cli import _STREAM_PAD, _terminal_columns
|
||||
return max(20, _terminal_columns() - len(_STREAM_PAD) - 2)
|
||||
|
||||
|
||||
def _render_final_assistant_content(text: str, mode: str = "render"):
|
||||
"""Render final assistant content as markdown, stripped text, or raw text."""
|
||||
from cli import _preserve_windows_dot_segments_for_markdown, _rich_text_from_ansi, _strip_markdown_syntax, _terminal_columns, realign_markdown_tables
|
||||
from rich.markdown import Markdown
|
||||
|
||||
# 1 border cell each side + margin so resize races don't push a borderline table into soft-wrap.
|
||||
panel_width = max(20, _terminal_columns() - 4)
|
||||
|
||||
normalized_mode = str(mode or "render").strip().lower()
|
||||
if normalized_mode == "strip":
|
||||
# Strip first (inline markdown changes cell width), then re-align padding.
|
||||
return _RichText(realign_markdown_tables(_strip_markdown_syntax(text), panel_width))
|
||||
if normalized_mode == "raw":
|
||||
return _rich_text_from_ansi(text or "")
|
||||
|
||||
# Normalising under-padded tables up front gives narrow-panel fallbacks consistent input.
|
||||
plain = _rich_text_from_ansi(text or "").plain
|
||||
plain = _preserve_windows_dot_segments_for_markdown(plain)
|
||||
plain = realign_markdown_tables(plain, panel_width)
|
||||
return Markdown(plain)
|
||||
|
||||
|
||||
def _post_stream_transform_output(response: str, result: dict | None) -> str:
|
||||
"""Text still to display after a streamed response transform: the suffix, or the whole response when replaced."""
|
||||
if not result or not result.get("response_transformed"):
|
||||
return ""
|
||||
|
||||
original = result.get("pre_transform_response") or ""
|
||||
if original and response.startswith(original):
|
||||
return response[len(original):]
|
||||
|
||||
return f"\n[Response transformed after streaming]\n{response}"
|
||||
|
||||
|
||||
def _coerce_output_history_limit(value) -> int:
|
||||
try:
|
||||
return max(10, int(value))
|
||||
except (TypeError, ValueError):
|
||||
return 200
|
||||
|
||||
|
||||
def _clear_output_history() -> None:
|
||||
_cli()._OUTPUT_HISTORY.clear()
|
||||
|
||||
|
||||
def _output_history_recording() -> bool:
|
||||
return _cli()._OUTPUT_HISTORY_ENABLED and not _cli()._OUTPUT_HISTORY_REPLAYING and not _cli()._OUTPUT_HISTORY_SUPPRESSED
|
||||
|
||||
|
||||
def _record_output_history_entry(entry) -> None:
|
||||
from cli import _output_history_recording
|
||||
if _output_history_recording():
|
||||
_cli()._OUTPUT_HISTORY.append(entry)
|
||||
|
||||
|
||||
def _record_output_history(text: str) -> None:
|
||||
from cli import _output_history_recording
|
||||
if _output_history_recording():
|
||||
_cli()._OUTPUT_HISTORY.extend(str(text).replace("\r", "").rstrip("\n").splitlines())
|
||||
|
||||
|
||||
def _pt_print_ansi(text: str) -> None:
|
||||
"""``_pt_print(ANSI(text))``, falling back to ``print`` when stdout is not a real console."""
|
||||
from cli import _PT_ANSI, _pt_print
|
||||
try:
|
||||
_pt_print(_PT_ANSI(text))
|
||||
except Exception:
|
||||
# NoConsoleScreenBufferError (Windows) / OSError when stdout is e.g. a worker log file.
|
||||
with suppress(Exception):
|
||||
print(text)
|
||||
|
||||
|
||||
def _cprint(text: str):
|
||||
"""Print ANSI text through prompt_toolkit's renderer (patch_stdout swallows raw ANSI).
|
||||
|
||||
From a background thread while an Application runs, a direct print races the input
|
||||
redraw and gets buried, so those go through ``run_in_terminal`` via ``call_soon_threadsafe``.
|
||||
"""
|
||||
from cli import _PT_ANSI, _pt_print, _pt_print_ansi, _record_output_history
|
||||
_record_output_history(text)
|
||||
|
||||
try:
|
||||
from prompt_toolkit.application import get_app_or_none, run_in_terminal
|
||||
except Exception:
|
||||
_pt_print(_PT_ANSI(text))
|
||||
return
|
||||
|
||||
try:
|
||||
app = get_app_or_none()
|
||||
except Exception:
|
||||
app = None
|
||||
|
||||
if app is None or not getattr(app, "_is_running", False):
|
||||
_pt_print_ansi(text)
|
||||
return
|
||||
|
||||
import asyncio as _asyncio
|
||||
|
||||
try:
|
||||
loop = app.loop # type: ignore[attr-defined]
|
||||
except Exception:
|
||||
loop = None
|
||||
try:
|
||||
# get_running_loop(): get_event_loop() warns from threads with no current loop.
|
||||
# Use get_running_loop() instead of get_event_loop() to avoid the DeprecationWarning /
|
||||
# RuntimeWarning emitted by Python 3.10+ when get_event_loop() is called from a thread that has no
|
||||
# current event loop set (e.g. the process_loop background thread). Fixes #19285.
|
||||
current_loop = _asyncio.get_running_loop()
|
||||
except Exception:
|
||||
current_loop = None
|
||||
if loop is None or (current_loop is loop and loop.is_running()):
|
||||
_pt_print(_PT_ANSI(text))
|
||||
return
|
||||
|
||||
def _schedule():
|
||||
# run_in_terminal() returns an awaitable (pt >= 3.0) that must be scheduled or the
|
||||
# output is dropped, or None (mocks / older pt) when it already ran synchronously.
|
||||
# Never fall back to a bare print on error: the sync path already printed.
|
||||
with suppress(Exception):
|
||||
import inspect as _inspect
|
||||
coro = run_in_terminal(lambda: _pt_print(_PT_ANSI(text)))
|
||||
if coro is not None and (_inspect.isawaitable(coro) or _inspect.iscoroutine(coro)):
|
||||
_asyncio.ensure_future(coro)
|
||||
|
||||
try:
|
||||
loop.call_soon_threadsafe(_schedule)
|
||||
except Exception:
|
||||
_pt_print_ansi(text)
|
||||
|
||||
|
||||
def _prepend_note_to_message(message, note: str):
|
||||
"""Prepend a one-shot note to a user message (str, or content-part list when an image is attached).
|
||||
|
||||
For lists the note is folded into the first text part or inserted as a leading one.
|
||||
Unknown shapes are returned unchanged.
|
||||
"""
|
||||
note = str(note or "").strip()
|
||||
if not note:
|
||||
return message
|
||||
if isinstance(message, str):
|
||||
return f"{note}\n\n{message}" if message else note
|
||||
if isinstance(message, list):
|
||||
parts = list(message)
|
||||
for i, part in enumerate(parts):
|
||||
if isinstance(part, dict) and part.get("type") == "text":
|
||||
text = part.get("text", "")
|
||||
parts[i] = {**part, "text": f"{note}\n\n{text}" if text else note}
|
||||
return parts
|
||||
return [{"type": "text", "text": note}, *parts]
|
||||
return message
|
||||
|
||||
|
||||
def _pt_app_is_running() -> bool:
|
||||
"""Whether a prompt_toolkit Application currently owns the live terminal."""
|
||||
try:
|
||||
from prompt_toolkit.application import get_app_or_none
|
||||
app = get_app_or_none()
|
||||
except Exception:
|
||||
return False
|
||||
return app is not None and bool(getattr(app, "_is_running", False))
|
||||
|
||||
|
||||
def _cli_visible_print(text: str = "") -> None:
|
||||
"""``print`` unless a prompt_toolkit Application owns the terminal (patch_stdout swallows bare prints)."""
|
||||
from cli import _cprint, _pt_app_is_running
|
||||
if _pt_app_is_running():
|
||||
_cprint(text)
|
||||
else:
|
||||
print(text)
|
||||
|
||||
|
||||
class ChatConsole:
|
||||
"""Rich Console drop-in routing rendered ANSI through ``_cprint`` so colors survive patch_stdout."""
|
||||
|
||||
def __init__(self):
|
||||
from io import StringIO
|
||||
self._buffer = StringIO()
|
||||
self._inner = Console(file=self._buffer, force_terminal=True, color_system="truecolor", highlight=False)
|
||||
|
||||
def print(self, *args, **kwargs):
|
||||
from cli import _OSC_ESCAPE_RE, _cprint
|
||||
self._buffer.seek(0)
|
||||
self._buffer.truncate()
|
||||
self._inner.width = shutil.get_terminal_size((80, 24)).columns
|
||||
self._inner.print(*args, **kwargs)
|
||||
for line in _OSC_ESCAPE_RE.sub("", self._buffer.getvalue()).rstrip("\n").split("\n"):
|
||||
_cprint(line)
|
||||
|
||||
@contextmanager
|
||||
def status(self, *_args, **_kwargs):
|
||||
"""No-op ``console.status`` so slash helpers don't duplicate ``_busy_command()``'s indicator."""
|
||||
yield self
|
||||
|
||||
|
||||
def _build_compact_banner() -> str:
|
||||
"""Build a compact banner that fits the current terminal width."""
|
||||
try:
|
||||
from hermes_cli.skin_engine import get_active_skin
|
||||
_skin = get_active_skin()
|
||||
except Exception:
|
||||
_skin = None
|
||||
|
||||
def _color(key, default):
|
||||
return _skin.get_color(key, default) if _skin else default
|
||||
|
||||
border_color = _color("banner_border", "#FFD700")
|
||||
title_color = _color("banner_title", "#FFBF00")
|
||||
dim_color = _color("banner_dim", "#B8860B")
|
||||
|
||||
if (getattr(_skin, "name", "default") if _skin else "default") == "default":
|
||||
tiny_line = "☤ NOUS HERMES"
|
||||
else:
|
||||
tiny_line = _skin.get_branding("agent_name", "Hermes Agent") if _skin else "Hermes Agent"
|
||||
line1 = f"{tiny_line} - AI Agent Framework"
|
||||
|
||||
if os.environ.get("HERMES_FAST_STARTUP_BANNER") == "1":
|
||||
from hermes_cli import __release_date__ as _release_date
|
||||
from hermes_cli import __version__ as _version
|
||||
|
||||
version_line = f"Hermes Agent v{_version} ({_release_date})"
|
||||
else:
|
||||
version_line = format_banner_version_label()
|
||||
|
||||
w = min(shutil.get_terminal_size().columns - 2, 88)
|
||||
if w < 30:
|
||||
return f"\n[{title_color}]{tiny_line}[/] [dim {dim_color}]- Nous Research[/]\n"
|
||||
|
||||
inner = w - 2 # inside the box border
|
||||
bar = "═" * w
|
||||
content_width = inner - 2
|
||||
|
||||
line1 = line1[:content_width].ljust(content_width)
|
||||
line2 = version_line[:content_width].ljust(content_width)
|
||||
|
||||
return (
|
||||
f"\n[bold {border_color}]╔{bar}╗[/]\n"
|
||||
f"[bold {border_color}]║[/] [{title_color}]{line1}[/] [bold {border_color}]║[/]\n"
|
||||
f"[bold {border_color}]║[/] [dim {dim_color}]{line2}[/] [bold {border_color}]║[/]\n"
|
||||
f"[bold {border_color}]╚{bar}╝[/]\n"
|
||||
)
|
||||
|
||||
|
||||
def _panel_box_width(title: str, content_lines: list[str], min_width: int = 46, max_width: int = 76) -> int:
|
||||
"""Stable TUI panel width wide enough for the title and content (incl. borders)."""
|
||||
term_cols = shutil.get_terminal_size((100, 20)).columns
|
||||
longest = max([len(title)] + [len(line) for line in content_lines] + [min_width - 4])
|
||||
inner = min(max(longest + 4, min_width - 2), max_width - 2, max(24, term_cols - 6))
|
||||
return inner + 2 # leading/trailing space inside the borders
|
||||
|
||||
|
||||
def _wrap_panel_text(text: str, width: int, subsequent_indent: str = "", *, keep_ws: bool = False) -> list[str]:
|
||||
"""Wrap panel text; ``keep_ws`` preserves whitespace (command/detail previews)."""
|
||||
kw = dict(replace_whitespace=False, drop_whitespace=False) if keep_ws else dict(break_long_words=False, break_on_hyphens=False)
|
||||
wrapped = textwrap.wrap(text, width=max(8, width), subsequent_indent=subsequent_indent, **kw)
|
||||
return wrapped or [""]
|
||||
|
||||
|
||||
_wrap_panel_text_keep_ws = functools.partial(_wrap_panel_text, keep_ws=True)
|
||||
|
||||
|
||||
def _append_panel_line(lines, border_style: str, content_style: str, text: str, box_width: int) -> None:
|
||||
lines.extend(((border_style, "│ "), (content_style, text.ljust(max(0, box_width - 2))), (border_style, " │\n")))
|
||||
|
||||
|
||||
def _append_blank_panel_line(lines, border_style: str, box_width: int) -> None:
|
||||
lines.append((border_style, "│" + (" " * box_width) + "│\n"))
|
||||
319
hermes_cli/cli_shutdown.py
Normal file
319
hermes_cli/cli_shutdown.py
Normal file
@@ -0,0 +1,319 @@
|
||||
"""Classic-CLI shutdown and one-shot finalize helpers: process session-id sync, deferred agent startup, exit watchdog, cleanup steps, session-finalize notifications and the terminal input-mode reset.
|
||||
|
||||
Split out of ``cli.py``; ``cli`` re-exports every public name and moved bodies late-bind
|
||||
cli-level names through ``from cli import ...`` at call time so facade monkeypatch seams hold.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
from contextlib import suppress
|
||||
|
||||
# Log-record parity with the origin module.
|
||||
logger = logging.getLogger("cli")
|
||||
|
||||
|
||||
def _cli():
|
||||
"""Late import of the ``cli`` facade: mutable CLI module state (and its test seams) lives there."""
|
||||
import cli
|
||||
|
||||
return cli
|
||||
|
||||
|
||||
def _sync_process_session_id(session_id: str) -> None:
|
||||
"""Keep process-local session-id consumers aligned after CLI switches."""
|
||||
from gateway.session_context import set_current_session_id
|
||||
|
||||
set_current_session_id(session_id)
|
||||
|
||||
|
||||
def _flush_logging_and_stdio() -> None:
|
||||
"""Best-effort ``logging.shutdown()`` + stdout/stderr flush before ``os._exit``."""
|
||||
with suppress(Exception):
|
||||
logging.shutdown()
|
||||
for _stream in (sys.stdout, sys.stderr):
|
||||
with suppress(Exception):
|
||||
_stream.flush()
|
||||
|
||||
|
||||
def _float_env(name: str, default: float) -> float:
|
||||
"""``float(os.getenv(name))``, or ``default`` when unset/unparseable."""
|
||||
try:
|
||||
return float(os.getenv(name, default))
|
||||
except (TypeError, ValueError):
|
||||
return default
|
||||
|
||||
|
||||
def _exit_watchdog_timeout() -> float:
|
||||
"""``HERMES_EXIT_WATCHDOG_S`` as a float (default 30; ``0`` disables)."""
|
||||
from cli import _float_env
|
||||
return _float_env("HERMES_EXIT_WATCHDOG_S", 30.0)
|
||||
|
||||
|
||||
def _arm_exit_watchdog(timeout_s: float | None = None, *, from_signal: bool = False) -> None:
|
||||
"""Daemon timer that ``os._exit(0)``s after ``timeout_s`` once shutdown has begun.
|
||||
|
||||
Backstop for a cleanup step wedged on network I/O and for interpreter teardown
|
||||
blocked joining non-daemon threads (ThreadPoolExecutor's atexit join). The daemon
|
||||
timer survives ``Py_FinalizeEx``'s joins. ``HERMES_EXIT_WATCHDOG_S=0`` disables.
|
||||
|
||||
1. 2. Interpreter teardown blocked joining non-daemon threads — stdlib ``ThreadPoolExecutor`` workers
|
||||
are joined unconditionally by ``concurrent.futures``' atexit hook even after ``shutdown(wait=False)``,
|
||||
so one tool thread wedged on a socket held the process open forever (#27563 class).
|
||||
"""
|
||||
from cli import _exit_watchdog_timeout, _flush_logging_and_stdio
|
||||
if timeout_s is None:
|
||||
timeout_s = _exit_watchdog_timeout()
|
||||
if timeout_s <= 0:
|
||||
return
|
||||
# Never under pytest: a delayed os._exit(0) would silently kill the test worker.
|
||||
if os.environ.get("PYTEST_CURRENT_TEST"):
|
||||
return
|
||||
|
||||
def _watchdog():
|
||||
time.sleep(timeout_s)
|
||||
# The signal-armed watchdog yields to cleanup's own timer once cleanup is running.
|
||||
if from_signal and _cli()._cleanup_in_progress:
|
||||
return
|
||||
|
||||
try:
|
||||
logger.warning(
|
||||
"Exit watchdog fired after %.0fs — forcing process exit "
|
||||
"(a cleanup step or non-daemon thread is wedged).",
|
||||
timeout_s,
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
_flush_logging_and_stdio()
|
||||
os._exit(0)
|
||||
|
||||
with suppress(Exception): # never block shutdown on watchdog setup
|
||||
threading.Thread(target=_watchdog, daemon=True, name="exit-watchdog").start()
|
||||
|
||||
|
||||
def _shutdown_agent_memory_provider(agent) -> None:
|
||||
"""Memory-provider shutdown (on_session_end + shutdown_all) at the real session boundary."""
|
||||
if not (agent and hasattr(agent, 'shutdown_memory_provider')):
|
||||
return
|
||||
# A /new shortly before exit leaves an LLM-bound boundary task queued; shutdown_all()'s
|
||||
# ~5s drain would cancel it, so give it a bounded head start (watchdog is the backstop).
|
||||
_mm = getattr(agent, '_memory_manager', None)
|
||||
if _mm is not None and hasattr(_mm, 'flush_pending'):
|
||||
with suppress(Exception):
|
||||
_mm.flush_pending(timeout=10)
|
||||
# Forward the agent's transcript so on_session_end hooks see the real conversation;
|
||||
# no-arg fallback for stubs / partially-initialised agents.
|
||||
_session_msgs = getattr(agent, '_session_messages', None)
|
||||
_sid = getattr(agent, "session_id", None) or "<unknown>"
|
||||
# ``_session_messages`` is set on ``AIAgent.__init__`` and refreshed every turn via
|
||||
# ``_persist_session``. Fall back to no-arg on test stubs / partially-initialised agents where the
|
||||
# attribute is missing. See #15165.
|
||||
if isinstance(_session_msgs, list):
|
||||
logger.info("CLI cleanup calling memory shutdown for session %s with %d message(s)", _sid, len(_session_msgs))
|
||||
agent.shutdown_memory_provider(_session_msgs)
|
||||
else:
|
||||
logger.info("CLI cleanup calling memory shutdown for session %s without session message list", _sid)
|
||||
agent.shutdown_memory_provider()
|
||||
|
||||
|
||||
def _stop_cli_wake_word() -> None:
|
||||
from tools.wake_word import stop_listening
|
||||
if _cli()._cli_wake_owner is not None:
|
||||
stop_listening(owner=_cli()._cli_wake_owner)
|
||||
|
||||
|
||||
def _interrupt_async_delegations() -> None:
|
||||
from tools.async_delegation import interrupt_all
|
||||
interrupt_all(reason="CLI shutdown")
|
||||
|
||||
|
||||
def _shutdown_mcp_servers() -> None:
|
||||
from tools.mcp_tool_lifecycle import shutdown_mcp_servers
|
||||
shutdown_mcp_servers()
|
||||
|
||||
|
||||
def _shutdown_cached_aux_clients() -> None:
|
||||
# Otherwise AsyncHttpxClientWrapper.__del__ fires on a closed loop ("Press ENTER to continue...").
|
||||
from agent.auxiliary_client import shutdown_cached_clients
|
||||
shutdown_cached_clients()
|
||||
|
||||
|
||||
# Ordered teardown steps (attribute names, resolved at call time so tests can patch them)
|
||||
# and the exception class each swallows.
|
||||
_CLEANUP_STEPS = (
|
||||
("_stop_cli_wake_word", Exception), ("_cleanup_all_terminals", Exception),
|
||||
("_interrupt_async_delegations", Exception), ("_cleanup_all_browsers", Exception),
|
||||
("_shutdown_mcp_servers", BaseException), ("_shutdown_cached_aux_clients", Exception),
|
||||
)
|
||||
|
||||
|
||||
def _should_emit_cleanup_session_finalize(session_id: str | None) -> bool:
|
||||
# A handed-off session is owned by the gateway process — never finalize it here.
|
||||
# The CLI must not finalize it on exit — that sets end_reason on a row the gateway reopened and is
|
||||
# actively writing to, causing the handoff leg to vanish from session history (#88234).
|
||||
if session_id is not None and session_id in _cli()._handed_off_session_ids:
|
||||
return False
|
||||
if not _cli()._single_query_finalize_attempted_session_ids:
|
||||
return True
|
||||
if session_id is None:
|
||||
return False
|
||||
return session_id not in _cli()._single_query_finalize_attempted_session_ids
|
||||
|
||||
|
||||
def _notify_session_finalize(*, session_id: str | None, platform: str = "cli", reason: str = "shutdown") -> None:
|
||||
with suppress(Exception):
|
||||
from hermes_cli.lifecycle import finalize_session
|
||||
finalize_session(session_id=session_id, platform=platform, reason=reason)
|
||||
|
||||
|
||||
def _oneshot_agent_and_session(cli):
|
||||
"""``(agent, session_id)`` for a one-shot run; the agent's id wins over the CLI's."""
|
||||
agent = getattr(cli, "agent", None)
|
||||
return agent, getattr(agent, "session_id", None) or getattr(cli, "session_id", None)
|
||||
|
||||
|
||||
def _invoke_interrupted_session_end(agent, session_id, reason: str, **extra) -> None:
|
||||
"""Best-effort ``on_session_end`` hook for a turn cut short (never raises)."""
|
||||
with suppress(Exception):
|
||||
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
|
||||
_invoke_hook(
|
||||
"on_session_end", session_id=session_id, completed=False, interrupted=True,
|
||||
model=getattr(agent, "model", None), platform=getattr(agent, "platform", None) or "cli",
|
||||
reason=reason, **extra,
|
||||
)
|
||||
|
||||
|
||||
def _emit_interrupted_session_end(cli, *, reason: str = "keyboard_interrupt") -> None:
|
||||
"""Best-effort on_session_end hook for interrupted non-interactive runs."""
|
||||
from cli import _invoke_interrupted_session_end, _oneshot_agent_and_session
|
||||
agent, session_id = _oneshot_agent_and_session(cli)
|
||||
if agent is None:
|
||||
return
|
||||
|
||||
with suppress(Exception):
|
||||
agent.interrupt(reason.replace("_", " "))
|
||||
|
||||
if session_id in _cli()._handed_off_session_ids: # gateway owns the lifecycle now
|
||||
return
|
||||
if session_id:
|
||||
with suppress(Exception):
|
||||
cli.session_id = session_id
|
||||
|
||||
_invoke_interrupted_session_end(
|
||||
agent, session_id, reason,
|
||||
task_id=getattr(agent, "_current_task_id", "") or "",
|
||||
turn_id=getattr(agent, "_current_turn_id", "") or "",
|
||||
api_request_id=getattr(agent, "_current_api_request_id", "") or "",
|
||||
)
|
||||
|
||||
|
||||
def _notify_single_query_session_finalize(cli, *, reason: str = "shutdown") -> None:
|
||||
from cli import _notify_session_finalize, _oneshot_agent_and_session
|
||||
agent, session_id = _oneshot_agent_and_session(cli)
|
||||
if session_id in _cli()._single_query_finalize_attempted_session_ids:
|
||||
return
|
||||
if session_id in _cli()._handed_off_session_ids: # gateway owns the lifecycle now
|
||||
return
|
||||
|
||||
try:
|
||||
_notify_session_finalize(session_id=session_id, platform=getattr(agent, "platform", None) or "cli", reason=reason)
|
||||
finally:
|
||||
_cli()._single_query_finalize_attempted_session_ids.add(session_id)
|
||||
|
||||
|
||||
def _flush_one_shot_session_store(cli) -> None:
|
||||
"""Durably flush + finalize the one-shot session row before exit (idempotent, best-effort).
|
||||
|
||||
One-shot runs get a single turn, so nothing retries a transiently-failed transcript
|
||||
flush, closes the session row, or drains token deltas the kanban ``os._exit(0)``
|
||||
path skips. Handed-off sessions are left alone.
|
||||
|
||||
- a turn whose in-loop ``_flush_messages_to_session_db`` failed under write-lock contention (e.g. a busy
|
||||
multiplex gateway sharing state.db) was silently lost — the reply reached stdout and agent.log but the
|
||||
resumed session's stored history never changed (#88583); - the resumed/created titled session row was
|
||||
left dangling open (``ended_at``/``end_reason`` NULL) on every one-shot exit; - queued async
|
||||
token-accounting deltas relied on interpreter-exit hooks, which the kanban SIGTERM path's
|
||||
``os._exit(0)`` skips entirely.
|
||||
Idempotent and best-effort: ``_persist_session`` dedupes via the per-message ``_DB_PERSISTED_MARKER``
|
||||
stamps (already-written turns are not re-written) and ``end_session`` no-ops on an already-ended row.
|
||||
See #88234.
|
||||
"""
|
||||
from cli import _oneshot_agent_and_session
|
||||
agent, session_id = _oneshot_agent_and_session(cli)
|
||||
if agent is None or not session_id or session_id in _cli()._handed_off_session_ids:
|
||||
return
|
||||
if getattr(agent, "_persist_disabled", False):
|
||||
return
|
||||
# Passing cli.conversation_history keeps resumed messages identity-skipped even when
|
||||
# the failed flush never stamped them.
|
||||
try:
|
||||
msgs = getattr(agent, "_session_messages", None)
|
||||
if isinstance(msgs, list) and msgs and hasattr(agent, "_persist_session"):
|
||||
agent._persist_session(msgs, getattr(cli, "conversation_history", None))
|
||||
except Exception:
|
||||
logger.debug("one-shot final session persist retry failed", exc_info=True)
|
||||
db = getattr(agent, "_session_db", None) or getattr(cli, "_session_db", None)
|
||||
if db is None:
|
||||
return
|
||||
try:
|
||||
db.flush_token_counts()
|
||||
except Exception:
|
||||
logger.debug("one-shot token-count drain failed", exc_info=True)
|
||||
try:
|
||||
db.end_session(session_id, "cli_close")
|
||||
except Exception:
|
||||
logger.debug("one-shot end_session failed", exc_info=True)
|
||||
|
||||
|
||||
def _wait_for_oneshot_background_completions(cli) -> None:
|
||||
"""Bounded linger for notify_on_complete background processes (children write to our pipes).
|
||||
|
||||
Waits on the whole registry: a one-shot process hosts one agent, and task_id
|
||||
filtering would skip processes registered before the session id settled.
|
||||
|
||||
Skipped when the quiet -Q notify-resume loop already consumed the run's linger
|
||||
budget: it calls wait_for_pending_completions with a shared deadline, so a
|
||||
re-wait here would double-block on the same stuck notify_on_complete child.
|
||||
|
||||
See #90879.
|
||||
"""
|
||||
from cli import _oneshot_agent_and_session
|
||||
from tools.process_registry import process_registry
|
||||
|
||||
if getattr(cli, "_quiet_notify_linger_done", False):
|
||||
return
|
||||
_agent, task_id = _oneshot_agent_and_session(cli)
|
||||
result = process_registry.wait_for_pending_completions(None)
|
||||
if result.get("waited"):
|
||||
logger.info(
|
||||
"One-shot exit linger for session %s: completed=%s timed_out=%s",
|
||||
task_id or "<unknown>",
|
||||
result.get("completed"),
|
||||
result.get("timed_out"),
|
||||
)
|
||||
|
||||
|
||||
def _finalize_single_query(cli) -> None:
|
||||
"""Close one-shot CLI resources before releasing the active session lease."""
|
||||
from cli import _flush_one_shot_session_store, _notify_single_query_session_finalize, _run_cleanup, _wait_for_oneshot_background_completions
|
||||
try:
|
||||
# Order matters: linger for spawned background work BEFORE any teardown (the
|
||||
# parent owns those children's stdout pipes); then the durable flush, since
|
||||
# memory-provider shutdown inside _run_cleanup can issue aux-LLM calls and
|
||||
# nothing after it may fail in a way that loses the turn.
|
||||
for step, what in (
|
||||
(_wait_for_oneshot_background_completions, "background completion wait"),
|
||||
(_flush_one_shot_session_store, "session store flush"),
|
||||
):
|
||||
try:
|
||||
step(cli)
|
||||
except Exception:
|
||||
logger.debug("one-shot %s failed", what, exc_info=True)
|
||||
_notify_single_query_session_finalize(cli)
|
||||
_run_cleanup(notify_session_finalize=False)
|
||||
finally:
|
||||
cli._release_active_session()
|
||||
505
hermes_cli/cli_single_query.py
Normal file
505
hermes_cli/cli_single_query.py
Normal file
@@ -0,0 +1,505 @@
|
||||
"""Single-query (``-q`` / one-shot) helpers: kanban goal loops, exit-code mapping, quiet single-query runner, image routing, signal handlers and the single-query mode orchestrator.
|
||||
|
||||
Split out of ``cli.py``; ``cli`` re-exports every public name and moved bodies late-bind
|
||||
cli-level names through ``from cli import ...`` at call time so facade monkeypatch seams hold.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
from agent.interrupt_compat import request_hard_interrupt
|
||||
from contextlib import suppress
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
# Log-record parity with the origin module.
|
||||
logger = logging.getLogger("cli")
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from cli import HermesCLI
|
||||
|
||||
|
||||
def _int_or(value, default: int) -> int:
|
||||
"""``int(value)``, or ``default`` when it does not parse."""
|
||||
try:
|
||||
return int(value)
|
||||
except (TypeError, ValueError):
|
||||
return default
|
||||
|
||||
|
||||
def _interrupt_agent_for_signal(agent, signum) -> None:
|
||||
"""Hard-interrupt ``agent`` for a shutdown signal, then sleep ``HERMES_SIGTERM_GRACE`` (1.5 s).
|
||||
|
||||
The grace lets the agent thread kill the tool's setsid subprocess group before the
|
||||
main thread unwinds (else an orphan child). Never raises.
|
||||
"""
|
||||
from cli import _float_env
|
||||
try:
|
||||
if agent is not None:
|
||||
request_hard_interrupt(agent, f"received signal {signum}")
|
||||
_grace = _float_env("HERMES_SIGTERM_GRACE", 1.5)
|
||||
if _grace > 0:
|
||||
time.sleep(_grace)
|
||||
except Exception:
|
||||
pass # never block signal handling
|
||||
|
||||
|
||||
def _run_kanban_goal_loop_q(cli: "HermesCLI", first_response: str, run_turn=None, log=None) -> None:
|
||||
"""Drive a kanban goal_mode worker through ``goals.run_kanban_goal_loop`` after its first turn.
|
||||
|
||||
``run_turn`` defaults to the bare ``-Q`` turn (final answer only). The ``-q`` worker path
|
||||
passes ``cli.chat`` so every follow-up turn keeps the tool activity feed that the Kanban
|
||||
worker log is made of. The caller swallows all errors: a broken loop must never wedge a worker.
|
||||
"""
|
||||
from cli import _int_or, _sync_cli_session_id_from_agent
|
||||
task_id = (os.environ.get("HERMES_KANBAN_TASK") or "").strip()
|
||||
if not task_id:
|
||||
return
|
||||
raw_run_id = (os.environ.get("HERMES_KANBAN_RUN_ID") or "").strip()
|
||||
worker_run_id = _int_or(raw_run_id, None) if raw_run_id else None
|
||||
if raw_run_id and worker_run_id is None:
|
||||
logger.warning("invalid HERMES_KANBAN_RUN_ID=%r", raw_run_id)
|
||||
|
||||
from hermes_cli import kanban_db as _kb
|
||||
from hermes_cli import kanban_db_connect as _kbc
|
||||
from hermes_cli.goals import run_kanban_goal_loop as _run_loop, DEFAULT_MAX_TURNS as _DEF_TURNS
|
||||
|
||||
# Goal text = title + body (the acceptance criteria the judge evaluates against).
|
||||
with _kbc.connect_closing() as conn:
|
||||
task = _kb.get_task(conn, task_id)
|
||||
if task is None:
|
||||
return
|
||||
|
||||
goal_text = "\n\n".join(p for p in (task.title or "", task.body) if p).strip()
|
||||
if not goal_text:
|
||||
return
|
||||
|
||||
def _quiet_turn(prompt: str) -> str:
|
||||
result = cli.agent.run_conversation(user_message=prompt, conversation_history=cli.conversation_history)
|
||||
_sync_cli_session_id_from_agent(cli)
|
||||
resp = result.get("final_response", "") if isinstance(result, dict) else str(result)
|
||||
if resp:
|
||||
print(resp)
|
||||
return resp or ""
|
||||
|
||||
def _task_status() -> "str | None":
|
||||
with _kbc.connect_closing() as c:
|
||||
return _kb.goal_run_status(c, task_id, worker_run_id)
|
||||
|
||||
def _block(reason: str) -> None:
|
||||
with _kbc.connect_closing() as c:
|
||||
_kb.block_task(c, task_id, reason=reason, expected_run_id=worker_run_id)
|
||||
|
||||
_run_loop(
|
||||
task_id=task_id, goal_text=goal_text, run_turn=run_turn or _quiet_turn,
|
||||
task_status_fn=_task_status, block_fn=_block,
|
||||
max_turns=task.goal_max_turns or _DEF_TURNS, first_response=first_response or "",
|
||||
log=log or (lambda m: logger.info("%s", m)),
|
||||
)
|
||||
|
||||
|
||||
def _run_kanban_goal_loop_chat(cli: "HermesCLI", first_response: str) -> None:
|
||||
"""``-q`` worker variant: follow-up turns go through ``cli.chat`` (tool feed stays on stdout,
|
||||
which is the Kanban worker log) and judge verdicts are printed there too, so a goal_mode card's
|
||||
log reads like any other worker's instead of staying blank until the final answer."""
|
||||
from cli import _run_kanban_goal_loop_q
|
||||
|
||||
def _log(msg: str) -> None:
|
||||
logger.info("%s", msg)
|
||||
print(msg, flush=True)
|
||||
|
||||
_run_kanban_goal_loop_q(cli, first_response, run_turn=lambda p: cli.chat(p) or "", log=_log)
|
||||
|
||||
|
||||
def _sync_cli_session_id_from_agent(cli) -> None:
|
||||
"""Keep ``cli.session_id`` in sync when mid-run compression rotated the agent's session."""
|
||||
if getattr(cli.agent, "session_id", None) and cli.agent.session_id != cli.session_id:
|
||||
cli.session_id = cli.agent.session_id
|
||||
|
||||
|
||||
# ``failure_reason`` values that say nothing about the task itself: the provider is walled,
|
||||
# down or unreachable, or the account is out of credit, so a Kanban worker signals "try
|
||||
# later" instead of "I failed" and the dispatcher does not spend the task's retry budget on it.
|
||||
_TRANSIENT_PROVIDER_REASONS = frozenset({
|
||||
"rate_limit", "upstream_rate_limit", "billing", "overloaded", "server_error", "timeout",
|
||||
})
|
||||
|
||||
|
||||
# ``failure_reason`` values a retry can never heal: the credential was rejected, the model does
|
||||
# not exist for this account, or the TLS chain is broken. A Kanban worker exits
|
||||
# ``KANBAN_TERMINAL_PROVIDER_EXIT_CODE`` so the dispatcher parks the card after ONE spawn with
|
||||
# the provider's words as the reason, instead of re-spawning into the same wall until
|
||||
# ``kanban.failure_limit`` is spent. ``billing`` stays transient: credit comes back.
|
||||
# ``upstream_blocked`` (a WAF/CDN refusing the SDK's User-Agent) is terminal too: only a
|
||||
# header change heals it, never a retry.
|
||||
_TERMINAL_PROVIDER_REASONS = frozenset({
|
||||
"auth", "auth_permanent", "model_not_found", "ssl_cert_verification", "upstream_blocked",
|
||||
})
|
||||
|
||||
|
||||
def _single_query_exit_code(result, *, credentials_rate_limited: bool = False) -> int:
|
||||
"""Map a one-shot turn result onto a process exit code, for both `-q` and `-Q`.
|
||||
|
||||
0 only when the turn completed; 130 when it was interrupted; 1 when it failed, stopped
|
||||
partway (`partial`, `completed: False`) or never ran at all (credentials / agent init
|
||||
failed, so ``result`` is not a dict). A Kanban worker (``HERMES_KANBAN_TASK`` set) that
|
||||
failed purely on a provider rate-limit / billing wall exits ``KANBAN_RATE_LIMIT_EXIT_CODE``
|
||||
(EX_TEMPFAIL): the dispatcher books that run ``rate_limited`` and requeues the task
|
||||
WITHOUT counting a failure, so a quota window or a provider outage cannot trip the breaker.
|
||||
The same sentinel applies when credential resolution itself is a quota/rate-limit
|
||||
AuthError (no turn result object is produced). One that failed on a terminal provider
|
||||
error (credential revoked, model gone) exits ``KANBAN_TERMINAL_PROVIDER_EXIT_CODE``
|
||||
(EX_CONFIG): the dispatcher blocks the card at once.
|
||||
"""
|
||||
from cli import _TERMINAL_PROVIDER_REASONS, _TRANSIENT_PROVIDER_REASONS
|
||||
if not isinstance(result, dict):
|
||||
if credentials_rate_limited and os.environ.get("HERMES_KANBAN_TASK"):
|
||||
from hermes_cli.kanban_db import KANBAN_RATE_LIMIT_EXIT_CODE
|
||||
return KANBAN_RATE_LIMIT_EXIT_CODE
|
||||
return 1
|
||||
if result.get("interrupted"):
|
||||
return 130
|
||||
if not (result.get("failed") or result.get("partial") or result.get("completed") is False):
|
||||
return 0
|
||||
if os.environ.get("HERMES_KANBAN_TASK"):
|
||||
reason = result.get("failure_reason")
|
||||
if reason in _TRANSIENT_PROVIDER_REASONS:
|
||||
from hermes_cli.kanban_db import KANBAN_RATE_LIMIT_EXIT_CODE
|
||||
return KANBAN_RATE_LIMIT_EXIT_CODE
|
||||
if reason in _TERMINAL_PROVIDER_REASONS:
|
||||
from hermes_cli.kanban_db import KANBAN_TERMINAL_PROVIDER_EXIT_CODE
|
||||
return KANBAN_TERMINAL_PROVIDER_EXIT_CODE
|
||||
return 1
|
||||
|
||||
|
||||
def _run_quiet_single_query(cli, effective_query, emitter=None):
|
||||
"""Quiet (-Q) one-shot turn: run, print the response (stderr for errors/session_id), then sys.exit with the automation exit code.
|
||||
With a ``StreamJsonEmitter`` the final answer and the exit line become the terminal ``result`` JSONL record instead.
|
||||
HERMES_TURN_AUTHOR (set only by a bot-to-bot dispatcher) is consumed here so tool subprocesses do not inherit it.
|
||||
Nested Bot Mode notifies bind this session's key (not the dispatcher's) and resume in-process
|
||||
before stdout is printed, so a teammate reply is the quiet run's final answer rather than a
|
||||
stranded receipt."""
|
||||
from cli import _emit_interrupted_session_end, _run_kanban_goal_loop_q, _single_query_exit_code, _sync_cli_session_id_from_agent
|
||||
from agent.interrupt_compat import _accepts_keyword
|
||||
from agent.turn_author import take_turn_author_from_env
|
||||
from hermes_cli.quiet_single_query import (
|
||||
adopt_unanswered_turn, bind_quiet_session_key, continue_quiet_notify_completions,
|
||||
exit_single_query, quiet_notify_linger_seconds, take_turn_report_path, write_turn_report,
|
||||
)
|
||||
|
||||
author = take_turn_author_from_env()
|
||||
# A spawner that bounds only the turn (cron Bot Chat lane) learns the outcome from this
|
||||
# report, written before the linger below; popped so tool subprocesses do not inherit it.
|
||||
turn_report_path = take_turn_report_path()
|
||||
# A dispatcher's re-run of a failed bot delivery resumes the DM row its first attempt persisted.
|
||||
adopt_unanswered_turn(cli, effective_query)
|
||||
author_kwargs = {"turn_author": author} if author is not None and _accepts_keyword(cli.agent.run_conversation, "turn_author") else {}
|
||||
with bind_quiet_session_key(getattr(cli, "session_id", "") or "default"):
|
||||
try:
|
||||
result = cli.agent.run_conversation(
|
||||
user_message=effective_query, conversation_history=cli.conversation_history, **author_kwargs,
|
||||
)
|
||||
except KeyboardInterrupt:
|
||||
_emit_interrupted_session_end(cli, reason="keyboard_interrupt")
|
||||
if emitter is not None:
|
||||
exit_single_query(emitter.emit_result({"failed": True, "error": "Interrupted"}, session_id=cli.session_id or "", exit_code=130))
|
||||
print(f"\nsession_id: {cli.session_id}", file=sys.stderr)
|
||||
exit_single_query(130)
|
||||
# The exit line below reports session_id to stderr for automation wrappers;
|
||||
# without this sync it would point at the ended parent after compression.
|
||||
_sync_cli_session_id_from_agent(cli)
|
||||
# The turn is over and persisted: the one-shot exit linger that follows protects nested
|
||||
# notify_on_complete replies and is NOT part of the spawner's delivery (#113608). The
|
||||
# report carries what this run will print, so a spawner booking a child still lingering
|
||||
# at its cap relays the answer instead of a timeout (#114980).
|
||||
def _report_turn(res) -> None:
|
||||
write_turn_report(
|
||||
turn_report_path, exit_code=_single_query_exit_code(res),
|
||||
error=str(res.get("error") or "") if isinstance(res, dict) else "agent turn did not run",
|
||||
reply=res.get("final_response", "") if isinstance(res, dict) else str(res),
|
||||
)
|
||||
|
||||
_report_turn(result)
|
||||
if isinstance(result, dict) and not result.get("failed"):
|
||||
history = result.get("messages") or cli.conversation_history
|
||||
|
||||
def _follow_up(text):
|
||||
nonlocal history
|
||||
follow = cli.agent.run_conversation(
|
||||
user_message=text, conversation_history=history, **author_kwargs,
|
||||
)
|
||||
if isinstance(follow, dict) and follow.get("messages"):
|
||||
history = follow["messages"]
|
||||
# Same sync contract as the main turn: a compression rotation during a
|
||||
# follow-up must not leave a stale id on the exit line / drain key.
|
||||
_sync_cli_session_id_from_agent(cli)
|
||||
return follow
|
||||
|
||||
# One shared linger budget for the whole run: the loop below and the later
|
||||
# _wait_for_oneshot_background_completions pass must not each wait the full
|
||||
# oneshot_completion_wait_seconds on the same stuck notify_on_complete child.
|
||||
# Flagged after the loop (finally-equivalent): the wait is the loop's first
|
||||
# statement, so anything raising past that point has consumed budget the
|
||||
# finalize pass must not re-wait.
|
||||
try:
|
||||
continued = continue_quiet_notify_completions(
|
||||
getattr(cli, "session_id", "") or "",
|
||||
_follow_up,
|
||||
owns_event=getattr(cli, "_owns_process_notification", None),
|
||||
linger_budget=quiet_notify_linger_seconds(),
|
||||
)
|
||||
finally:
|
||||
cli._quiet_notify_linger_done = True
|
||||
if isinstance(continued, dict):
|
||||
result = continued
|
||||
# A teammate's reply displaced the answer this run prints; tell the spawner.
|
||||
_report_turn(result)
|
||||
response = result.get("final_response", "") if isinstance(result, dict) else str(result)
|
||||
# Surface backend errors that produced no visible output (e.g. invalid model slug
|
||||
# -> provider 4xx) on stderr so piped stdout stays clean.
|
||||
if emitter is not None:
|
||||
pass # the result record below carries text/error; nothing else may touch stdout
|
||||
elif (
|
||||
not response and isinstance(result, dict) and result.get("error")
|
||||
and (result.get("failed") or result.get("partial"))
|
||||
):
|
||||
print(f"Error: {result['error']}", file=sys.stderr)
|
||||
elif response:
|
||||
print(response)
|
||||
|
||||
# Kanban goal_mode: keep working in THIS session until a judge agrees the card is
|
||||
# done, the worker terminates it, or the turn budget runs out (sticky block).
|
||||
if os.environ.get("HERMES_KANBAN_GOAL_MODE") == "1":
|
||||
try:
|
||||
_run_kanban_goal_loop_q(cli, response)
|
||||
except Exception as _goal_exc:
|
||||
logger.debug("kanban goal loop failed: %s", _goal_exc)
|
||||
|
||||
if emitter is None:
|
||||
print(f"\nsession_id: {cli.session_id}", file=sys.stderr)
|
||||
|
||||
_exit_code = _single_query_exit_code(result)
|
||||
if emitter is not None:
|
||||
_exit_code = emitter.emit_result(result, session_id=cli.session_id or "", exit_code=_exit_code)
|
||||
exit_single_query(_exit_code)
|
||||
|
||||
|
||||
def _route_single_query_images(cli, query, effective_query, single_query_images, single_query_image_urls):
|
||||
"""Attach one-shot images natively when the model supports vision, else pre-describe them as text."""
|
||||
if not (single_query_images or single_query_image_urls):
|
||||
return effective_query
|
||||
# Same image-routing decision as the interactive path: a vision-capable model
|
||||
# (incl. custom-provider models declaring `model.supports_vision: true`) gets
|
||||
# native image_url parts; otherwise the text pipeline (vision_analyze
|
||||
# pre-description).
|
||||
_img_mode = "text"
|
||||
_build_parts = None
|
||||
try:
|
||||
from agent.image_routing import build_native_content_parts as _build_parts # noqa: F811
|
||||
from agent.image_routing import decide_image_input_mode
|
||||
from hermes_cli.config import load_config
|
||||
|
||||
_img_mode = decide_image_input_mode(
|
||||
(cli.provider or "").strip(), (cli.model or "").strip(), load_config(),
|
||||
requested_provider=(cli.requested_provider or "").strip(),
|
||||
)
|
||||
except Exception:
|
||||
_img_mode = "text"
|
||||
|
||||
def _text_fallback():
|
||||
# ``_preprocess_images_with_vision`` only knows local files; when only URLs
|
||||
# were supplied keep the original query text intact.
|
||||
if single_query_images:
|
||||
return cli._preprocess_images_with_vision(query, single_query_images, announce=False)
|
||||
return effective_query
|
||||
|
||||
if _img_mode != "native" or _build_parts is None:
|
||||
return _text_fallback()
|
||||
try:
|
||||
_parts, _skipped = _build_parts(
|
||||
query if isinstance(query, str) else "",
|
||||
[str(p) for p in single_query_images],
|
||||
image_urls=list(single_query_image_urls) or None,
|
||||
)
|
||||
if any(p.get("type") == "image_url" for p in _parts):
|
||||
return _parts
|
||||
return _text_fallback() # all images unreadable
|
||||
except Exception:
|
||||
return _text_fallback()
|
||||
|
||||
|
||||
def _collect_kanban_task_images(single_query_images):
|
||||
"""Kanban workers: image paths/URLs in the task body join the first turn's attachments."""
|
||||
single_query_image_urls: list[str] = []
|
||||
_kanban_task_id = os.environ.get("HERMES_KANBAN_TASK", "").strip()
|
||||
if not _kanban_task_id:
|
||||
return single_query_image_urls
|
||||
try:
|
||||
from hermes_cli import kanban_db as _kb
|
||||
from hermes_cli import kanban_db_connect as _kbc
|
||||
from agent.image_routing import extract_image_refs as _extract_refs
|
||||
|
||||
with _kbc.connect_closing() as _conn:
|
||||
_task = _kb.get_task(_conn, _kanban_task_id)
|
||||
_body = getattr(_task, "body", "") if _task is not None else ""
|
||||
if _body:
|
||||
_kb_paths, _kb_urls = _extract_refs(_body)
|
||||
# Dedupe against any --image the user already passed.
|
||||
_seen = {str(p) for p in single_query_images}
|
||||
for _p in _kb_paths:
|
||||
if _p not in _seen:
|
||||
_seen.add(_p)
|
||||
single_query_images.append(Path(_p))
|
||||
single_query_image_urls.extend(_kb_urls)
|
||||
except Exception as _exc:
|
||||
# Best-effort enrichment; never block worker startup on it.
|
||||
logger.debug("kanban image-ref extraction failed: %s", _exc)
|
||||
return single_query_image_urls
|
||||
|
||||
|
||||
def _install_single_query_signal_handlers(cli):
|
||||
"""Route SIGINT/SIGTERM/SIGHUP through agent.interrupt() before unwinding; kanban workers hard-exit.
|
||||
|
||||
A plain KeyboardInterrupt only unwinds the main thread, so tool worker threads
|
||||
would orphan the setsid child; the interrupt + grace window lets them kill it.
|
||||
"""
|
||||
from cli import _arm_exit_watchdog_on_shutdown_signal, _flush_logging_and_stdio, _flush_one_shot_session_store, _interrupt_agent_for_signal
|
||||
import signal as _signal
|
||||
|
||||
def _signal_handler_q(signum, frame):
|
||||
logger.debug("Received signal %s in single-query mode", signum)
|
||||
_arm_exit_watchdog_on_shutdown_signal() # covers wedges in the unwind below
|
||||
_interrupt_agent_for_signal(getattr(cli, "agent", None), signum)
|
||||
# Kanban: a non-daemon worker blocked in _wait_for_process survives KeyboardInterrupt
|
||||
# and the dispatcher sees 'running' forever, so os._exit(0) (SIGALRM deadman guards
|
||||
# a blocking flush). That skips atexit + the token-drain hook, hence the explicit flush.
|
||||
# Kanban worker exit path (#28181): SIGTERM hits a dispatcher-spawned worker that's likely in a
|
||||
# non-daemon thread waiting on a child subprocess in _wait_for_process. Raising KeyboardInterrupt
|
||||
# only unwinds the main thread; the worker thread keeps running, the process gets reparented to
|
||||
# init, and the dispatcher's _pid_alive check returns True forever — task stuck in 'running'
|
||||
# indefinitely. Skip the controlled-unwind dance and call os._exit(0) so the kernel reclaims the PID
|
||||
# immediately and detect_crashed_workers can reclaim the stale claim on the next tick. Flush logging
|
||||
# + stdout/stderr first so the final debug trace isn't lost; SIGALRM deadman guards the flush
|
||||
# against any rare blocking-I/O case (the reporter measured flush in <1ms; the alarm is a failsafe,
|
||||
# not the common path).
|
||||
if os.environ.get("HERMES_KANBAN_TASK"):
|
||||
with suppress(Exception):
|
||||
if hasattr(_signal, "SIGALRM"):
|
||||
_signal.signal(_signal.SIGALRM, lambda *_: os._exit(0))
|
||||
_signal.alarm(5)
|
||||
with suppress(Exception):
|
||||
# Durable flush FIRST: memory-provider shutdown inside _run_cleanup can issue aux-LLM calls,
|
||||
# and nothing after it may fail in a way that loses the turn (#88583).
|
||||
# os._exit(0) skips atexit AND SessionDB's token-drain hook, so flush + finalize the session
|
||||
# store here or the worker's turn (and its usage deltas) never become durable (#88583 /
|
||||
# #50881 class). Best-effort under the SIGALRM deadman above.
|
||||
_flush_one_shot_session_store(cli)
|
||||
_flush_logging_and_stdio()
|
||||
os._exit(0)
|
||||
raise KeyboardInterrupt()
|
||||
with suppress(Exception): # restricted environments
|
||||
for _name in ("SIGINT", "SIGTERM", "SIGHUP"):
|
||||
if hasattr(_signal, _name):
|
||||
_signal.signal(getattr(_signal, _name), _signal_handler_q)
|
||||
|
||||
|
||||
def _configure_quiet_agent(agent) -> None:
|
||||
"""Neutralize every stdout-writing callback so -Q stdout carries only the final response."""
|
||||
agent.quiet_mode = True
|
||||
agent.suppress_status_output = True
|
||||
agent.stream_delta_callback = None
|
||||
agent.tool_gen_callback = None
|
||||
agent.reasoning_callback = None
|
||||
# The diff/progress callbacks print directly and are gated by neither quiet_mode nor
|
||||
# tool_progress_mode, so they must go too; "off" also covers the executor's direct prints.
|
||||
agent.tool_progress_callback = None
|
||||
agent.tool_start_callback = None
|
||||
agent.tool_complete_callback = None
|
||||
agent.tool_progress_mode = "off"
|
||||
|
||||
|
||||
def _run_single_query_mode(cli, query, image, quiet, oneshot, stream_json: bool = False):
|
||||
"""``-q``/``--image`` entry: seed an interactive session on a TTY, else run the one-shot turn and exit.
|
||||
``stream_json`` (implies quiet) swaps the plain-text final answer for the JSONL event protocol."""
|
||||
from cli import _SeededQueryMessage, _collect_kanban_task_images, _collect_query_images, _configure_quiet_agent, _finalize_single_query, _route_single_query_images, _run_kanban_goal_loop_chat, _run_quiet_single_query, _should_seed_interactive, _single_query_exit_code
|
||||
if _should_seed_interactive(query, image, quiet, oneshot):
|
||||
seeded_query, seeded_images = _collect_query_images(query, image)
|
||||
logger.info(
|
||||
"Seeding interactive session with -q prompt (%d chars, %d images)",
|
||||
len(seeded_query or ""), len(seeded_images),
|
||||
)
|
||||
cli._seeded_first_message = _SeededQueryMessage(seeded_query, seeded_images)
|
||||
return cli.run()
|
||||
cli._single_query_mode = True # agent waits the full MCP cold-start before its only tool snapshot
|
||||
# No user can answer approval prompts: the approval gate takes the deterministic path.
|
||||
# One-shot mode: no between-turns MCP late-binding refresh, so the agent must wait the full MCP
|
||||
# cold-start bound before its first (and only) tool snapshot. See #51316.
|
||||
# Mark single-query for the approval gate. cli.py sets HERMES_INTERACTIVE earlier for interactive sudo
|
||||
# prompts, but a -q run has NO user waiting to answer approval prompts. The gate reads this marker (via
|
||||
# gateway.session_context.get_session_env, which falls back to os.environ when the session-context layer
|
||||
# isn't engaged) and takes the deterministic approvals.single_query_mode path instead of waiting the
|
||||
# full timeout. See #86878.
|
||||
os.environ["HERMES_SINGLE_QUERY_SESSION"] = "1"
|
||||
from hermes_cli.quiet_single_query import exit_single_query
|
||||
if not cli._claim_active_session("cli", stderr=bool(quiet)):
|
||||
exit_single_query(1)
|
||||
try:
|
||||
query, single_query_images = _collect_query_images(query, image)
|
||||
single_query_image_urls = _collect_kanban_task_images(single_query_images)
|
||||
if quiet:
|
||||
# Quiet mode: suppress banner, spinner, tool previews.
|
||||
cli.tool_progress_mode = "off"
|
||||
emitter = None
|
||||
if stream_json:
|
||||
# Built BEFORE credentials/agent init so a failed start still closes the protocol
|
||||
# (init + result) instead of exiting 1 with an empty stdout.
|
||||
from hermes_cli.stream_json import StreamJsonEmitter
|
||||
emitter = StreamJsonEmitter(model=getattr(cli, "model", "") or "", session_id=cli.session_id or "")
|
||||
if cli._ensure_runtime_credentials():
|
||||
effective_query: Any = _route_single_query_images(
|
||||
cli, query, query, single_query_images, single_query_image_urls
|
||||
)
|
||||
turn_route = cli._resolve_turn_agent_config(effective_query)
|
||||
if turn_route["signature"] != cli._active_agent_route_signature:
|
||||
cli.agent = None
|
||||
if cli._init_agent(
|
||||
model_override=turn_route["model"],
|
||||
runtime_override=turn_route["runtime"],
|
||||
request_overrides=turn_route.get("request_overrides"),
|
||||
):
|
||||
_configure_quiet_agent(cli.agent)
|
||||
if emitter is not None:
|
||||
emitter.attach(cli.agent)
|
||||
_run_quiet_single_query(cli, effective_query, emitter=emitter)
|
||||
|
||||
fail_code = _single_query_exit_code(
|
||||
None, credentials_rate_limited=getattr(cli, "_credentials_rate_limited", False))
|
||||
if emitter is not None:
|
||||
emitter.emit_result({"failed": True, "error": "credentials or agent init failed"},
|
||||
session_id=cli.session_id or "", exit_code=fail_code)
|
||||
exit_single_query(fail_code) # credentials or agent init failed
|
||||
# No welcome banner (~420 ms cold); session id / resume hint come from _print_exit_summary().
|
||||
_query_label = query or ("[image attached]" if single_query_images else "")
|
||||
if _query_label:
|
||||
cli.console.print(f"[bold blue]Query:[/] {_query_label}")
|
||||
cli._show_security_advisories()
|
||||
response = cli.chat(query, images=single_query_images or None)
|
||||
# Kanban goal_mode on the `-q` path: same judge loop as `-Q`, but each follow-up turn
|
||||
# runs through cli.chat so the worker log keeps its live tool feed (the dispatcher
|
||||
# used to force -Q here, which left goal_mode cards with a blank Worker log).
|
||||
if os.environ.get("HERMES_KANBAN_GOAL_MODE") == "1":
|
||||
try:
|
||||
_run_kanban_goal_loop_chat(cli, response or "")
|
||||
except Exception as _goal_exc:
|
||||
logger.debug("kanban goal loop failed: %s", _goal_exc)
|
||||
cli._print_exit_summary(clear_screen=False)
|
||||
# Same exit contract as `-Q`: scripts and the Kanban dispatcher read the outcome from
|
||||
# the exit code. This path used to fall through to an implicit 0 for every outcome.
|
||||
exit_single_query(_single_query_exit_code(cli._last_turn_result))
|
||||
finally:
|
||||
_finalize_single_query(cli)
|
||||
600
hermes_cli/cli_terminal_input.py
Normal file
600
hermes_cli/cli_terminal_input.py
Normal file
@@ -0,0 +1,600 @@
|
||||
"""Classic-CLI terminal input helpers: file-drop/attachment path resolution, bracketed-paste patching, extended Enter-key sequences, CPR leak guards and TUI input-height estimation.
|
||||
|
||||
Split out of ``cli.py``; ``cli`` re-exports every public name and moved bodies late-bind
|
||||
cli-level names through ``from cli import ...`` at call time so facade monkeypatch seams hold.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
import shutil
|
||||
import sys
|
||||
import time
|
||||
from contextlib import suppress
|
||||
from pathlib import Path
|
||||
from typing import Dict, Any, Optional, Mapping
|
||||
from urllib.parse import unquote, urlparse
|
||||
|
||||
# Log-record parity with the origin module.
|
||||
logger = logging.getLogger("cli")
|
||||
|
||||
|
||||
def _cli():
|
||||
"""Late import of the ``cli`` facade: mutable CLI module state (and its test seams) lives there."""
|
||||
import cli
|
||||
|
||||
return cli
|
||||
|
||||
|
||||
_IMAGE_EXTENSIONS = frozenset({
|
||||
'.png', '.jpg', '.jpeg', '.gif', '.webp',
|
||||
'.bmp', '.tiff', '.tif', '.svg', '.ico',
|
||||
})
|
||||
|
||||
|
||||
def _termux_example_image_path(filename: str = "cat.png") -> str:
|
||||
"""Return a realistic example media path for the current Termux setup."""
|
||||
candidates = [
|
||||
os.path.expanduser("~/storage/shared"),
|
||||
"/sdcard",
|
||||
"/storage/emulated/0",
|
||||
"/storage/self/primary",
|
||||
]
|
||||
# Literal "/" so the Android hint is right even on Windows.
|
||||
for root in candidates:
|
||||
if os.path.isdir(root):
|
||||
return f"{root}/Pictures/{filename}"
|
||||
return f"~/storage/shared/Pictures/{filename}"
|
||||
|
||||
|
||||
def _split_path_input(raw: str) -> tuple[str, str]:
|
||||
r"""Split a leading path token (quoted or with ``\ `` escapes) from trailing free-form text."""
|
||||
raw = str(raw or "").strip()
|
||||
if not raw:
|
||||
return "", ""
|
||||
|
||||
if raw[0] in {'"', "'"}:
|
||||
quote = raw[0]
|
||||
pos = 1
|
||||
while pos < len(raw):
|
||||
ch = raw[pos]
|
||||
if ch == '\\' and pos + 1 < len(raw):
|
||||
pos += 2
|
||||
continue
|
||||
if ch == quote:
|
||||
return raw[1:pos], raw[pos + 1 :].strip()
|
||||
pos += 1
|
||||
return raw[1:], ""
|
||||
|
||||
pos = 0
|
||||
while pos < len(raw):
|
||||
ch = raw[pos]
|
||||
if ch == '\\' and pos + 1 < len(raw) and raw[pos + 1] == ' ':
|
||||
pos += 2
|
||||
elif ch == ' ':
|
||||
break
|
||||
else:
|
||||
pos += 1
|
||||
|
||||
return raw[:pos].replace('\\ ', ' '), raw[pos:].strip()
|
||||
|
||||
|
||||
def _resolve_attachment_path(raw_path: str) -> Path | None:
|
||||
"""Resolve a user-supplied attachment path (quotes, ``~``, env vars, ``file://``; relative to TERMINAL_CWD).
|
||||
|
||||
Returns ``None`` unless it resolves to an existing file.
|
||||
"""
|
||||
token = str(raw_path or "").strip()
|
||||
if not token:
|
||||
return None
|
||||
|
||||
if token[0] == token[-1] and token[0] in {'"', "'"}:
|
||||
token = token[1:-1].strip()
|
||||
token = token.replace('\\ ', ' ')
|
||||
if not token:
|
||||
return None
|
||||
|
||||
expanded = token
|
||||
if token.startswith("file://"):
|
||||
try:
|
||||
parsed = urlparse(token)
|
||||
if parsed.scheme == "file":
|
||||
expanded = unquote(parsed.path or "")
|
||||
if parsed.netloc and os.name == "nt":
|
||||
expanded = f"//{parsed.netloc}{expanded}"
|
||||
elif os.name == "nt" and len(expanded) >= 3 and expanded[0] == "/" and expanded[1].isalpha() and expanded[2] == ":":
|
||||
# file:///C:/... parses to path "/C:/..." — drop the leading slash
|
||||
# so it resolves as a drive-letter path.
|
||||
expanded = expanded[1:]
|
||||
except Exception:
|
||||
expanded = token
|
||||
expanded = os.path.expandvars(os.path.expanduser(expanded))
|
||||
if os.name != "nt":
|
||||
normalized = expanded.replace("\\", "/")
|
||||
if len(normalized) >= 3 and normalized[1] == ":" and normalized[2] == "/" and normalized[0].isalpha():
|
||||
expanded = f"/mnt/{normalized[0].lower()}/{normalized[3:]}"
|
||||
path = Path(expanded)
|
||||
if not path.is_absolute():
|
||||
base_dir = Path(os.getenv("TERMINAL_CWD", os.getcwd()))
|
||||
path = base_dir / path
|
||||
|
||||
try:
|
||||
resolved = path.resolve()
|
||||
except Exception:
|
||||
resolved = path
|
||||
|
||||
# ENAMETOOLONG for a pasted `/goal <long prose>` that passed the `/` prefilter
|
||||
# would otherwise reach process_loop and silently lose the input.
|
||||
try:
|
||||
if not resolved.exists() or not resolved.is_file():
|
||||
return None
|
||||
except OSError:
|
||||
return None
|
||||
return resolved
|
||||
|
||||
|
||||
def _file_drop_result(path: Path, remainder: str) -> dict:
|
||||
from cli import _IMAGE_EXTENSIONS
|
||||
return {"path": path, "is_image": path.suffix.lower() in _IMAGE_EXTENSIONS, "remainder": remainder}
|
||||
|
||||
|
||||
def _detect_file_drop(user_input: str) -> "dict | None":
|
||||
"""Detect a dragged/pasted file path at the start of *user_input* -> ``{path, is_image, remainder}`` or None."""
|
||||
from cli import _file_drop_result, _resolve_attachment_path, _split_path_input
|
||||
if not isinstance(user_input, str):
|
||||
return None
|
||||
|
||||
stripped = user_input.strip()
|
||||
if not stripped:
|
||||
return None
|
||||
|
||||
# Optionally quoted; then /, ~, ./, ../, a Windows drive prefix, or (unquoted) file://.
|
||||
quoted = stripped[:1] in {"'", '"'}
|
||||
unquoted = stripped[1:] if quoted else stripped
|
||||
starts_like_path = (
|
||||
unquoted.startswith(("/", "~", "./", "../"))
|
||||
or (not quoted and unquoted.startswith("file://"))
|
||||
or (len(unquoted) >= 3 and unquoted[1] == ":" and unquoted[2] in {"\\", "/"} and unquoted[0].isalpha())
|
||||
)
|
||||
if not starts_like_path:
|
||||
return None
|
||||
|
||||
direct_path = _resolve_attachment_path(stripped)
|
||||
if direct_path is not None:
|
||||
return _file_drop_result(direct_path, "")
|
||||
|
||||
first_token, remainder = _split_path_input(stripped)
|
||||
drop_path = _resolve_attachment_path(first_token)
|
||||
if drop_path is None and " " in stripped and not quoted:
|
||||
for pos in reversed([idx for idx, ch in enumerate(stripped) if ch == " "]):
|
||||
drop_path = _resolve_attachment_path(stripped[:pos].rstrip())
|
||||
if drop_path is not None:
|
||||
remainder = stripped[pos + 1 :].strip()
|
||||
break
|
||||
if drop_path is None:
|
||||
return None
|
||||
return _file_drop_result(drop_path, remainder)
|
||||
|
||||
|
||||
def _format_image_attachment_badges(attached_images: list[Path], image_counter: int, width: int | None = None) -> str:
|
||||
"""Attached-image badge row: compact summary on narrow terminals, per-image badges otherwise."""
|
||||
if not attached_images:
|
||||
return ""
|
||||
|
||||
width = width or shutil.get_terminal_size((80, 24)).columns
|
||||
|
||||
def _trunc(name: str, limit: int) -> str:
|
||||
return name if len(name) <= limit else name[: max(1, limit - 3)] + "..."
|
||||
|
||||
if width < 52:
|
||||
if len(attached_images) == 1:
|
||||
return f"[📎 {_trunc(attached_images[0].name, 20)}]"
|
||||
return f"[📎 {len(attached_images)} images attached]"
|
||||
|
||||
if width < 80:
|
||||
if len(attached_images) == 1:
|
||||
return f"[📎 {_trunc(attached_images[0].name, 32)}]"
|
||||
return f"[📎 {_trunc(attached_images[0].name, 20)}] [+{len(attached_images) - 1}]"
|
||||
|
||||
base = image_counter - len(attached_images) + 1
|
||||
return " ".join(f"[📎 Image #{base + i}]" for i in range(len(attached_images)))
|
||||
|
||||
|
||||
def _should_auto_attach_clipboard_image_on_paste(pasted_text: str) -> bool:
|
||||
"""Auto-attach clipboard images only for image-only paste gestures."""
|
||||
return not pasted_text.strip()
|
||||
|
||||
|
||||
def _hermes_call_output_screen_diff(
|
||||
orig_osd, app, output, screen, current_pos, color_depth, previous_screen, last_style, is_done, full_screen,
|
||||
attrs_for_style_string, style_string_has_style, size, previous_width,
|
||||
):
|
||||
"""prompt_toolkit ``_output_screen_diff`` with resize guards.
|
||||
|
||||
Inflates ``previous_screen.height`` when the new screen is taller so pt skips the
|
||||
cursor move that stamps chrome into scrollback; on a corrupt previous paint buffer
|
||||
(tmux re-attach) retries once as a first paint instead of crashing the loop.
|
||||
|
||||
1. 2. On AttributeError/TypeError from a corrupt previous paint buffer (classic after tmux attach with
|
||||
same width), retry once with ``previous_screen=None`` so pt first-paints cleanly instead of crashing the
|
||||
event loop with ``'cell' object has no attribute 'char'``. See #26137.
|
||||
"""
|
||||
try:
|
||||
if previous_screen is not None and hasattr(previous_screen, "height") and previous_screen.height < screen.height:
|
||||
previous_screen.height = screen.height
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
common = (app, output, screen, current_pos, color_depth)
|
||||
tail = (is_done, full_screen, attrs_for_style_string, style_string_has_style, size)
|
||||
try:
|
||||
return orig_osd(*common, previous_screen, last_style, *tail, previous_width)
|
||||
except (AttributeError, TypeError):
|
||||
# Corrupt previous_screen / row cells after client reattach: previous_screen=None
|
||||
# takes the first-paint erase path, previous_width=0 treats the width as changed.
|
||||
return orig_osd(*common, None, None, *tail, 0)
|
||||
|
||||
|
||||
def _apply_bracketed_paste_timeout_patch() -> None:
|
||||
"""Patch ``Vt100Parser.feed`` to flush a bracketed paste whose ESC[201~ end mark never arrives.
|
||||
|
||||
Without it a dropped end mark (SSH glitch, sleep/wake) freezes input forever. Idempotent.
|
||||
"""
|
||||
try:
|
||||
import prompt_toolkit.input.vt100_parser as _vt100_mod
|
||||
from prompt_toolkit.keys import Keys as _PtKeys
|
||||
from prompt_toolkit.key_binding.key_processor import KeyPress as _PtKeyPress
|
||||
|
||||
if getattr(_vt100_mod, "_hermes_bp_timeout_patched", False):
|
||||
return
|
||||
|
||||
_BP_TIMEOUT_S = 2.0
|
||||
|
||||
def _patched_vt100_feed(self_parser, data: str) -> None:
|
||||
if self_parser._in_bracketed_paste:
|
||||
self_parser._paste_buffer += data
|
||||
end_mark = "\x1b[201~"
|
||||
|
||||
if end_mark in self_parser._paste_buffer:
|
||||
end_index = self_parser._paste_buffer.index(end_mark)
|
||||
paste_content = self_parser._paste_buffer[:end_index]
|
||||
self_parser.feed_key_callback(_PtKeyPress(_PtKeys.BracketedPaste, paste_content))
|
||||
self_parser._in_bracketed_paste = False
|
||||
remaining = self_parser._paste_buffer[end_index + len(end_mark):]
|
||||
self_parser._paste_buffer = ""
|
||||
self_parser._hermes_bp_start = None
|
||||
if remaining:
|
||||
_patched_vt100_feed(self_parser, remaining)
|
||||
else:
|
||||
bp_start = getattr(self_parser, "_hermes_bp_start", None)
|
||||
now = time.monotonic()
|
||||
if bp_start is None:
|
||||
self_parser._hermes_bp_start = now
|
||||
elif now - bp_start > _BP_TIMEOUT_S:
|
||||
paste_content = self_parser._paste_buffer
|
||||
self_parser._in_bracketed_paste = False
|
||||
self_parser._paste_buffer = ""
|
||||
self_parser._hermes_bp_start = None
|
||||
if paste_content:
|
||||
self_parser.feed_key_callback(_PtKeyPress(_PtKeys.BracketedPaste, paste_content))
|
||||
logger.warning(
|
||||
"Bracketed-paste timeout (%.1fs) — flushed %d bytes "
|
||||
"without end mark. Terminal may have dropped ESC[201~ "
|
||||
"(see #16263).",
|
||||
now - bp_start, len(paste_content),
|
||||
)
|
||||
else:
|
||||
# Re-inlined: calling the original would double-buffer after entering paste mode.
|
||||
for i, c in enumerate(data):
|
||||
if self_parser._in_bracketed_paste:
|
||||
_patched_vt100_feed(self_parser, data[i:])
|
||||
break
|
||||
self_parser._input_parser.send(c)
|
||||
|
||||
_vt100_mod.Vt100Parser.feed = _patched_vt100_feed
|
||||
_vt100_mod._hermes_bp_timeout_patched = True
|
||||
logger.debug("Applied Vt100Parser bracketed-paste timeout patch (#16263)")
|
||||
except Exception as exc: # noqa: BLE001 — defensive: never break startup
|
||||
logger.debug("Bracketed-paste timeout patch skipped: %s", exc)
|
||||
|
||||
|
||||
# CPR replies (``ESC[<row>;<col>R``) can race past the input parser under resize storms
|
||||
# and land as literal text; the ``^[[...R`` form appears when a filter stripped the ESC.
|
||||
# Cursor Position Report (CPR / DSR) response, format ``ESC[<row>;<col>R``. prompt_toolkit's _on_resize() +
|
||||
# renderer send ``ESC[6n`` queries to the terminal; under resize storms or tab switches the terminal's reply
|
||||
# can race past the input parser and end up in the input buffer as literal text (see issue #14692). Also
|
||||
# matches the visible-form ``^[[<row>;<col>R`` that appears when the ESC byte was stripped by a prior
|
||||
# filter.
|
||||
_DSR_CPR_ESC_RE = re.compile(r"\x1b\[\d+;\d+R")
|
||||
|
||||
|
||||
_DSR_CPR_VISIBLE_RE = re.compile(r"\^\[\[\d+;\d+R")
|
||||
|
||||
|
||||
_SGR_MOUSE_ESC_RE = re.compile(r"\x1b\[<\d+;\d+;\d+[Mm]")
|
||||
|
||||
|
||||
_SGR_MOUSE_VISIBLE_RE = re.compile(r"\^\[\[<\d+;\d+;\d+[Mm]")
|
||||
|
||||
|
||||
# Bare "<btn;col;rowM" fragments; deliberately broad, they are almost never intentional input.
|
||||
_SGR_MOUSE_BARE_RE = re.compile(r"<\d+;\d+;\d+[Mm]")
|
||||
|
||||
|
||||
_TERMINAL_INPUT_MODE_RESET_SEQ = (
|
||||
"\x1b[?1006l\x1b[?1003l\x1b[?1002l\x1b[?1000l" # mouse: SGR, any-motion, button-motion, click
|
||||
"\x1b[?1004l" # focus events
|
||||
"\x1b[?2004l" # bracketed paste
|
||||
"\x1b[?1049l" # leave alt screen
|
||||
"\x1b[<u" # pop kitty keyboard mode
|
||||
"\x1b[>4m" # reset modifyOtherKeys
|
||||
"\x1b[0m\x1b[?25h" # reset attributes, show cursor
|
||||
)
|
||||
|
||||
|
||||
_KITTY_KEYBOARD_PUSH_SEQ = "\x1b[>1u"
|
||||
|
||||
|
||||
_MODIFY_OTHER_KEYS_SEQ = "\x1b[>4;2m"
|
||||
|
||||
|
||||
_EXTENDED_ENTER_KEYS_SEQ = _KITTY_KEYBOARD_PUSH_SEQ + _MODIFY_OTHER_KEYS_SEQ
|
||||
|
||||
|
||||
_BACKSLASH_LINE_CONTINUATION_RE = re.compile(r"\\[ \t]*$")
|
||||
|
||||
|
||||
def _is_ghostty_terminal(env: Optional[Mapping[str, str]] = None) -> bool:
|
||||
"""Whether the terminal is Ghostty.
|
||||
|
||||
Ghostty gets ONLY modifyOtherKeys: its Kitty disambiguate mode strips Alt from
|
||||
Backspace (upstream bug), breaking backward-kill-word.
|
||||
|
||||
Ghostty implements modifyOtherKeys correctly (it then emits ``\\x1b[27;3;127~``, which the alias table
|
||||
also maps). See #87630.
|
||||
"""
|
||||
env = os.environ if env is None else env
|
||||
return (env.get("TERM_PROGRAM") or "").strip() == "ghostty" or (env.get("TERM") or "").strip().lower() == "xterm-ghostty"
|
||||
|
||||
|
||||
def _terminal_supports_extended_enter_keys(env: Optional[Mapping[str, str]] = None) -> bool:
|
||||
"""Allowlist of terminals where requesting modified-Enter reporting is safe (aligned with the Ink TUI)."""
|
||||
env = os.environ if env is None else env
|
||||
term_program = (env.get("TERM_PROGRAM") or "").strip()
|
||||
term = (env.get("TERM") or "").strip().lower()
|
||||
return bool(
|
||||
env.get("WT_SESSION")
|
||||
or term_program in {"iTerm.app", "WezTerm", "ghostty", "vscode"}
|
||||
or env.get("KITTY_WINDOW_ID") or "kitty" in term
|
||||
or term == "xterm-ghostty"
|
||||
or term.startswith("tmux") or term_program.lower() == "tmux"
|
||||
)
|
||||
|
||||
|
||||
def _enable_extended_enter_keys(output=None, env: Optional[Mapping[str, str]] = None) -> bool:
|
||||
"""Ask allowlisted terminals to report modified keys distinctly.
|
||||
|
||||
Pushes BOTH kitty keyboard protocol and xterm modifyOtherKeys (kitty dropped the
|
||||
latter; tmux/VS Code only accept it). Both re-encode modified keys as sequences
|
||||
stock prompt_toolkit barely maps (Ctrl+C once arrived as ``ESC[99;5u``), so
|
||||
``install_modify_other_keys_aliases()`` must have run first. Ghostty gets only
|
||||
modifyOtherKeys. The exit reset pops both modes.
|
||||
|
||||
Under either protocol the terminal re-encodes modified keys as escape sequences — Kitty disambiguate
|
||||
mode as ``ESC[<codepoint>;<mod>u`` (plus the Esc key as ``ESC[27u``), modifyOtherKeys=2 as
|
||||
``ESC[27;<mod>;<codepoint>~``. Stock prompt_toolkit 3.x maps almost none of these, which is why the CSI
|
||||
>1u push was temporarily removed in 87074 (Ctrl+C arrived as ``ESC[99;5u`` and died, #56684).
|
||||
``install_modify_other_keys_aliases()`` (called at CLI startup from ``hermes_cli.pt_input_extras``) now
|
||||
populates ``ANSI_SEQUENCES`` with the full Ctrl/Alt/Shift/multi-modifier and functional-key tables under
|
||||
BOTH formats, so every existing key binding continues to fire — including Ctrl+C, which is handled by
|
||||
prompt_toolkit's ``c-c`` binding (raw mode clears ISIG, so the kernel INTR path was never in play for
|
||||
the CLI).
|
||||
See #87630.
|
||||
"""
|
||||
from cli import _EXTENDED_ENTER_KEYS_SEQ, _MODIFY_OTHER_KEYS_SEQ, _is_ghostty_terminal, _terminal_supports_extended_enter_keys
|
||||
if not _terminal_supports_extended_enter_keys(env):
|
||||
return False
|
||||
seq = _MODIFY_OTHER_KEYS_SEQ if _is_ghostty_terminal(env) else _EXTENDED_ENTER_KEYS_SEQ
|
||||
try:
|
||||
if output is not None and hasattr(output, "write_raw"):
|
||||
output.write_raw(seq)
|
||||
output.flush()
|
||||
return True
|
||||
if sys.stdout is not None and sys.stdout.isatty():
|
||||
sys.stdout.write(seq)
|
||||
sys.stdout.flush()
|
||||
return True
|
||||
except Exception:
|
||||
pass
|
||||
return False
|
||||
|
||||
|
||||
def _cli_multiline_shortcuts_enabled(config: Optional[Dict[str, Any]] = None) -> bool:
|
||||
"""``display.cli_multiline_shortcuts`` (default on: Ctrl+J = newline; off restores the legacy c-j submit)."""
|
||||
if config is None:
|
||||
config = _cli().CLI_CONFIG
|
||||
display = config.get("display") if isinstance(config, dict) else None
|
||||
value = display.get("cli_multiline_shortcuts", True) if isinstance(display, dict) else True
|
||||
if isinstance(value, bool):
|
||||
return value
|
||||
return not (isinstance(value, str) and value.strip().lower() in {"0", "false", "no", "off", "disabled"})
|
||||
|
||||
|
||||
def _is_backslash_line_continuation(text: str) -> bool:
|
||||
"""True when Enter should turn a trailing backslash into a newline."""
|
||||
from cli import _BACKSLASH_LINE_CONTINUATION_RE
|
||||
return bool(_BACKSLASH_LINE_CONTINUATION_RE.search(text or ""))
|
||||
|
||||
|
||||
def _apply_backslash_line_continuation(text: str) -> str:
|
||||
"""Replace a trailing ``\\`` marker with an actual newline."""
|
||||
from cli import _BACKSLASH_LINE_CONTINUATION_RE
|
||||
return _BACKSLASH_LINE_CONTINUATION_RE.sub("", text or "") + "\n"
|
||||
|
||||
|
||||
def _preserve_ctrl_enter_newline() -> bool:
|
||||
"""Environments delivering Ctrl+Enter as bare LF (Windows Terminal, WSL, SSH, Ghostty): c-j must stay newline.
|
||||
|
||||
See issue #22379.
|
||||
"""
|
||||
env = os.environ
|
||||
if (
|
||||
sys.platform == "win32"
|
||||
or any(env.get(v) for v in ("SSH_CONNECTION", "SSH_CLIENT", "SSH_TTY", "WT_SESSION",
|
||||
"GHOSTTY_RESOURCES_DIR", "GHOSTTY_BIN_DIR"))
|
||||
or env.get("TERM", "").lower() == "xterm-ghostty" or env.get("TERM_PROGRAM", "").lower() == "ghostty"
|
||||
or "microsoft" in env.get("WSL_DISTRO_NAME", "").lower()
|
||||
):
|
||||
return True
|
||||
# WSL env vars can be scrubbed under sudo; also peek /proc.
|
||||
for p in ("/proc/version", "/proc/sys/kernel/osrelease"):
|
||||
try:
|
||||
with open(p, "r", encoding="utf-8", errors="ignore") as f:
|
||||
if "microsoft" in f.read().lower():
|
||||
return True
|
||||
except OSError:
|
||||
continue
|
||||
return False
|
||||
|
||||
|
||||
def _bind_prompt_submit_keys(kb, handler, *, multiline_shortcuts_enabled: Optional[bool] = None) -> None:
|
||||
"""Enter always submits; c-j submits only with multiline shortcuts off AND where Ctrl+Enter isn't c-j.
|
||||
|
||||
Even when the setting is disabled, environments where Ctrl+Enter is known to arrive as c-j (Windows,
|
||||
WSL, SSH, Windows Terminal, Ghostty) keep c-j reserved for newline; otherwise Ctrl+Enter submits instead
|
||||
of composing. See _preserve_ctrl_enter_newline() and issue #22379.
|
||||
"""
|
||||
from cli import _cli_multiline_shortcuts_enabled, _preserve_ctrl_enter_newline
|
||||
if multiline_shortcuts_enabled is None:
|
||||
multiline_shortcuts_enabled = _cli_multiline_shortcuts_enabled()
|
||||
kb.add("enter")(handler)
|
||||
if sys.platform != "win32" and not multiline_shortcuts_enabled and not _preserve_ctrl_enter_newline():
|
||||
kb.add("c-j")(handler)
|
||||
|
||||
|
||||
def _disable_prompt_toolkit_cpr_warning(app) -> None:
|
||||
"""Let prompt_toolkit fall back from CPR without printing into the prompt."""
|
||||
with suppress(Exception):
|
||||
app.renderer.cpr_not_supported_callback = None
|
||||
|
||||
|
||||
def _terminal_may_leak_cpr() -> bool:
|
||||
"""Suppress prompt_toolkit CPR queries (delayed replies leak into input); Windows keeps pt's default.
|
||||
|
||||
Delayed CPR replies (``ESC[<row>;<col>R`` / visible ``^[[<row>;<col>R``) leak into the status line and
|
||||
can freeze input when the reply is slow (#13870 on SSH/slow PTYs). The same race hits local POSIX TTYs
|
||||
under heavy subagent / status-line load — see ``tests/hermes_cli/test_cpr_local_leak.py``.
|
||||
"""
|
||||
return os.environ.get("PROMPT_TOOLKIT_NO_CPR", "") == "1" or sys.platform != "win32"
|
||||
|
||||
|
||||
def _build_cpr_disabled_output(stdout):
|
||||
"""Vt100_Output with ``enable_cpr=False`` (``from_pty()`` doesn't expose it), or None on failure.
|
||||
|
||||
prompt_toolkit's renderer sends ``ESC[6n`` (Device Status Report) to learn the cursor row before
|
||||
painting in non-fullscreen mode; the terminal replies ``ESC[<row>;<col>R``. When that reply is delayed
|
||||
it races into the display as raw ``^[[39;1R`` and can stall the renderer's pending-CPR future (#13870;
|
||||
also local POSIX under heavy subagent load).
|
||||
"""
|
||||
try:
|
||||
import io as _io
|
||||
from prompt_toolkit.output.vt100 import Vt100_Output, _get_size
|
||||
from prompt_toolkit.data_structures import Size
|
||||
|
||||
def _get_term_size():
|
||||
rows = columns = None
|
||||
try:
|
||||
rows, columns = _get_size(stdout.fileno())
|
||||
except (OSError, _io.UnsupportedOperation, AttributeError, ValueError):
|
||||
pass
|
||||
return Size(rows=rows or 24, columns=columns or 80)
|
||||
|
||||
return Vt100_Output(stdout, _get_term_size, enable_cpr=False)
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
def _select_classic_cli_pt_output(stdout):
|
||||
"""CPR-disabled ``Vt100_Output`` when CPR may leak, else None (Application keeps pt's default)."""
|
||||
from cli import _build_cpr_disabled_output, _terminal_may_leak_cpr
|
||||
return _build_cpr_disabled_output(stdout) if _terminal_may_leak_cpr() else None
|
||||
|
||||
|
||||
def _strip_leaked_terminal_responses_with_meta(text: str) -> tuple[str, bool]:
|
||||
"""Strip leaked CPR replies and mouse-report fragments -> ``(cleaned, had_mouse_reports)``."""
|
||||
from cli import _DSR_CPR_ESC_RE, _DSR_CPR_VISIBLE_RE, _SGR_MOUSE_BARE_RE, _SGR_MOUSE_ESC_RE, _SGR_MOUSE_VISIBLE_RE
|
||||
if not text:
|
||||
return text, False
|
||||
|
||||
had_mouse_reports = False
|
||||
for present, cpr_re, mouse_re in (
|
||||
("\x1b[" in text, _DSR_CPR_ESC_RE, _SGR_MOUSE_ESC_RE),
|
||||
("^[" in text, _DSR_CPR_VISIBLE_RE, _SGR_MOUSE_VISIBLE_RE),
|
||||
("<" in text and ";" in text and ("M" in text or "m" in text), None, _SGR_MOUSE_BARE_RE),
|
||||
):
|
||||
if not present:
|
||||
continue
|
||||
if cpr_re is not None:
|
||||
text = cpr_re.sub("", text)
|
||||
text, count = mouse_re.subn("", text)
|
||||
had_mouse_reports = had_mouse_reports or count > 0
|
||||
return text, had_mouse_reports
|
||||
|
||||
|
||||
def _estimate_tui_input_height(
|
||||
lines: list[str] | tuple[str, ...], prompt_text: str, terminal_columns: int, *, max_height: int = 8,
|
||||
) -> int:
|
||||
"""Input rows from live terminal cells; the BeforeInput prompt consumes cells only on line 0.
|
||||
|
||||
Never substitute a fake wide fallback: a mis-sized TextArea leaves stale cells at the bottom.
|
||||
"""
|
||||
from cli import _int_or
|
||||
try:
|
||||
from prompt_toolkit.utils import get_cwidth
|
||||
except Exception:
|
||||
get_cwidth = lambda value: len(value or "") # type: ignore[assignment]
|
||||
|
||||
columns = max(1, _int_or(terminal_columns or 0, 0))
|
||||
prompt_width = max(0, get_cwidth(prompt_text or ""))
|
||||
|
||||
visual_lines = 0
|
||||
for index, line in enumerate(lines or [""]):
|
||||
display_width = get_cwidth(line or "") + (prompt_width if index == 0 else 0)
|
||||
visual_lines += max(1, -(-display_width // columns))
|
||||
|
||||
return min(max(visual_lines, 1), max(1, int(max_height or 1)))
|
||||
|
||||
|
||||
def _status_bar_visible_from_display_config(display_config: object) -> bool:
|
||||
"""Initial status-bar visibility; both YAML ``off`` (False) and strings like ``"hidden"`` mean off."""
|
||||
if not isinstance(display_config, dict):
|
||||
display_config = {}
|
||||
statusbar_config = display_config.get("statusbar", display_config.get("tui_statusbar", "top"))
|
||||
if isinstance(statusbar_config, str):
|
||||
return statusbar_config.strip().lower() not in {"0", "false", "hidden", "no", "off"}
|
||||
return statusbar_config is not False
|
||||
|
||||
|
||||
def _collect_query_images(query: str | None, image_arg: str | None = None) -> tuple[str, list[Path]]:
|
||||
"""Collect local image attachments for single-query CLI flows."""
|
||||
from cli import _IMAGE_EXTENSIONS, _detect_file_drop, _resolve_attachment_path
|
||||
message = query or ""
|
||||
images: list[Path] = []
|
||||
|
||||
if isinstance(message, str):
|
||||
dropped = _detect_file_drop(message)
|
||||
if dropped and dropped.get("is_image"):
|
||||
images.append(dropped["path"])
|
||||
message = dropped["remainder"] or f"[User attached image: {dropped['path'].name}]"
|
||||
|
||||
if image_arg:
|
||||
explicit_path = _resolve_attachment_path(image_arg)
|
||||
if explicit_path is None:
|
||||
raise ValueError(f"Image file not found: {image_arg}")
|
||||
if explicit_path.suffix.lower() not in _IMAGE_EXTENSIONS:
|
||||
raise ValueError(f"Not a supported image file: {explicit_path}")
|
||||
images.append(explicit_path)
|
||||
|
||||
return message, list(dict.fromkeys(images))
|
||||
@@ -15,11 +15,11 @@ from prompt_toolkit.keys import Keys
|
||||
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[2]
|
||||
CLI_PATH = ROOT / "cli.py"
|
||||
CLI_PATH = ROOT / "hermes_cli" / "cli_terminal_input.py"
|
||||
|
||||
|
||||
def _load_production_patch_helper():
|
||||
"""Load cli._apply_bracketed_paste_timeout_patch without importing cli.
|
||||
"""Load _apply_bracketed_paste_timeout_patch (hermes_cli/cli_terminal_input.py) without importing cli.
|
||||
|
||||
Importing cli.py pulls optional runtime deps that aren't required for this
|
||||
parser-level regression. AST-loading the exact helper keeps the test tied
|
||||
@@ -38,7 +38,7 @@ def _load_production_patch_helper():
|
||||
None,
|
||||
)
|
||||
assert helper_node is not None, (
|
||||
"cli.py must define _apply_bracketed_paste_timeout_patch()"
|
||||
"hermes_cli/cli_terminal_input.py must define _apply_bracketed_paste_timeout_patch()"
|
||||
)
|
||||
helper_source = ast.get_source_segment(source, helper_node)
|
||||
namespace = {"time": time, "logger": logging.getLogger("test.cli")}
|
||||
|
||||
Reference in New Issue
Block a user