From c134a16e198bf3157902a311bb6de8d3beeabc82 Mon Sep 17 00:00:00 2001 From: Xi Zhang <106144707+X-iZhang@users.noreply.github.com> Date: Fri, 24 Apr 2026 17:10:25 +0200 Subject: [PATCH] Fix/channel slash rich cli (#184) * feat(cli): implement slash command dispatch for channel messages * fix(cli): make EvoSci serve exit on Ctrl+C and hot-swap /model (#181) * fix(cli): streamline debug logging and formatting in channel command handling * fix(cli): ensure proper handling of asyncio event loop in slash command processing * fix(cli): add error handling for unexpected exceptions in slash command dispatch * fix(cli): improve error messaging for slash command dispatch failures * fix(tests): enhance test setup by restoring channel globals and simplifying assertions * fix(cli): enhance slash command handling across UI surfaces and improve resume command warnings --- EvoScientist/cli/channel.py | 192 +++++++++++ EvoScientist/cli/commands.py | 239 +++++++++++++- EvoScientist/cli/interactive.py | 58 ++++ EvoScientist/cli/tui_interactive.py | 66 ++-- tests/test_cli_channel_slash.py | 473 ++++++++++++++++++++++++++++ tests/test_serve_agent_holder.py | 286 +++++++++++++++++ 6 files changed, 1251 insertions(+), 63 deletions(-) create mode 100644 tests/test_cli_channel_slash.py create mode 100644 tests/test_serve_agent_holder.py diff --git a/EvoScientist/cli/channel.py b/EvoScientist/cli/channel.py index 3deb25f..7c9b73b 100644 --- a/EvoScientist/cli/channel.py +++ b/EvoScientist/cli/channel.py @@ -15,6 +15,7 @@ import queue import threading import time import uuid +from collections.abc import Awaitable, Callable from dataclasses import dataclass from typing import Any @@ -105,6 +106,197 @@ def _pop_channel_response(msg_id: str, *, cancel_pending: bool = False) -> str | return slot["response"] +# --------------------------------------------------------------------------- +# Slash command dispatch for channel messages +# --------------------------------------------------------------------------- +# Shared by all three UI surfaces that accept inbound channel messages: +# Rich CLI (``cli/interactive.py::_process_channel_message``), Textual +# TUI (``cli/tui_interactive.py``'s channel handler), and headless +# serve (``cli/commands.py::_serve_process_message``). They all route +# ``/foo`` text through ``cmd_manager`` instead of feeding it to the +# LLM as a plain prompt. + + +async def dispatch_channel_slash_command( + msg: ChannelMessage, + *, + agent: Any, + thread_id: str, + workspace_dir: str | None, + checkpointer: Any, + append_system: Callable[[str, str], None], + start_new_session_cb: Callable[[], None] | None = None, + handle_session_resume_cb: Callable[..., Awaitable[None]] | None = None, + await_agent_ready: Callable[[], Awaitable[Any]] | None = None, + on_cmd_completed: Callable[..., Awaitable[None]] | None = None, +) -> bool: + """Dispatch a slash command from a channel message. + + Returns True if the helper handled the message (successfully or with + an error) — the caller must then return without streaming anything + to the agent. Returns False for non-slash content or unresolved + slash commands, so the caller can fall through to the agent + streaming path (matches TUI behavior). + + Parameters + ---------- + msg: + The inbound ``ChannelMessage`` to inspect. + agent: + Default agent handle for the ``CommandContext``. Commands that + do not need the agent use this value directly. + thread_id, workspace_dir, checkpointer: + Populate ``CommandContext``. + append_system: + ``(text, style)`` callback for local CLI/TUI log output. Used + by ``ChannelCommandUI`` to surface system breadcrumbs and by + this helper to print the "Executed command from ..." line. + start_new_session_cb, handle_session_resume_cb: + Optional lifecycle callbacks forwarded to ``ChannelCommandUI``. + Headless serve passes ``None`` — ``/new`` and ``/resume`` degrade + gracefully via the default ``ChannelCommandUI`` messages. + await_agent_ready: + Optional async resolver that blocks until the background agent + load finishes. Called only when ``cmd.needs_agent(args)`` is + True. Headless serve passes ``None`` because the agent is + loaded up-front before the bus starts. + on_cmd_completed: + Optional ``async (ctx, original_agent, cmd) -> None`` callback + fired only after ``cmd_manager.execute`` returns True. The + ``original_agent`` argument is the agent handle command execution + started against: ``agent_for_ctx`` after any ``await_agent_ready`` + resolution, or the dispatcher's input agent when no resolver is + supplied. Callers can compare ``ctx.agent`` with + ``original_agent`` to detect command-driven swaps. Used by Rich + CLI to (a) adopt an agent swap (``/model``) back into the + running session and (b) refresh the status snapshot for + commands that mutate session-level state (``/new``, + ``/compact``) — mirrors the REPL dispatch at + ``cli/interactive.py:1002-1030``. Headless serve passes + ``None`` since it cannot hot-swap its polling-loop agent. + """ + if not msg.content.strip().startswith("/"): + return False + + try: + return await _dispatch_channel_slash_impl( + msg, + agent=agent, + thread_id=thread_id, + workspace_dir=workspace_dir, + checkpointer=checkpointer, + append_system=append_system, + start_new_session_cb=start_new_session_cb, + handle_session_resume_cb=handle_session_resume_cb, + await_agent_ready=await_agent_ready, + on_cmd_completed=on_cmd_completed, + ) + except Exception as exc: + # Last-ditch safety: any uncaught exception from inside the + # dispatch pipeline (lazy import failure, ChannelCommandUI + # construction, terminal I/O from ``append_system``, bus + # publish races, ...) must not take down the caller's polling + # loop — a crashed serve / dead channel queue task is worse + # than one failed command. + _channel_logger.exception( + "Unexpected slash dispatch failure for %s (msg=%s)", + msg.channel_type, + msg.msg_id, + ) + try: + _set_channel_response(msg.msg_id, f"Command error: {exc}") + except Exception: # pragma: no cover — defensive + pass + # Return True so the caller treats the message as handled and + # does not fall through to the agent streaming path. + return True + + +async def _dispatch_channel_slash_impl( + msg: ChannelMessage, + *, + agent: Any, + thread_id: str, + workspace_dir: str | None, + checkpointer: Any, + append_system: Callable[[str, str], None], + start_new_session_cb: Callable[[], None] | None, + handle_session_resume_cb: Callable[..., Awaitable[None]] | None, + await_agent_ready: Callable[[], Awaitable[Any]] | None, + on_cmd_completed: Callable[..., Awaitable[None]] | None, +) -> bool: + """Inner body of ``dispatch_channel_slash_command``. + + Split from the public wrapper so the wrapper can guard with a + top-level try/except without visually obscuring the main flow. + """ + # Lazy imports: avoid coupling the channel module to ``commands`` at + # import time (tui_interactive.py does the same). + from ..commands.base import CommandContext + from ..commands.channel_ui import ChannelCommandUI + from ..commands.manager import manager as cmd_manager + + parsed = cmd_manager.resolve(msg.content) + if parsed is None: + # Unknown slash command — let the agent handle it (matches TUI). + return False + cmd, cmd_args = parsed + + agent_for_ctx = agent + if cmd.needs_agent(cmd_args) and await_agent_ready is not None: + try: + agent_for_ctx = await await_agent_ready() + except Exception as exc: + _set_channel_response(msg.msg_id, f"Command error: {exc}") + return True + + ui = ChannelCommandUI( + msg, + append_system_callback=append_system, + start_new_session_callback=start_new_session_cb, + handle_session_resume_callback=handle_session_resume_cb, + ) + ctx = CommandContext( + agent=agent_for_ctx, + thread_id=thread_id, + ui=ui, + workspace_dir=workspace_dir, + checkpointer=checkpointer, + ) + + try: + cmd_executed = await cmd_manager.execute(msg.content, ctx) + except Exception as exc: + _channel_logger.debug(f"Channel command error: {exc}", exc_info=True) + _set_channel_response(msg.msg_id, f"Command error: {exc}") + return True # must return — do NOT fall through to the agent + + if cmd_executed: + if on_cmd_completed is not None: + try: + # Command output already flushed by ``cmd_manager.execute`` + # via ``ctx.ui.flush()`` — the hook does internal state + # sync (agent adoption, status snapshot refresh) only, + # so swallowing its errors keeps the user-visible reply + # intact even if the sync path is broken. + await on_cmd_completed(ctx, agent_for_ctx, cmd) + except Exception as exc: + _channel_logger.debug( + f"Channel command post-exec callback error: {exc}", + exc_info=True, + ) + append_system( + f"[{msg.channel_type}: Executed command from {msg.sender}]", + "dim", + ) + _set_channel_response(msg.msg_id, f"Command executed: {msg.content}") + return True + + # ``cmd_manager.execute`` returned False (empty / unparseable input). + # Fall through to the agent streaming path. + return False + + # --------------------------------------------------------------------------- # HITL approval intercept: bus thread ⇄ main CLI thread # --------------------------------------------------------------------------- diff --git a/EvoScientist/cli/commands.py b/EvoScientist/cli/commands.py index f4857ba..09a7db6 100644 --- a/EvoScientist/cli/commands.py +++ b/EvoScientist/cli/commands.py @@ -4,6 +4,7 @@ import logging import os import queue import re +from collections.abc import Awaitable, Callable from datetime import datetime from importlib.metadata import version as _pkg_version from pathlib import Path @@ -32,6 +33,7 @@ from .channel import ( _start_channels_bus_mode, channel_ask_user_prompt, channel_hitl_prompt, + dispatch_channel_slash_command, ) from .mcp_ui import ( _mcp_add_server_from_kwargs, @@ -496,19 +498,122 @@ async def compact_conversation( _serve_logger = logging.getLogger(__name__) +def _make_serve_start_new_session_cb(agent_holder: dict[str, Any]): + """Build the ``start_new_session_cb`` used by serve mode. + + ``/new`` delegates session rotation entirely to this callback: it + does not mutate ``ctx.thread_id`` itself, it just calls + ``ctx.ui.start_new_session()`` and expects the surface to issue a + fresh thread id. Without a wired callback the channel user gets + ``ChannelCommandUI``'s fallback "restart the channel link" message + and nothing actually rotates. This helper generates a new thread + id, updates the shared holder, and syncs the channel-module global + so subsequent messages land on the new thread. + """ + + def _cb() -> None: + from ..sessions import generate_thread_id + + new_tid = generate_thread_id() + agent_holder["thread_id"] = new_tid + try: + import EvoScientist.cli.channel as _ch_mod + + _ch_mod._cli_thread_id = new_tid + except Exception: # pragma: no cover — defensive + pass + console.print(f"[dim][serve] New thread: {new_tid}[/dim]") + + return _cb + + +def _make_serve_cmd_completed_hook(agent_holder: dict[str, Any]): + """Build the ``on_cmd_completed`` hook used by serve mode. + + Adopts ``/model`` agent swaps and ``/resume`` thread swaps back + into ``agent_holder`` so the outer poll loop picks up the new + handles on subsequent messages. Also keeps + ``EvoScientist.cli.channel`` globals in sync so other readers + (e.g. the bus) see the new values. + + For ``/resume`` specifically, surface a user-visible warning via + ``ctx.ui``: serve uses ``InMemorySaver`` (not the SQLite + checkpointer the interactive CLI uses), so historical state for + any persisted thread is not available — the resumed thread will + start fresh. Without this the ``/resume`` command appears to + succeed silently from the channel user's POV. + + Extracted from ``_serve_process_message`` so it can be unit tested + without spinning up the whole serve loop. + """ + + async def _hook(ctx: Any, original_agent: Any, cmd: Any) -> None: + import EvoScientist.cli.channel as _ch_mod + + if ctx.agent is not None and ctx.agent is not original_agent: + agent_holder["agent"] = ctx.agent + try: + _ch_mod._cli_agent = ctx.agent + except Exception: # pragma: no cover — defensive + pass + + # ``/resume`` mutates ``ctx.thread_id`` directly (its UI callback + # is a no-op in serve mode since there's no REPL to reset). Pick + # up the new id here so subsequent messages run on the resumed + # thread instead of the one captured at serve startup. A bare + # ``/resume`` with no argument just prints usage and leaves + # ``ctx.thread_id`` unchanged — ``thread_changed`` gates both + # the adoption and the user-facing warning so neither fires in + # that case. + new_tid = getattr(ctx, "thread_id", None) + thread_changed = bool(new_tid) and new_tid != agent_holder.get("thread_id") + if thread_changed: + agent_holder["thread_id"] = new_tid + try: + _ch_mod._cli_thread_id = new_tid + except Exception: # pragma: no cover — defensive + pass + + # Surface the in-memory-state limitation to the channel user + # for ``/resume`` so the missing history isn't silent. Flush + # is required because ``cmd_manager.execute`` already flushed + # the command's own output before calling this hook. + if getattr(cmd, "name", None) == "/resume" and thread_changed: + try: + ctx.ui.append_system( + "Note: serve mode uses in-memory state — " + f"thread {new_tid[:8]} starts without prior history.", + style="yellow", + ) + await ctx.ui.flush() + except Exception: # pragma: no cover — defensive + pass + + return _hook + + def _serve_process_message( msg: ChannelMessage, *, - agent: Any, - thread_id: str, + agent_holder: dict[str, Any], model: str | None, workspace_dir: str, show_thinking: bool, + on_cmd_completed: Callable[..., Awaitable[None]] | None = None, + start_new_session_cb: Callable[[], None] | None = None, ) -> None: """Process a single channel message in headless serve mode. Headless equivalent of interactive.py's ``_process_channel_message``. No CLI prompt manipulation — just log lines for monitoring. + + ``agent_holder`` is a mutable dict (keys: ``agent``, ``thread_id``) + shared with the outer ``serve()`` loop. ``on_cmd_completed`` (the + agent-swap / thread-swap adoption hook) and ``start_new_session_cb`` + (thread rotation for ``/new``) are constructed once in ``serve()`` + — if omitted, they're rebuilt per message (backward compat for + existing tests). ``/resume`` lands via the ``on_cmd_completed`` + hook because the command mutates ``ctx.thread_id`` directly. """ import asyncio @@ -573,13 +678,71 @@ def _serve_process_message( def _ask_user_prompt(ask_user_data: dict) -> dict: return channel_ask_user_prompt(ask_user_data, msg) + # ---- Slash command dispatch (cmd_manager, not the agent) ---- + # Headless equivalent of the Rich CLI / TUI slash branch so channel + # commands like ``/evoskills`` actually execute in serve mode instead + # of being fed to the LLM as a plain prompt. ``await_agent_ready`` is + # None because the agent is always loaded before the serve loop polls. + # Uses a dedicated event loop (not ``asyncio.run``) so SIGINT handling + # installed by ``serve()`` remains authoritative — ``asyncio.run`` + # swaps ``signal.set_wakeup_fd`` and can leave it dangling on edge + # cases, which breaks Ctrl+C between messages. + # ``set_event_loop`` is needed because some downstream commands + # (e.g. ``/install-mcp``) call ``asyncio.get_event_loop()``, which + # raises ``RuntimeError`` on Python 3.12+ when the thread has no + # current loop set. The prior loop (often ``None``) is restored in + # the ``finally`` below so subsequent messages start from a clean + # slate. Loop creation lives inside the try so an exception between + # creation and ``set_event_loop`` still closes the loop. + _prev_loop: asyncio.AbstractEventLoop | None + try: + _prev_loop = asyncio.get_event_loop_policy().get_event_loop() + except RuntimeError: + _prev_loop = None + _slash_loop: asyncio.AbstractEventLoop | None = None + _slash_handled = False + _slash_error: Exception | None = None + try: + _slash_loop = asyncio.new_event_loop() + asyncio.set_event_loop(_slash_loop) + _slash_handled = _slash_loop.run_until_complete( + dispatch_channel_slash_command( + msg, + agent=agent_holder["agent"], + thread_id=agent_holder["thread_id"], + workspace_dir=workspace_dir, + checkpointer=None, + append_system=lambda t, s="dim": console.print(t, style=s), + start_new_session_cb=start_new_session_cb + or _make_serve_start_new_session_cb(agent_holder), + on_cmd_completed=on_cmd_completed + or _make_serve_cmd_completed_hook(agent_holder), + ) + ) + except Exception as exc: + _slash_error = exc + _serve_logger.exception("Slash dispatch failed for %s", msg.channel_type) + finally: + if _slash_loop is not None: + _slash_loop.close() + asyncio.set_event_loop(_prev_loop) + + if _slash_error is not None: + _set_channel_response(msg.msg_id, f"Command error: {_slash_error}") + console.print(f"[red]Slash command error: {escape(str(_slash_error))}[/red]") + return + + if _slash_handled: + console.print(f"[dim][{msg.channel_type}] Replied to {msg.sender}[/dim]") + return + meta = build_metadata(workspace_dir, model) try: response = run_streaming( ui_backend="cli", - agent=agent, + agent=agent_holder["agent"], message=msg.content, - thread_id=thread_id, + thread_id=agent_holder["thread_id"], show_thinking=show_thinking, interactive=True, metadata=meta, @@ -695,6 +858,18 @@ def serve( tid = generate_thread_id() + # Mutable holder shared with _serve_process_message so ``/model`` + # invoked over a channel can hot-swap the agent for subsequent + # messages. A pass-by-value parameter gets captured once at startup + # and never updated. + agent_holder: dict[str, Any] = {"agent": agent, "thread_id": tid} + + # Build the slash-dispatch callbacks once; the poll loop reuses + # them for every inbound message. Without this hoist each message + # would allocate a fresh closure pair. + _serve_on_cmd_completed = _make_serve_cmd_completed_hook(agent_holder) + _serve_start_new_session_cb = _make_serve_start_new_session_cb(agent_holder) + _start_channels_bus_mode( config, agent, @@ -707,23 +882,57 @@ def serve( console.print(f"[dim]Workspace: {_shorten_path(ws)}[/dim]") console.print("[dim]Press Ctrl+C to stop.[/dim]\n") + # Explicit SIGINT/SIGTERM handlers. Python's default SIGINT raises + # KeyboardInterrupt in the main thread, which ought to unblock + # ``_message_queue.get(timeout=...)`` and land in the ``except`` + # below — but edge cases (e.g. an asyncio ``set_wakeup_fd`` left + # dangling by a nested ``asyncio.run``) can silently swallow the + # signal. Setting a ``threading.Event`` in addition gives us a + # second gate that the poll loop always observes. + import signal + import threading + + shutdown_event = threading.Event() + + def _handle_shutdown(signum: int, _frame: Any) -> None: + shutdown_event.set() + # Fall back to Python's default SIGINT behavior (raises + # KeyboardInterrupt) so blocking I/O inside ``run_streaming`` + # is still interrupted. For SIGTERM there's no default that + # raises, so the event check below is the only gate. + if signum == signal.SIGINT: + signal.default_int_handler(signum, _frame) + + _orig_sigint = signal.signal(signal.SIGINT, _handle_shutdown) + _orig_sigterm = signal.signal(signal.SIGTERM, _handle_shutdown) + try: - while True: + while not shutdown_event.is_set(): try: - msg = _message_queue.get(timeout=1.0) + msg = _message_queue.get(timeout=0.5) except queue.Empty: continue - _serve_process_message( - msg, - agent=agent, - thread_id=tid, - model=config.model, - workspace_dir=ws, - show_thinking=effective_channel_thinking, - ) + if shutdown_event.is_set(): + break + try: + _serve_process_message( + msg, + agent_holder=agent_holder, + model=config.model, + workspace_dir=ws, + show_thinking=effective_channel_thinking, + on_cmd_completed=_serve_on_cmd_completed, + start_new_session_cb=_serve_start_new_session_cb, + ) + except KeyboardInterrupt: + shutdown_event.set() + break except KeyboardInterrupt: - console.print("\n[dim]Shutting down...[/dim]") + shutdown_event.set() finally: + signal.signal(signal.SIGINT, _orig_sigint) + signal.signal(signal.SIGTERM, _orig_sigterm) + console.print("\n[dim]Shutting down...[/dim]") _channels_stop() console.print("[dim]Stopped.[/dim]") diff --git a/EvoScientist/cli/interactive.py b/EvoScientist/cli/interactive.py index 59bd1b6..9235444 100644 --- a/EvoScientist/cli/interactive.py +++ b/EvoScientist/cli/interactive.py @@ -51,6 +51,7 @@ from .channel import ( _channels_is_running, _message_queue, _set_channel_response, + dispatch_channel_slash_command, ) from .file_mentions import complete_file_mention, resolve_file_mentions from .rich_command_ui import RichCLICommandUI @@ -801,6 +802,63 @@ def cmd_interactive( """Send ask_user questions to channel user and wait for reply.""" return _ch_mod.channel_ask_user_prompt(ask_user_data, msg) + # ---- Slash command dispatch (cmd_manager, not the agent) ---- + # Mirrors the TUI's behavior so ``/evoskills``, ``/mcp list`` + # etc. sent via iMessage actually execute instead of being + # fed to the LLM as a plain prompt. + async def _on_channel_cmd_completed( + ctx: Any, original_agent: Any, cmd: Any + ) -> None: + """Mirror the REPL adoption block at + ``interactive.py:1005-1030`` so ``/model`` and similar + state-mutating commands invoked via a channel actually + rebind the running session and keep the status bar + in sync.""" + nonlocal model + agent_swapped = ( + ctx.agent is not None and ctx.agent is not original_agent + ) + if agent_swapped: + from ..EvoScientist import _ensure_config + + agent_loader.adopt(ctx.agent) + cfg = _ensure_config() + model = cfg.model + state["status_base_snapshot"] = make_empty_status_snapshot( + model + ) + if _channels_is_running(): + _ch_mod._cli_agent = ctx.agent + _ch_mod._cli_thread_id = state["thread_id"] + # ``/new`` rotates ``state["thread_id"]`` / workspace, + # ``/compact`` reduces token usage — both need the + # status snapshot re-rendered even when the agent + # didn't swap. ``/resume`` refreshes inline in its + # own async callback. + if agent_swapped or getattr(cmd, "name", None) in ( + "/compact", + "/new", + ): + await _refresh_status_snapshot(reset_streaming_text=True) + + _slash_handled = await dispatch_channel_slash_command( + msg, + agent=agent_loader.agent, + thread_id=state["thread_id"], + workspace_dir=state["workspace_dir"], + checkpointer=checkpointer, + append_system=lambda t, s="dim": console.print(t, style=s), + start_new_session_cb=_on_start_new_session, + handle_session_resume_cb=_on_handle_session_resume, + await_agent_ready=_await_agent_ready, + on_cmd_completed=_on_channel_cmd_completed, + ) + if _slash_handled: + _print_separator() + sys.stdout.write("\033[34;1m❯\033[0m ") + sys.stdout.flush() + return + try: ready_agent = await _await_agent_ready() meta = build_metadata(state["workspace_dir"], model) diff --git a/EvoScientist/cli/tui_interactive.py b/EvoScientist/cli/tui_interactive.py index 6375d49..0025ad7 100644 --- a/EvoScientist/cli/tui_interactive.py +++ b/EvoScientist/cli/tui_interactive.py @@ -44,6 +44,7 @@ from .channel import ( _channels_stop, _message_queue, _set_channel_response, + dispatch_channel_slash_command, ) from .file_mentions import complete_file_mention, resolve_file_mentions from .history_suggester import HistorySuggester @@ -1819,54 +1820,23 @@ def run_textual_interactive( """ return _ch_mod.channel_ask_user_prompt(ask_user_data, msg) - from ..commands.channel_ui import ChannelCommandUI - - # Handle slash commands from channel - if msg.content.strip().startswith("/"): - # Only wait for the agent if the command actually - # needs it — otherwise ``/mcp add`` & friends would - # hang behind a failing MCP load they're meant to fix. - cmd, cmd_args = cmd_manager.resolve(msg.content) or (None, []) - agent = None - if cmd is not None and cmd.needs_agent(cmd_args): - try: - agent = await self._await_agent_ready() - except Exception as exc: - _set_channel_response(msg.msg_id, f"Error: {exc}") - return - ctx = CommandContext( - agent=agent, - thread_id=self._conversation_tid, - ui=ChannelCommandUI( - msg, - append_system_callback=self._append_system, - start_new_session_callback=self.start_new_session, - handle_session_resume_callback=self.handle_session_resume, - ), - workspace_dir=self._workspace_dir, - checkpointer=self._checkpointer, - ) - try: - cmd_executed = await cmd_manager.execute(msg.content, ctx) - except Exception as _cmd_exc: - # Command raised — report the error and do NOT fall through - # to _stream_with_widgets (which would treat the slash - # command text as a plain user message to the agent). - _channel_logger.debug( - f"Channel command error: {_cmd_exc}", exc_info=True - ) - _set_channel_response(msg.msg_id, f"Command error: {_cmd_exc}") - return # outer finally handles _busy / widget cleanup - - if cmd_executed: - self._append_system( - f"[{msg.channel_type}: Executed command from {msg.sender}]", - style="dim", - ) - _set_channel_response( - msg.msg_id, f"Command executed: {msg.content}" - ) - return # outer finally handles _busy / widget cleanup + # Handle slash commands from channel via the shared + # dispatcher (same path Rich CLI and headless serve use). + # Returns True when the command was handled (or errored) + # and we must NOT fall through to agent streaming. + _slash_handled = await dispatch_channel_slash_command( + msg, + agent=None, # resolved via await_agent_ready on demand + thread_id=self._conversation_tid, + workspace_dir=self._workspace_dir, + checkpointer=self._checkpointer, + append_system=self._append_system, + start_new_session_cb=self.start_new_session, + handle_session_resume_cb=self.handle_session_resume, + await_agent_ready=self._await_agent_ready, + ) + if _slash_handled: + return # outer finally handles _busy / widget cleanup # Non-slash message — streams through the agent, so wait # for readiness now. diff --git a/tests/test_cli_channel_slash.py b/tests/test_cli_channel_slash.py new file mode 100644 index 0000000..3121b90 --- /dev/null +++ b/tests/test_cli_channel_slash.py @@ -0,0 +1,473 @@ +"""Tests for ``dispatch_channel_slash_command`` in ``cli/channel.py``. + +Regression coverage for issue #181 — slash commands arriving over a +messaging channel must route through ``cmd_manager`` instead of being +fed to the LLM as a plain prompt, on every UI surface (Rich CLI, TUI, +headless ``serve``). +""" + +from unittest.mock import AsyncMock, MagicMock, patch + +from EvoScientist.cli.channel import ( + ChannelMessage, + dispatch_channel_slash_command, +) +from tests.conftest import run_async as _run + + +def _make_msg( + content: str = "/evoskills core", msg_id: str = "msg-1" +) -> ChannelMessage: + return ChannelMessage( + msg_id=msg_id, + content=content, + sender="+44XXXXXX", + channel_type="imessage", + metadata={}, + channel_ref=None, + bus_ref=None, + chat_id="+44XXXXXX", + message_id="ts-1", + ) + + +def test_non_slash_returns_false(): + """Plain text messages must fall through to the agent.""" + msg = _make_msg(content="hello agent") + append = MagicMock() + handled = _run( + dispatch_channel_slash_command( + msg, + agent=None, + thread_id="t1", + workspace_dir=None, + checkpointer=None, + append_system=append, + ) + ) + assert handled is False + append.assert_not_called() + + +def test_unresolved_slash_returns_false(): + """Unknown slash commands must fall through (matches TUI behavior).""" + msg = _make_msg(content="/unknown-cmd") + append = MagicMock() + with patch( + "EvoScientist.commands.manager.manager.resolve", + return_value=None, + ): + handled = _run( + dispatch_channel_slash_command( + msg, + agent=None, + thread_id="t1", + workspace_dir=None, + checkpointer=None, + append_system=append, + ) + ) + assert handled is False + + +def test_successful_slash_execution_sets_response_and_breadcrumb(): + """Known slash command: cmd_manager.execute ran, helper returns True, + sends a confirmation to the channel user, and appends a local log line.""" + msg = _make_msg() + fake_cmd = MagicMock() + fake_cmd.needs_agent.return_value = False + append = MagicMock() + with ( + patch( + "EvoScientist.commands.manager.manager.resolve", + return_value=(fake_cmd, ["core"]), + ), + patch( + "EvoScientist.commands.manager.manager.execute", + new=AsyncMock(return_value=True), + ) as mock_execute, + patch("EvoScientist.cli.channel._set_channel_response") as mock_set_resp, + ): + handled = _run( + dispatch_channel_slash_command( + msg, + agent="fake-agent", + thread_id="t1", + workspace_dir="/tmp", + checkpointer=None, + append_system=append, + ) + ) + assert handled is True + mock_execute.assert_awaited_once() + mock_set_resp.assert_called_once() + assert mock_set_resp.call_args[0][0] == "msg-1" + assert "Command executed" in mock_set_resp.call_args[0][1] + breadcrumbs = [call.args[0] for call in append.call_args_list] + assert any("Executed command from" in t for t in breadcrumbs) + + +def test_needs_agent_awaits_loader_and_passes_result(): + """Commands with needs_agent=True must await the loader and the + resulting agent must flow through the CommandContext.""" + msg = _make_msg() + fake_cmd = MagicMock() + fake_cmd.needs_agent.return_value = True + append = MagicMock() + await_called = MagicMock() + + async def _await_ready(): + await_called() + return "ready-agent" + + with ( + patch( + "EvoScientist.commands.manager.manager.resolve", + return_value=(fake_cmd, []), + ), + patch( + "EvoScientist.commands.manager.manager.execute", + new=AsyncMock(return_value=True), + ) as mock_execute, + patch("EvoScientist.cli.channel._set_channel_response"), + ): + handled = _run( + dispatch_channel_slash_command( + msg, + agent=None, + thread_id="t1", + workspace_dir=None, + checkpointer=None, + append_system=append, + await_agent_ready=_await_ready, + ) + ) + assert handled is True + await_called.assert_called_once() + ctx_arg = mock_execute.await_args.args[1] + assert ctx_arg.agent == "ready-agent" + + +def test_await_agent_ready_failure_sets_error_response(): + msg = _make_msg() + fake_cmd = MagicMock() + fake_cmd.needs_agent.return_value = True + append = MagicMock() + + async def _await_ready(): + raise RuntimeError("agent blew up") + + with ( + patch( + "EvoScientist.commands.manager.manager.resolve", + return_value=(fake_cmd, []), + ), + patch("EvoScientist.cli.channel._set_channel_response") as mock_set_resp, + ): + handled = _run( + dispatch_channel_slash_command( + msg, + agent=None, + thread_id="t1", + workspace_dir=None, + checkpointer=None, + append_system=append, + await_agent_ready=_await_ready, + ) + ) + assert handled is True + mock_set_resp.assert_called_once() + resp_text = mock_set_resp.call_args[0][1] + assert "Command error" in resp_text + assert "agent blew up" in resp_text + + +def test_cmd_manager_raises_returns_true_with_error(): + """If cmd_manager.execute raises past its own try/except, the helper + must absorb it, return True, and report via _set_channel_response.""" + msg = _make_msg() + fake_cmd = MagicMock() + fake_cmd.needs_agent.return_value = False + append = MagicMock() + with ( + patch( + "EvoScientist.commands.manager.manager.resolve", + return_value=(fake_cmd, []), + ), + patch( + "EvoScientist.commands.manager.manager.execute", + new=AsyncMock(side_effect=RuntimeError("boom")), + ), + patch("EvoScientist.cli.channel._set_channel_response") as mock_set_resp, + ): + handled = _run( + dispatch_channel_slash_command( + msg, + agent=None, + thread_id="t1", + workspace_dir=None, + checkpointer=None, + append_system=append, + ) + ) + assert handled is True + mock_set_resp.assert_called_once() + resp_text = mock_set_resp.call_args[0][1] + assert "Command error" in resp_text + assert "boom" in resp_text + + +def test_on_cmd_completed_awaited_with_ctx_original_agent_and_cmd(): + """After a successful slash execute, the on_cmd_completed hook must + be awaited with (ctx, original_agent, cmd) so Rich CLI can adopt an + ``/model`` agent swap and refresh status for state-mutating commands.""" + msg = _make_msg() + fake_cmd = MagicMock() + fake_cmd.needs_agent.return_value = False + fake_cmd.name = "/model" + append = MagicMock() + captured: dict[str, object] = {} + + async def _fake_execute(command_str, ctx): + # Simulate /model swapping ctx.agent to a new handle. + ctx.agent = "swapped-agent" + return True + + async def _on_completed(ctx, original_agent, cmd): + captured["ctx_agent"] = ctx.agent + captured["original_agent"] = original_agent + captured["cmd_name"] = getattr(cmd, "name", None) + + with ( + patch( + "EvoScientist.commands.manager.manager.resolve", + return_value=(fake_cmd, []), + ), + patch( + "EvoScientist.commands.manager.manager.execute", + new=_fake_execute, + ), + patch("EvoScientist.cli.channel._set_channel_response"), + ): + handled = _run( + dispatch_channel_slash_command( + msg, + agent="original-agent", + thread_id="t1", + workspace_dir=None, + checkpointer=None, + append_system=append, + on_cmd_completed=_on_completed, + ) + ) + assert handled is True + assert captured["ctx_agent"] == "swapped-agent" + assert captured["original_agent"] == "original-agent" + assert captured["cmd_name"] == "/model" + + +def test_on_cmd_completed_receives_cmd_for_new_and_compact(): + """``/new`` / ``/compact`` invoked via channel must flow the cmd into + the hook so the callback can still refresh status when the agent + didn't swap — mirrors REPL ``interactive.py:1027-1030``.""" + for cmd_name in ("/new", "/compact"): + fake_cmd = MagicMock() + fake_cmd.needs_agent.return_value = False + fake_cmd.name = cmd_name + captured: dict[str, str | None] = {} + + async def _on_completed(ctx, original_agent, cmd, _captured=captured): + _captured["cmd_name"] = getattr(cmd, "name", None) + + with ( + patch( + "EvoScientist.commands.manager.manager.resolve", + return_value=(fake_cmd, []), + ), + patch( + "EvoScientist.commands.manager.manager.execute", + new=AsyncMock(return_value=True), + ), + patch("EvoScientist.cli.channel._set_channel_response"), + ): + _run( + dispatch_channel_slash_command( + _make_msg(content=cmd_name), + agent="same-agent", + thread_id="t1", + workspace_dir=None, + checkpointer=None, + append_system=MagicMock(), + on_cmd_completed=_on_completed, + ) + ) + assert captured["cmd_name"] == cmd_name, cmd_name + + +def test_on_cmd_completed_skipped_on_fall_through_and_error(): + """The hook must NOT fire for unresolved slash, non-slash text, or + when cmd_manager.execute raised.""" + fake_cmd = MagicMock() + fake_cmd.needs_agent.return_value = False + completed = MagicMock() + + async def _noop(*args, **kwargs): + completed(*args, **kwargs) + + # Non-slash + with patch("EvoScientist.cli.channel._set_channel_response"): + _run( + dispatch_channel_slash_command( + _make_msg(content="hi"), + agent=None, + thread_id="t1", + workspace_dir=None, + checkpointer=None, + append_system=MagicMock(), + on_cmd_completed=_noop, + ) + ) + # Unresolved slash + with ( + patch( + "EvoScientist.commands.manager.manager.resolve", + return_value=None, + ), + patch("EvoScientist.cli.channel._set_channel_response"), + ): + _run( + dispatch_channel_slash_command( + _make_msg(content="/nope"), + agent=None, + thread_id="t1", + workspace_dir=None, + checkpointer=None, + append_system=MagicMock(), + on_cmd_completed=_noop, + ) + ) + # Execute raises + with ( + patch( + "EvoScientist.commands.manager.manager.resolve", + return_value=(fake_cmd, []), + ), + patch( + "EvoScientist.commands.manager.manager.execute", + new=AsyncMock(side_effect=RuntimeError("boom")), + ), + patch("EvoScientist.cli.channel._set_channel_response"), + ): + _run( + dispatch_channel_slash_command( + _make_msg(), + agent=None, + thread_id="t1", + workspace_dir=None, + checkpointer=None, + append_system=MagicMock(), + on_cmd_completed=_noop, + ) + ) + + completed.assert_not_called() + + +def test_on_cmd_completed_exception_is_absorbed(): + """A raising hook must NOT prevent the channel response from being set.""" + msg = _make_msg() + fake_cmd = MagicMock() + fake_cmd.needs_agent.return_value = False + + async def _boom(ctx, original_agent, cmd): + raise RuntimeError("hook blew up") + + with ( + patch( + "EvoScientist.commands.manager.manager.resolve", + return_value=(fake_cmd, []), + ), + patch( + "EvoScientist.commands.manager.manager.execute", + new=AsyncMock(return_value=True), + ), + patch("EvoScientist.cli.channel._set_channel_response") as mock_set_resp, + ): + handled = _run( + dispatch_channel_slash_command( + msg, + agent="orig", + thread_id="t1", + workspace_dir=None, + checkpointer=None, + append_system=MagicMock(), + on_cmd_completed=_boom, + ) + ) + assert handled is True + mock_set_resp.assert_called_once() + assert "Command executed" in mock_set_resp.call_args[0][1] + + +def test_top_level_exception_is_absorbed(): + """Last-ditch safety net: if anything inside the dispatch pipeline + raises unexpectedly (lazy import failure, ChannelCommandUI ctor, + terminal I/O from append_system, ...), the helper must NOT + propagate — it sets an error response and returns True so the + caller's polling loop stays alive and doesn't fall through to the + agent streaming path.""" + msg = _make_msg() + with ( + patch( + "EvoScientist.commands.manager.manager.resolve", + side_effect=RuntimeError("exploded during resolve"), + ), + patch("EvoScientist.cli.channel._set_channel_response") as mock_set_resp, + ): + handled = _run( + dispatch_channel_slash_command( + msg, + agent=None, + thread_id="t1", + workspace_dir=None, + checkpointer=None, + append_system=MagicMock(), + ) + ) + assert handled is True + mock_set_resp.assert_called_once() + resp_text = mock_set_resp.call_args[0][1] + assert "Command error" in resp_text + assert "exploded during resolve" in resp_text + + +def test_cmd_execute_returning_false_falls_through(): + """When cmd_manager.execute returns False (empty/unparseable input), + the helper must return False so the caller falls through to the agent.""" + msg = _make_msg(content="/") + fake_cmd = MagicMock() + fake_cmd.needs_agent.return_value = False + append = MagicMock() + with ( + patch( + "EvoScientist.commands.manager.manager.resolve", + return_value=(fake_cmd, []), + ), + patch( + "EvoScientist.commands.manager.manager.execute", + new=AsyncMock(return_value=False), + ), + patch("EvoScientist.cli.channel._set_channel_response") as mock_set_resp, + ): + handled = _run( + dispatch_channel_slash_command( + msg, + agent=None, + thread_id="t1", + workspace_dir=None, + checkpointer=None, + append_system=append, + ) + ) + assert handled is False + mock_set_resp.assert_not_called() diff --git a/tests/test_serve_agent_holder.py b/tests/test_serve_agent_holder.py new file mode 100644 index 0000000..722ae7c --- /dev/null +++ b/tests/test_serve_agent_holder.py @@ -0,0 +1,286 @@ +"""Tests for the serve-mode ``on_cmd_completed`` hook factory. + +Regression coverage for the follow-up to issue #181 — `/model` invoked +over a channel in ``EvoSci serve`` must swap the running agent for +subsequent messages, not silently keep the stale one the while-loop +captured at startup. +""" + +from unittest.mock import AsyncMock, MagicMock, patch + +import pytest + +from EvoScientist.cli.channel import ChannelMessage +from EvoScientist.cli.commands import ( + _make_serve_cmd_completed_hook, + _make_serve_start_new_session_cb, + _serve_process_message, +) +from tests.conftest import run_async as _run + + +@pytest.fixture(autouse=True) +def _restore_channel_globals(): + import EvoScientist.cli.channel as _ch_mod + + prev_agent = _ch_mod._cli_agent + prev_tid = _ch_mod._cli_thread_id + try: + yield + finally: + _ch_mod._cli_agent = prev_agent + _ch_mod._cli_thread_id = prev_tid + + +def test_hook_updates_holder_on_agent_swap(): + """``/model`` mutates ``ctx.agent`` to a new handle — the hook must + push that handle into the shared holder so the outer poll loop sees + it on the next message.""" + holder = {"agent": "original-agent"} + hook = _make_serve_cmd_completed_hook(holder) + + ctx = MagicMock() + ctx.agent = "new-agent" + cmd = MagicMock() + cmd.name = "/model" + + _run(hook(ctx, "original-agent", cmd)) + + assert holder["agent"] == "new-agent" + + +def test_hook_syncs_channel_module_global(): + """Other readers (the bus) look at ``cli.channel._cli_agent``; the + hook keeps that global in sync with the holder update.""" + import EvoScientist.cli.channel as _ch_mod + + holder = {"agent": "original-agent"} + hook = _make_serve_cmd_completed_hook(holder) + + ctx = MagicMock() + ctx.agent = "new-agent" + cmd = MagicMock() + cmd.name = "/model" + + _run(hook(ctx, "original-agent", cmd)) + + assert _ch_mod._cli_agent == "new-agent" + + +def test_hook_noop_when_agent_unchanged(): + """Commands like ``/evoskills`` don't touch ``ctx.agent`` — the + holder must stay put.""" + holder = {"agent": "original-agent"} + hook = _make_serve_cmd_completed_hook(holder) + + ctx = MagicMock() + ctx.agent = "original-agent" # no swap + cmd = MagicMock() + cmd.name = "/evoskills" + + _run(hook(ctx, "original-agent", cmd)) + + assert holder["agent"] == "original-agent" + + +def test_hook_noop_when_ctx_agent_is_none(): + """Guard against commands that reset ``ctx.agent`` to ``None`` — + we never want to write ``None`` into the holder.""" + holder = {"agent": "original-agent"} + hook = _make_serve_cmd_completed_hook(holder) + + ctx = MagicMock() + ctx.agent = None + cmd = MagicMock() + cmd.name = "/whatever" + + _run(hook(ctx, "original-agent", cmd)) + + assert holder["agent"] == "original-agent" + + +def test_hook_updates_thread_id_on_resume(): + """``/resume`` mutates ``ctx.thread_id`` — the hook must push the + new id into the holder so the outer poll loop runs subsequent + messages on the resumed thread.""" + holder = {"agent": "a", "thread_id": "original-tid"} + hook = _make_serve_cmd_completed_hook(holder) + + ctx = MagicMock() + ctx.agent = "a" # no agent swap + ctx.thread_id = "new-tid" + cmd = MagicMock() + cmd.name = "/resume" + + _run(hook(ctx, "a", cmd)) + + assert holder["thread_id"] == "new-tid" + + +def test_hook_syncs_channel_module_thread_id(): + """The bus reads ``cli.channel._cli_thread_id``; hook must sync it + alongside the holder update.""" + import EvoScientist.cli.channel as _ch_mod + + holder = {"agent": "a", "thread_id": "original-tid"} + hook = _make_serve_cmd_completed_hook(holder) + + ctx = MagicMock() + ctx.agent = "a" + ctx.thread_id = "new-tid" + cmd = MagicMock() + cmd.name = "/resume" + + _run(hook(ctx, "a", cmd)) + + assert _ch_mod._cli_thread_id == "new-tid" + + +def test_hook_noop_when_thread_id_unchanged(): + """Most commands don't touch thread_id — holder stays put.""" + holder = {"agent": "a", "thread_id": "same-tid"} + hook = _make_serve_cmd_completed_hook(holder) + + ctx = MagicMock() + ctx.agent = "a" + ctx.thread_id = "same-tid" + cmd = MagicMock() + cmd.name = "/evoskills" + + _run(hook(ctx, "a", cmd)) + + assert holder["thread_id"] == "same-tid" + + +def test_hook_skips_resume_warning_when_thread_unchanged(): + """Bare ``/resume`` with no argument prints usage but leaves + ``ctx.thread_id`` unchanged — the in-memory-state warning must NOT + fire because no resume actually happened.""" + holder = {"agent": "a", "thread_id": "original-tid"} + hook = _make_serve_cmd_completed_hook(holder) + + ctx = MagicMock() + ctx.agent = "a" + ctx.thread_id = "original-tid" # unchanged — bare /resume case + cmd = MagicMock() + cmd.name = "/resume" + + _run(hook(ctx, "a", cmd)) + + ctx.ui.append_system.assert_not_called() + ctx.ui.flush.assert_not_called() + + +def test_hook_emits_resume_warning_when_thread_changed(): + """``/resume `` that actually changes thread_id must surface + the in-memory-state warning via ``ctx.ui``.""" + holder = {"agent": "a", "thread_id": "original-tid"} + hook = _make_serve_cmd_completed_hook(holder) + + ctx = MagicMock() + # Mock out async flush so the test can synchronously run the hook. + ctx.ui.flush = AsyncMock() + ctx.agent = "a" + ctx.thread_id = "abc12345-resumed-tid" + cmd = MagicMock() + cmd.name = "/resume" + + _run(hook(ctx, "a", cmd)) + + ctx.ui.append_system.assert_called_once() + warn_text, warn_kwargs = ( + ctx.ui.append_system.call_args.args, + ctx.ui.append_system.call_args.kwargs, + ) + assert "in-memory state" in warn_text[0] + assert "abc12345" in warn_text[0] + assert warn_kwargs.get("style") == "yellow" + ctx.ui.flush.assert_awaited_once() + + +def test_start_new_session_cb_rotates_thread_id(): + """``/new`` via channel calls this callback — must generate a new + thread id, push into holder, and sync the channel-module global.""" + import EvoScientist.cli.channel as _ch_mod + + holder = {"agent": "a", "thread_id": "old-tid"} + + with patch( + "EvoScientist.sessions.generate_thread_id", + return_value="freshly-generated-tid", + ): + cb = _make_serve_start_new_session_cb(holder) + cb() + + assert holder["thread_id"] == "freshly-generated-tid" + assert _ch_mod._cli_thread_id == "freshly-generated-tid" + + +def test_start_new_session_cb_leaves_agent_alone(): + """``/new`` rotates thread only — agent handle must stay put + (serve's agent is a single pre-loaded instance, not per-thread).""" + holder = {"agent": "a", "thread_id": "old-tid"} + + with patch( + "EvoScientist.sessions.generate_thread_id", + return_value="new-tid", + ): + cb = _make_serve_start_new_session_cb(holder) + cb() + + assert holder["agent"] == "a" + + +def test_hook_handles_both_agent_and_thread_swap(): + """Edge case: a command that changes both (hypothetical). Both + updates must land in the holder.""" + holder = {"agent": "old-agent", "thread_id": "old-tid"} + hook = _make_serve_cmd_completed_hook(holder) + + ctx = MagicMock() + ctx.agent = "new-agent" + ctx.thread_id = "new-tid" + cmd = MagicMock() + + _run(hook(ctx, "old-agent", cmd)) + + assert holder["agent"] == "new-agent" + assert holder["thread_id"] == "new-tid" + + +def test_serve_process_message_reports_slash_dispatch_error_without_fallback(): + """Defensive: if ``dispatch_channel_slash_command`` ever leaks an + exception past its own wrapper, ``_serve_process_message`` must set + one error response and not fall through to ``run_streaming``. + """ + msg = ChannelMessage( + msg_id="msg-1", + content="/evoskills core", + sender="channel-user", + channel_type="imessage", + metadata={}, + channel_ref=None, + bus_ref=None, + chat_id="channel-user", + message_id="ts-1", + ) + holder = {"agent": "agent", "thread_id": "tid"} + + with ( + patch( + "EvoScientist.cli.commands.dispatch_channel_slash_command", + new=AsyncMock(side_effect=RuntimeError("slash broke")), + ), + patch("EvoScientist.cli.commands._set_channel_response") as mock_set_resp, + patch("EvoScientist.cli.tui_runtime.run_streaming") as mock_run_streaming, + ): + _serve_process_message( + msg, + agent_holder=holder, + model="model", + workspace_dir="/tmp", + show_thinking=False, + ) + + mock_set_resp.assert_called_once_with("msg-1", "Command error: slash broke") + mock_run_streaming.assert_not_called()