refactor(tui): extract the prompt turn into prompt_turn.py; compact methods_*/session_*/agent_callbacks; split _start_agent_build into scope/wiring helpers
- server.py 7663 -> 5319: _run_prompt_submit and its goal/loop/voice/scope phases move to tui_gateway/prompt_turn.py (bound via method_ctx.bind_module); _start_agent_build split into _bind/_release_build_profile_scopes, _deferred_build_agent_kwargs, _wire_session_agent, _start_session_services; _load_enabled_toolsets split (_enabled_mcp_server_names, _resolve_explicit_toolsets). - methods_slash: _LIVE_SLASH_OUTPUT dispatch table; methods_tools: _SLASH_BUILTINS, _guarded; tool_progress: _PROGRESS_HANDLERS; methods_config_set: _REASONING_DISPLAY_WORDS; methods_voice: _VOICE_TOGGLE_ACTIONS. - Unified: _denied_source (methods_session) replaces _WORKER_SOURCES in methods_profiles; _compress_live_with_feedback / _compute_host_slash shared by the slash mirror and /compress; _end_voice_chat shared by stop-phrase paths; _watcher_mtime_ns; _reaper_session_is_detached_idle; _notif_* helpers. - Dead: _profile_dir_or_err, _resume_info, _slash_builtin_table (refs.py: 0 hits). - acp_adapter/server.py: docstring compaction only (AST-identical). - Every file semantically reviewed hunk-by-hunk for wire/log/lock/order parity.
This commit is contained in:
+134
-318
@@ -13,25 +13,17 @@ from .method_ctx import HandlerRegistry, bind_module
|
||||
_registry = HandlerRegistry()
|
||||
|
||||
|
||||
# ── Child-session live mirror ────────────────────────────────────────
|
||||
# A delegated child is not a live gateway session — it runs synchronously
|
||||
# inside the parent's turn, and its activity reaches the gateway only as
|
||||
# relayed ``subagent.*`` events on the PARENT sid. When a UI opens the child's
|
||||
# own session (session.resume on ``child_session_id``, e.g. the desktop's
|
||||
# open-in-new-window), that window would otherwise sit silent until the run
|
||||
# persists. Translate the relayed events into the native stream events the
|
||||
# window already renders — emitted on the CHILD sid, routed to its transport
|
||||
# by write_json — so the window shows a real midstream turn.
|
||||
# Child-session live mirror: a delegated child's activity reaches the gateway only
|
||||
# as relayed ``subagent.*`` events on the PARENT sid, so a window opened on the
|
||||
# child's own session would sit silent until the run persists. Translate them into
|
||||
# the native stream events emitted on the CHILD sid (write_json routes by sid).
|
||||
_child_mirrors: dict[str, dict] = {}
|
||||
_child_mirrors_lock = threading.Lock()
|
||||
# Stored child session ids with a delegation run currently in flight (refreshed
|
||||
# on every relayed subagent.* event, popped on subagent.complete). Lets a lazy
|
||||
# watch resume report running=true so the window shows a busy indicator even
|
||||
# while the child is silent inside a long tool call (no events for 25s+).
|
||||
# Child session ids with a run in flight (refreshed per relayed event, popped on
|
||||
# complete) so a lazy watch resume reports running=true during a silent long tool.
|
||||
_active_child_runs: dict[str, float] = {}
|
||||
# Staleness bound for the registry: entries refresh on every relayed event, so
|
||||
# anything this quiet means the completion event was lost (callback raised,
|
||||
# parent crashed) — don't let a leaked entry pin "running" forever.
|
||||
# Anything quiet this long lost its completion event (callback raised, parent
|
||||
# crashed) — don't pin "running".
|
||||
_CHILD_RUN_STALE_S = 3600.0
|
||||
|
||||
|
||||
@@ -44,41 +36,36 @@ def _mirror_subagent_to_child(event_type: str, payload: dict) -> None:
|
||||
child_key = str(payload.get("child_session_id") or "")
|
||||
if not child_key:
|
||||
return
|
||||
# Liveness registry first — it must be accurate even when no window is
|
||||
# open, so a window opened mid-run can immediately know the child is busy.
|
||||
# Liveness registry first: accurate with no window open, so one opened mid-run
|
||||
# immediately knows the child is busy.
|
||||
if event_type == "subagent.complete":
|
||||
_active_child_runs.pop(child_key, None)
|
||||
else:
|
||||
_active_child_runs[child_key] = time.time()
|
||||
# Mirror only into a live watch session (keyed by session_key; its live sid
|
||||
# differs from the stored id) that has NOT been upgraded to a full agent.
|
||||
# No window / closed → nothing to mirror; an upgraded session owns a real
|
||||
# native stream and mirroring on top would interleave two turns on one sid.
|
||||
# Either way drop state so a reopened window starts a fresh synthetic turn.
|
||||
# Mirror only into a live watch session NOT upgraded to a full agent: an
|
||||
# upgraded one owns a real native stream and mirroring would interleave two
|
||||
# turns on one sid. Either way drop state so a reopened window starts fresh.
|
||||
live = _find_live_session_by_key(child_key)
|
||||
if live is None or live[1].get("agent") is not None:
|
||||
with _child_mirrors_lock:
|
||||
_child_mirrors.pop(child_key, None)
|
||||
return
|
||||
csid = live[0]
|
||||
text = str(payload.get("text") or "")
|
||||
with _child_mirrors_lock:
|
||||
st = _child_mirrors.setdefault(child_key, {"seq": 0, "open_tool": None, "started": False})
|
||||
if not st["started"]:
|
||||
st["started"] = True
|
||||
_emit("message.start", csid)
|
||||
if event_type == "subagent.thinking":
|
||||
if text := str(payload.get("text") or ""):
|
||||
if text:
|
||||
_emit("reasoning.delta", csid, {"text": text})
|
||||
elif event_type == "subagent.text":
|
||||
# The child's streamed reply text — the actual "agent talking".
|
||||
# Relayed token-by-token from the child's run_conversation
|
||||
# stream_callback, so the watch window streams the reply live.
|
||||
if text := str(payload.get("text") or ""):
|
||||
if text:
|
||||
_emit("message.delta", csid, {"text": text})
|
||||
elif event_type == "subagent.start":
|
||||
# One-time header line (the child's goal) so a freshly opened window
|
||||
# shows immediate context before the first reply token streams.
|
||||
if text := str(payload.get("text") or ""):
|
||||
# One-time header (the child's goal) so a fresh window has context first.
|
||||
if text:
|
||||
_emit("message.delta", csid, {"text": f"{text}\n"})
|
||||
elif event_type == "subagent.tool":
|
||||
if st["open_tool"]:
|
||||
@@ -102,124 +89,59 @@ def _mirror_subagent_to_child(event_type: str, payload: dict) -> None:
|
||||
|
||||
|
||||
def _agent_cbs(sid: str) -> dict:
|
||||
def _read_block(event: str, timeout: int):
|
||||
# read_terminal / read_preview (desktop GUI): blocking bridge like clarify; the
|
||||
# preview read gets longer since a URL tab extracts text from a live page.
|
||||
return lambda start=None, count=None: _block(
|
||||
event, sid, {k: v for k, v in (("start", start), ("count", count)) if v is not None}, timeout=timeout
|
||||
)
|
||||
|
||||
callbacks = {
|
||||
"tool_start_callback": lambda tc_id, name, args: _on_tool_start(
|
||||
sid, tc_id, name, args
|
||||
),
|
||||
"tool_complete_callback": lambda tc_id, name, args, result: _on_tool_complete(
|
||||
sid, tc_id, name, args, result
|
||||
),
|
||||
"tool_start_callback": lambda tc_id, name, args: _on_tool_start(sid, tc_id, name, args),
|
||||
"tool_complete_callback": lambda tc_id, name, args, result: _on_tool_complete(sid, tc_id, name, args, result),
|
||||
"tool_progress_callback": lambda event_type, name=None, preview=None, args=None, **kwargs: _on_tool_progress(
|
||||
sid, event_type, name, preview, args, **kwargs
|
||||
),
|
||||
"tool_gen_callback": lambda name: _tool_progress_enabled(sid)
|
||||
and _emit("tool.generating", sid, {"name": name}),
|
||||
"tool_gen_callback": lambda name: _tool_progress_enabled(sid) and _emit("tool.generating", sid, {"name": name}),
|
||||
"thinking_callback": lambda text: _emit("thinking.delta", sid, {"text": text}),
|
||||
# Affection reaction (ily / <3 / good bot) → hearts. Core-detected, so
|
||||
# the TUI heart and desktop floating hearts share one signal.
|
||||
# Affection reaction (ily / <3 / good bot) → hearts; core-detected so TUI and desktop share it.
|
||||
"reaction_callback": lambda kind: _emit("reaction", sid, {"kind": kind}),
|
||||
"reasoning_callback": lambda text: _emit(
|
||||
"reasoning.delta",
|
||||
sid,
|
||||
{"text": text, **({"verbose": True} if _session_verbose(sid) else {})},
|
||||
"reasoning.delta", sid, {"text": text, **({"verbose": True} if _session_verbose(sid) else {})}
|
||||
),
|
||||
"status_callback": lambda kind, text=None: _status_update(
|
||||
sid, str(kind), None if text is None else str(text)
|
||||
),
|
||||
# Credits/notice spine (L1): an AgentNotice fired by the agent becomes a
|
||||
# notification.show WS event; a recovery clear becomes notification.clear.
|
||||
# Snake_case payload to match the existing gateway-event convention.
|
||||
"status_callback": lambda kind, text=None: _status_update(sid, str(kind), None if text is None else str(text)),
|
||||
# Credits/notice spine: AgentNotice → notification.show; recovery clear → notification.clear.
|
||||
"notice_callback": lambda n: _emit(
|
||||
"notification.show",
|
||||
sid,
|
||||
{
|
||||
"text": n.text,
|
||||
"level": n.level,
|
||||
"kind": n.kind,
|
||||
"ttl_ms": n.ttl_ms,
|
||||
"key": n.key,
|
||||
"id": n.id,
|
||||
},
|
||||
),
|
||||
"notice_clear_callback": lambda key: _emit(
|
||||
"notification.clear", sid, {"key": key}
|
||||
{"text": n.text, "level": n.level, "kind": n.kind, "ttl_ms": n.ttl_ms, "key": n.key, "id": n.id},
|
||||
),
|
||||
"notice_clear_callback": lambda key: _emit("notification.clear", sid, {"key": key}),
|
||||
"clarify_callback": lambda q, c, multi_select=False, questions=None: (
|
||||
_clarify_block(sid, q, c, multi_select=multi_select, questions=questions)
|
||||
),
|
||||
# read_terminal tool (desktop GUI): same blocking bridge as clarify — the
|
||||
# renderer answers terminal.read.respond with the serialized buffer.
|
||||
"read_terminal_callback": lambda start=None, count=None: _block(
|
||||
"terminal.read.request",
|
||||
sid,
|
||||
{k: v for k, v in (("start", start), ("count", count)) if v is not None},
|
||||
timeout=30,
|
||||
),
|
||||
# read_preview tool (desktop GUI): the renderer serializes the active
|
||||
# preview tab (a Browser webview's readable text, a file's identity)
|
||||
# and answers preview.read.respond. Longer timeout than the terminal
|
||||
# read — a URL tab extracts text from a live page.
|
||||
"read_preview_callback": lambda start=None, count=None: _block(
|
||||
"preview.read.request",
|
||||
sid,
|
||||
{k: v for k, v in (("start", start), ("count", count)) if v is not None},
|
||||
timeout=45,
|
||||
),
|
||||
# drive_preview tool (desktop GUI): the renderer injects the interaction
|
||||
# engine into the preview pane's webview (or drives the pane's history)
|
||||
# and answers preview.act.respond with the outcome plus a refreshed
|
||||
# element inventory. Same budget as the preview read, which it ends
|
||||
# with — a click on a slow page pays for the settle and the re-scan.
|
||||
# annotate_preview rides this same callback: it resolves a target
|
||||
# through the same engine and differs only in the verb it sends, so it
|
||||
# needs a tool of its own but not a channel of its own.
|
||||
"drive_preview_callback": lambda payload: _block(
|
||||
"preview.act.request",
|
||||
sid,
|
||||
dict(payload),
|
||||
timeout=45,
|
||||
),
|
||||
# read_window_below tool (desktop GUI): the renderer asks its main
|
||||
# process (which owns native window enumeration) which OS window sits
|
||||
# directly underneath the Hermes window, and answers
|
||||
# window.read.respond with the serialized metadata.
|
||||
"read_window_below_callback": lambda: _block(
|
||||
"window.read.request",
|
||||
sid,
|
||||
{},
|
||||
timeout=30,
|
||||
),
|
||||
# setup_mcp tool (desktop GUI): the renderer shows an inline consent
|
||||
# card and walks the user through install/enable/OAuth via the REST
|
||||
# endpoints, then answers mcp.setup.respond with the JSON outcome.
|
||||
# Long timeout on purpose — the flow can include typing an API key or
|
||||
# a browser OAuth round-trip. Same lifecycle as clarify: on timeout
|
||||
# the tool returns "unanswered" and a late answer is tolerated.
|
||||
"read_terminal_callback": _read_block("terminal.read.request", 30),
|
||||
"read_preview_callback": _read_block("preview.read.request", 45),
|
||||
# drive_preview / annotate_preview (desktop GUI): renderer drives the preview webview and
|
||||
# answers with outcome + refreshed element inventory; same budget as the preview read it ends with.
|
||||
"drive_preview_callback": lambda payload: _block("preview.act.request", sid, dict(payload), timeout=45),
|
||||
# read_window_below (desktop GUI): main process enumerates native windows.
|
||||
"read_window_below_callback": lambda: _block("window.read.request", sid, {}, timeout=30),
|
||||
# setup_mcp (desktop GUI): consent card + install/enable/OAuth. Long timeout on purpose (typing
|
||||
# an API key, browser OAuth); like clarify, timeout returns "unanswered" and a late answer is tolerated.
|
||||
"setup_mcp_callback": lambda server, action, reason: _block(
|
||||
"mcp.setup.request",
|
||||
sid,
|
||||
{"server": server, "action": action, "reason": reason},
|
||||
timeout=600,
|
||||
"mcp.setup.request", sid, {"server": server, "action": action, "reason": reason}, timeout=600
|
||||
),
|
||||
# tour tool (desktop GUI): the renderer drives driver.js — highlighting
|
||||
# elements in the app's own DOM or injecting the engine into the
|
||||
# preview pane's webview — and answers tour.respond with the outcome
|
||||
# (did the selector match, which step is active).
|
||||
# tour (desktop GUI): renderer drives driver.js and answers tour.respond.
|
||||
"tour_callback": lambda payload: _tour_request(sid, payload),
|
||||
}
|
||||
|
||||
# Interim assistant commentary (text alongside tool calls, or the attempted
|
||||
# final answer before a verify-on-stop nudge). Gated on
|
||||
# display.interim_assistant_messages (default true). Also set per-turn in
|
||||
# _run_prompt_submit as defense-in-depth — the per-turn set overwrites
|
||||
# this, and the finally block clears it so a stale closure can't fire.
|
||||
# Interim assistant commentary (text alongside tool calls). Gated on
|
||||
# display.interim_assistant_messages (default true); _run_prompt_submit overwrites
|
||||
# it per turn and clears it in its finally so a stale closure can't fire.
|
||||
if _load_interim_assistant_messages():
|
||||
callbacks["interim_assistant_callback"] = (
|
||||
lambda text, *, already_streamed=False: _emit(
|
||||
"message.interim",
|
||||
sid,
|
||||
{"text": str(text), "already_streamed": bool(already_streamed)},
|
||||
)
|
||||
callbacks["interim_assistant_callback"] = lambda text, *, already_streamed=False: _emit(
|
||||
"message.interim", sid, {"text": str(text), "already_streamed": bool(already_streamed)}
|
||||
)
|
||||
|
||||
return callbacks
|
||||
@@ -227,17 +149,13 @@ def _agent_cbs(sid: str) -> dict:
|
||||
|
||||
def _apply_project_workspace(task_id: str, path: str, _name: str = "") -> None:
|
||||
"""Intentional workspace move from the project_* tools: re-anchor the live
|
||||
session's cwd to the chosen project's folder and push session.info so the
|
||||
desktop follows (refresh tree + scope into the project). This is the ONLY
|
||||
session's cwd and push session.info so the desktop follows. This is the ONLY
|
||||
auto-cwd path — driven by an explicit tool call, never a terminal `cd`."""
|
||||
if not path:
|
||||
return
|
||||
|
||||
# The tool's task_id is the durable session_key, but _sessions is keyed by a
|
||||
# short sid uuid (and the desktop routes events by that sid). Resolve it.
|
||||
# task_id is the durable session_key; _sessions (and desktop event routing) key by sid.
|
||||
key = str(task_id or "")
|
||||
sid = ""
|
||||
session = None
|
||||
sid, session = "", None
|
||||
with _sessions_lock:
|
||||
if key in _sessions:
|
||||
sid, session = key, _sessions[key]
|
||||
@@ -246,34 +164,25 @@ def _apply_project_workspace(task_id: str, path: str, _name: str = "") -> None:
|
||||
if cand.get("session_key") == key or getattr(cand.get("agent"), "session_id", None) == key:
|
||||
sid, session = cand_sid, cand
|
||||
break
|
||||
|
||||
if session is None:
|
||||
return
|
||||
|
||||
resolved = os.path.abspath(os.path.expanduser(str(path)))
|
||||
if not os.path.isdir(resolved):
|
||||
return
|
||||
|
||||
session["cwd"] = resolved
|
||||
session["explicit_cwd"] = True
|
||||
# An explicit project switch supersedes any earlier settle-adopted cwd.
|
||||
session["cwd_from_settle"] = False
|
||||
session["cwd_from_settle"] = False # explicit switch supersedes a settle-adopted cwd
|
||||
_register_session_cwd(session)
|
||||
|
||||
_persist_session_cwd_and_schedule_git_meta(session, resolved)
|
||||
|
||||
try:
|
||||
agent = session.get("agent")
|
||||
info = (
|
||||
_session_info(agent, session)
|
||||
if agent is not None
|
||||
else {
|
||||
"cwd": resolved,
|
||||
"branch": _git_branch_for_cwd(resolved),
|
||||
"project": _project_info_for_cwd(resolved),
|
||||
"lazy": True,
|
||||
if agent is not None:
|
||||
info = _session_info(agent, session)
|
||||
else:
|
||||
info = {
|
||||
"cwd": resolved, "branch": _git_branch_for_cwd(resolved),
|
||||
"project": _project_info_for_cwd(resolved), "lazy": True,
|
||||
}
|
||||
)
|
||||
_emit("session.info", sid, info)
|
||||
except Exception:
|
||||
logger.debug("failed to emit session.info after project workspace move", exc_info=True)
|
||||
@@ -293,20 +202,10 @@ def _wire_callbacks(sid: str):
|
||||
pl["metadata"] = metadata
|
||||
val = _block("secret.request", sid, pl)
|
||||
if not val:
|
||||
return {
|
||||
"success": True,
|
||||
"stored_as": env_var,
|
||||
"validated": False,
|
||||
"skipped": True,
|
||||
"message": "skipped",
|
||||
}
|
||||
return {"success": True, "stored_as": env_var, "validated": False, "skipped": True, "message": "skipped"}
|
||||
from hermes_cli.config import save_env_value_secure
|
||||
|
||||
return {
|
||||
**save_env_value_secure(env_var, val),
|
||||
"skipped": False,
|
||||
"message": "ok",
|
||||
}
|
||||
return {**save_env_value_secure(env_var, val), "skipped": False, "message": "ok"}
|
||||
|
||||
set_secret_capture_callback(secret_cb)
|
||||
|
||||
@@ -322,18 +221,14 @@ def _available_personalities(cfg: dict | None = None) -> dict:
|
||||
"""Built-ins + user overrides, via hermes_cli.personality (single owner)."""
|
||||
from hermes_cli.personality import available_personalities
|
||||
|
||||
if cfg is None:
|
||||
cfg = _load_cfg()
|
||||
return available_personalities(cfg)
|
||||
return available_personalities(_load_cfg() if cfg is None else cfg)
|
||||
|
||||
|
||||
def _validate_personality(value: str, cfg: dict | None = None) -> tuple[str, str]:
|
||||
"""Resolve a requested personality against _available_personalities.
|
||||
"""Resolve a requested personality to (name, prompt) or raise ValueError.
|
||||
|
||||
Same contract as hermes_cli.personality.resolve_personality — (name,
|
||||
prompt) or ValueError — but resolves through the module-level
|
||||
_available_personalities so tests (and future gateway-side overrides)
|
||||
keep a single patch point.
|
||||
Same contract as hermes_cli.personality.resolve_personality, but goes through
|
||||
the module-level _available_personalities so tests keep a single patch point.
|
||||
"""
|
||||
from hermes_cli.personality import normalize_personality_name
|
||||
|
||||
@@ -343,17 +238,12 @@ def _validate_personality(value: str, cfg: dict | None = None) -> tuple[str, str
|
||||
personalities = _available_personalities(cfg)
|
||||
if name not in personalities:
|
||||
names = ", ".join(f"`{n}`" for n in sorted(personalities))
|
||||
raise ValueError(
|
||||
f"Unknown personality: `{str(value).strip()}`.\n\nAvailable: `none`, {names}"
|
||||
)
|
||||
raise ValueError(f"Unknown personality: `{str(value).strip()}`.\n\nAvailable: `none`, {names}")
|
||||
return name, _render_personality_prompt(personalities[name])
|
||||
|
||||
|
||||
def _prompt_text(value) -> str:
|
||||
"""Normalize config prompt values from YAML before handing them to AIAgent.
|
||||
|
||||
Delegates to hermes_cli.personality (single owner).
|
||||
"""
|
||||
"""Normalize config prompt values from YAML for AIAgent (hermes_cli.personality owns this)."""
|
||||
from hermes_cli.personality import prompt_text
|
||||
|
||||
return prompt_text(value)
|
||||
@@ -362,68 +252,50 @@ def _prompt_text(value) -> str:
|
||||
def _apply_personality_to_session(
|
||||
sid: str, session: dict, new_prompt: str, personality: str = ""
|
||||
) -> tuple[bool, dict | None]:
|
||||
"""Apply a personality change to an existing session without resetting history.
|
||||
"""Apply a personality change to a live session without resetting history.
|
||||
|
||||
Updates the agent's ephemeral system prompt in-place so the new personality
|
||||
takes effect on the next turn. The cached base system prompt is left intact
|
||||
(ephemeral_system_prompt is appended at API-call time, not baked into the
|
||||
cache), which preserves prompt-cache hits.
|
||||
|
||||
Also injects a system-role marker into the conversation history so the model
|
||||
knows to pivot its style from this point forward (without this, LLMs tend to
|
||||
continue the tone established by earlier messages in the transcript).
|
||||
|
||||
Returns (history_reset, info) — history_reset is always False since we
|
||||
preserve the conversation.
|
||||
Updates the ephemeral system prompt in place (appended at API-call time, so
|
||||
prompt-cache hits survive) and injects a pivot marker so the model stops
|
||||
pattern-matching its earlier tone. Returns (history_reset=False, info).
|
||||
"""
|
||||
if not session:
|
||||
return False, None
|
||||
session["personality"] = personality
|
||||
|
||||
agent = session.get("agent")
|
||||
if agent:
|
||||
agent.ephemeral_system_prompt = new_prompt or None
|
||||
# Inject a pivot marker into history so the model sees the change point.
|
||||
# This prevents it from pattern-matching its prior style.
|
||||
if new_prompt:
|
||||
marker = (
|
||||
"[System: The user has changed the assistant's personality. "
|
||||
"From this point forward, adopt the following persona and respond "
|
||||
f"accordingly: {new_prompt}]"
|
||||
)
|
||||
else:
|
||||
marker = (
|
||||
"[System: The user has cleared the personality overlay. "
|
||||
"From this point forward, respond in your normal default style.]"
|
||||
)
|
||||
# Tagged like the model-switch marker (`_append_model_switch_marker`):
|
||||
# the marker rides as role=user so strict OpenAI-compatible providers
|
||||
# accept it mid-conversation, but `display_kind` keeps it out of the
|
||||
# `truncate_before_user_ordinal` addressing space. Untagged, it counts
|
||||
# as a real user turn on the gateway side while no client counts it, so
|
||||
# every later rewind resolves one turn too early and `replace_messages`
|
||||
# hard-deletes the difference (#82756).
|
||||
with session["history_lock"]:
|
||||
session["history"].append(
|
||||
{"role": "user", "content": marker, "display_kind": "personality_switch"}
|
||||
)
|
||||
session["history_version"] = int(session.get("history_version", 0)) + 1
|
||||
info = _session_info(agent)
|
||||
_emit("session.info", sid, info)
|
||||
return False, info
|
||||
return False, None
|
||||
if not agent:
|
||||
return False, None
|
||||
agent.ephemeral_system_prompt = new_prompt or None
|
||||
if new_prompt:
|
||||
marker = (
|
||||
"[System: The user has changed the assistant's personality. "
|
||||
"From this point forward, adopt the following persona and respond "
|
||||
f"accordingly: {new_prompt}]"
|
||||
)
|
||||
else:
|
||||
marker = (
|
||||
"[System: The user has cleared the personality overlay. "
|
||||
"From this point forward, respond in your normal default style.]"
|
||||
)
|
||||
# Like the model-switch marker: role=user so strict providers accept it
|
||||
# mid-conversation, but `display_kind` keeps it out of the
|
||||
# `truncate_before_user_ordinal` addressing space (untagged, every rewind would
|
||||
# land one turn early and `replace_messages` hard-delete the difference).
|
||||
with session["history_lock"]:
|
||||
session["history"].append({"role": "user", "content": marker, "display_kind": "personality_switch"})
|
||||
session["history_version"] = int(session.get("history_version", 0)) + 1
|
||||
info = _session_info(agent)
|
||||
_emit("session.info", sid, info)
|
||||
return False, info
|
||||
|
||||
|
||||
def _cfg_max_turns(cfg: dict, default: int) -> int:
|
||||
from hermes_cli.config import resolve_turn_limit as _resolve_turn_limit
|
||||
# Env var override (highest priority)
|
||||
# Env override wins; resolve_turn_limit makes "none"/"unlimited"/0 first-class spellings.
|
||||
env_val = os.environ.get("HERMES_TUI_MAX_TURNS")
|
||||
if env_val:
|
||||
return _resolve_turn_limit(env_val, default=default)
|
||||
# Config file value — route through resolve_turn_limit so that
|
||||
# "none"/"unlimited"/0 are first-class spellings, not int() crashes.
|
||||
agent_cfg = cfg.get("agent") or {}
|
||||
raw = agent_cfg.get("max_turns")
|
||||
raw = (cfg.get("agent") or {}).get("max_turns")
|
||||
if raw is None:
|
||||
raw = cfg.get("max_turns")
|
||||
if raw is not None:
|
||||
@@ -434,24 +306,17 @@ def _cfg_max_turns(cfg: dict, default: int) -> int:
|
||||
def _parse_tui_skills_env() -> list[str]:
|
||||
raw = os.environ.get("HERMES_TUI_SKILLS", "")
|
||||
skills: list[str] = []
|
||||
seen: set[str] = set()
|
||||
for part in raw.replace("\n", ",").split(","):
|
||||
item = part.strip()
|
||||
if item and item not in seen:
|
||||
seen.add(item)
|
||||
if item and item not in skills:
|
||||
skills.append(item)
|
||||
return skills
|
||||
|
||||
|
||||
def _load_fallback_model():
|
||||
"""Return the configured fallback chain for TUI-created agents.
|
||||
|
||||
Delegates to the shared ``get_fallback_chain`` helper so the TUI path
|
||||
stays in parity with ``HermesCLI.__init__`` and ``gateway/run.py``:
|
||||
``fallback_providers`` is the primary source of truth and keeps its
|
||||
order, with legacy ``fallback_model`` entries merged in afterwards
|
||||
(deduped on provider/model/base_url).
|
||||
"""
|
||||
"""Configured fallback chain for TUI-created agents, via the shared
|
||||
``get_fallback_chain`` (parity with HermesCLI/gateway: ``fallback_providers``
|
||||
first in order, legacy ``fallback_model`` merged after, deduped)."""
|
||||
from hermes_cli.fallback_config import get_fallback_chain
|
||||
|
||||
return get_fallback_chain(_load_cfg())
|
||||
@@ -460,9 +325,9 @@ def _load_fallback_model():
|
||||
def _agent_fallback_model(agent):
|
||||
"""Return an agent's fallback chain without rehydrating deliberately empty chains."""
|
||||
if hasattr(agent, "_fallback_chain"):
|
||||
return getattr(agent, "_fallback_chain") or []
|
||||
return agent._fallback_chain or []
|
||||
if hasattr(agent, "_fallback_model"):
|
||||
return getattr(agent, "_fallback_model", None)
|
||||
return agent._fallback_model
|
||||
return _load_fallback_model()
|
||||
|
||||
|
||||
@@ -478,23 +343,17 @@ def _background_agent_kwargs(agent, task_id: str) -> dict:
|
||||
"acp_args": getattr(agent, "acp_args", None) or None,
|
||||
"model": getattr(agent, "model", None) or _resolve_model(),
|
||||
"max_iterations": _cfg_max_turns(cfg, 25),
|
||||
"enabled_toolsets": getattr(agent, "enabled_toolsets", None)
|
||||
# Detached background tasks declare platform="tui" below: they have no
|
||||
# UI session id, so a renderer-routed event has nowhere to land. Resolve
|
||||
# their toolsets against that same platform rather than the gateway
|
||||
# process's, so they never carry GUI schema they cannot use.
|
||||
or _load_enabled_toolsets("tui"),
|
||||
# Detached tasks declare platform="tui" (no UI sid for renderer-routed
|
||||
# events), so resolve toolsets against it — never GUI schema they can't use.
|
||||
"enabled_toolsets": getattr(agent, "enabled_toolsets", None) or _load_enabled_toolsets("tui"),
|
||||
"quiet_mode": True,
|
||||
"verbose_logging": False,
|
||||
"ephemeral_system_prompt": getattr(agent, "ephemeral_system_prompt", None)
|
||||
or None,
|
||||
"ephemeral_system_prompt": getattr(agent, "ephemeral_system_prompt", None) or None,
|
||||
"providers_allowed": getattr(agent, "providers_allowed", None),
|
||||
"providers_ignored": getattr(agent, "providers_ignored", None),
|
||||
"providers_order": getattr(agent, "providers_order", None),
|
||||
"provider_sort": getattr(agent, "provider_sort", None),
|
||||
"provider_require_parameters": getattr(
|
||||
agent, "provider_require_parameters", False
|
||||
),
|
||||
"provider_require_parameters": getattr(agent, "provider_require_parameters", False),
|
||||
"provider_data_collection": getattr(agent, "provider_data_collection", None),
|
||||
"openrouter_min_coding_score": getattr(agent, "openrouter_min_coding_score", None),
|
||||
"session_id": task_id,
|
||||
@@ -510,30 +369,15 @@ def _background_agent_kwargs(agent, task_id: str) -> dict:
|
||||
|
||||
def _ephemeral_preview_agent_kwargs(agent, task_id: str) -> dict:
|
||||
kwargs = _background_agent_kwargs(agent, task_id)
|
||||
kwargs.update(
|
||||
{
|
||||
"enabled_toolsets": ["terminal", "file"],
|
||||
"session_db": None,
|
||||
"skip_memory": True,
|
||||
}
|
||||
)
|
||||
kwargs.update({"enabled_toolsets": ["terminal", "file"], "session_db": None, "skip_memory": True})
|
||||
return kwargs
|
||||
|
||||
|
||||
def _preview_restart_history(session: dict, max_messages: int = 24, max_tool_chars: int = 1200) -> list[dict]:
|
||||
"""Distill the parent session's recent history into a context the
|
||||
ephemeral preview-restart agent can actually use.
|
||||
|
||||
The restart agent has no idea what app the user was building, what
|
||||
server they ran, what cwd was active, or which port belongs to which
|
||||
project. Without this, it would take the bare URL + console logs and
|
||||
guess — usually starting the wrong thing.
|
||||
|
||||
We keep the last ``max_messages`` messages from the parent session so
|
||||
the restart agent sees recent user prompts, assistant replies, and
|
||||
most importantly any terminal/tool calls. Tool result payloads are
|
||||
truncated so we don't blow the context window with file dumps.
|
||||
"""
|
||||
"""Distill recent parent history for the ephemeral preview-restart agent, which
|
||||
otherwise would guess app/server/cwd/port from the bare URL + console logs.
|
||||
Keeps the last ``max_messages`` (always back to the last user turn); tool
|
||||
results are truncated so file dumps don't blow the context window."""
|
||||
try:
|
||||
with session["history_lock"]:
|
||||
history = list(session.get("history", []) or [])
|
||||
@@ -542,20 +386,12 @@ def _preview_restart_history(session: dict, max_messages: int = 24, max_tool_cha
|
||||
|
||||
if not history:
|
||||
return []
|
||||
|
||||
# Anchor on the last user turn so we always include at least the most
|
||||
# recent request and the assistant/tool work that followed it. Then
|
||||
# extend backwards up to max_messages so we capture the prior context.
|
||||
last_user_idx = None
|
||||
start = max(0, len(history) - max_messages)
|
||||
for idx in range(len(history) - 1, -1, -1):
|
||||
if history[idx].get("role") == "user":
|
||||
last_user_idx = idx
|
||||
start = min(start, idx)
|
||||
break
|
||||
|
||||
start = max(0, len(history) - max_messages)
|
||||
if last_user_idx is not None:
|
||||
start = min(start, last_user_idx)
|
||||
|
||||
trimmed: list[dict] = []
|
||||
for msg in history[start:]:
|
||||
if not isinstance(msg, dict):
|
||||
@@ -563,19 +399,12 @@ def _preview_restart_history(session: dict, max_messages: int = 24, max_tool_cha
|
||||
role = msg.get("role")
|
||||
if role not in ("user", "assistant", "tool", "system"):
|
||||
continue
|
||||
|
||||
copy = {k: v for k, v in msg.items() if k != "reasoning"}
|
||||
# Truncate heavy tool outputs so a single 50KB file read doesn't
|
||||
# crowd out the rest of the context.
|
||||
if role == "tool":
|
||||
content = copy.get("content")
|
||||
if isinstance(content, str) and len(content) > max_tool_chars:
|
||||
copy["content"] = (
|
||||
content[:max_tool_chars]
|
||||
+ f"\n... (truncated, original {len(content)} chars)"
|
||||
)
|
||||
copy["content"] = content[:max_tool_chars] + f"\n... (truncated, original {len(content)} chars)"
|
||||
trimmed.append(copy)
|
||||
|
||||
return trimmed
|
||||
|
||||
|
||||
@@ -584,20 +413,16 @@ def _preview_tool_result_preview(name: str, result: str) -> str:
|
||||
data = json.loads(result)
|
||||
except Exception:
|
||||
return ""
|
||||
|
||||
if not isinstance(data, dict):
|
||||
return ""
|
||||
|
||||
if name == "terminal":
|
||||
output = str(data.get("output") or "").strip()
|
||||
exit_code = data.get("exit_code")
|
||||
if output:
|
||||
return output[-1200:]
|
||||
if data.get("session_id"):
|
||||
return f"Background process started: {data.get('session_id')}"
|
||||
if exit_code is not None:
|
||||
return f"terminal exited with code {exit_code}"
|
||||
|
||||
if data.get("exit_code") is not None:
|
||||
return f"terminal exited with code {data.get('exit_code')}"
|
||||
return str(data.get("error") or "").strip()[:1200]
|
||||
|
||||
|
||||
@@ -638,40 +463,31 @@ def _preview_restart_callbacks(parent: str, task_id: str) -> dict:
|
||||
def _reset_session_agent(sid: str, session: dict) -> dict:
|
||||
tokens = _set_session_context(session["session_key"])
|
||||
try:
|
||||
# /new is a full conversation boundary: session-scoped runtime
|
||||
# overrides (/model, /reasoning, /fast) do NOT carry forward — the
|
||||
# fresh agent re-derives model/provider, reasoning, and service tier
|
||||
# from config.yaml (#48055, #23131). Session pins are cleared below so
|
||||
# a rebuild can't resurrect them. (Global process state is still never
|
||||
# touched — see the cross-session-contamination note in
|
||||
# _apply_model_switch.)
|
||||
session.pop("model_override", None)
|
||||
session.pop("create_reasoning_override", None)
|
||||
session.pop("create_service_tier_override", None)
|
||||
session.pop("one_turn_model_restore", None)
|
||||
# /new is a full conversation boundary: session-scoped runtime overrides
|
||||
# (/model, /reasoning, /fast) do NOT carry forward — the fresh agent
|
||||
# re-derives them from config.yaml, and the pins are cleared so a rebuild
|
||||
# can't resurrect them. Global process state is never touched (see the
|
||||
# cross-session-contamination note in _apply_model_switch).
|
||||
for k in ("model_override", "create_reasoning_override", "create_service_tier_override", "one_turn_model_restore"):
|
||||
session.pop(k, None)
|
||||
new_agent = _make_agent(
|
||||
sid,
|
||||
session["session_key"],
|
||||
session_id=session["session_key"],
|
||||
platform_override=_session_source(session),
|
||||
context_cwd_is_launch_artifact=(
|
||||
_context_cwd_is_launch_artifact(session)
|
||||
),
|
||||
context_cwd_is_launch_artifact=_context_cwd_is_launch_artifact(session),
|
||||
)
|
||||
finally:
|
||||
_clear_session_context(tokens)
|
||||
session["agent"] = new_agent
|
||||
session["config_model_seen"] = _config_model_target()
|
||||
session["attached_images"] = []
|
||||
session["queued_prompt"] = None
|
||||
session.update(
|
||||
agent=new_agent, config_model_seen=_config_model_target(), attached_images=[], queued_prompt=None
|
||||
)
|
||||
session.pop("queued_prompts", None)
|
||||
session["_queued_prompt_generation"] = int(session.get("_queued_prompt_generation", 0)) + 1
|
||||
session["edit_snapshots"] = {}
|
||||
session["image_counter"] = 0
|
||||
session["running"] = False
|
||||
session["show_reasoning"] = _load_show_reasoning()
|
||||
session["tool_progress_mode"] = _load_tool_progress_mode()
|
||||
session["tool_started_at"] = {}
|
||||
session.update(
|
||||
edit_snapshots={}, image_counter=0, running=False, show_reasoning=_load_show_reasoning(),
|
||||
tool_progress_mode=_load_tool_progress_mode(), tool_started_at={},
|
||||
)
|
||||
with session["history_lock"]:
|
||||
session["history"] = []
|
||||
session["history_version"] = int(session.get("history_version", 0)) + 1
|
||||
|
||||
Reference in New Issue
Block a user