From 1449280fcd9cb92ed9fb80af018cdb6cd822b76a Mon Sep 17 00:00:00 2001 From: Athena Kaminsky Date: Fri, 21 Aug 2026 15:39:28 -0500 Subject: [PATCH 1/5] feat(cli): add interactive TUI chat Add a Textual chat interface with threaded SSE streaming, slash commands, interrupt handling, and bare nexus dispatch. Package it behind the tui extra, document usage, and cover command routing, dependencies, and headless interaction with tests. --- CLAUDE.md | 29 +-- docs/CLI.md | 3 + nexusos_cli/cli.py | 39 +++- nexusos_cli/tui_app.py | 426 +++++++++++++++++++++++++++++++++++ pyproject.toml | 2 + tests/test_cli_packaging.py | 17 +- tests/test_packaging_deps.py | 3 +- tests/test_tui.py | 85 +++++++ 8 files changed, 586 insertions(+), 18 deletions(-) create mode 100644 nexusos_cli/tui_app.py create mode 100644 tests/test_tui.py diff --git a/CLAUDE.md b/CLAUDE.md index 824d045..6bffd47 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -50,25 +50,28 @@ uvicorn synapse.main:sio_app --host 127.0.0.1 --port 8000 --reload cd interface/web && npm run dev ``` -**Management CLI** (`ncp`) — start/stop services with PID tracking, plus terminal -access to the same features as the web UI (all via the REST API on `:8000`): +**Management CLI** (`nexus` / `ncp`) — start/stop services with PID tracking, plus +terminal access to the same features as the web UI (REST API on `:8000`): ```bash ./management/nexus-cli.sh start # starts backend + frontend ./management/nexus-cli.sh stop ./management/nexus-cli.sh start --backend|-b / --frontend|-f / --memory|-m -# Feature commands (dispatch to nexusos_cli/nexus_api.py — httpx, no TUI): -ncp chat "" # stream a reply (POST /chat/stream) -ncp memory list|add |rm -ncp playbook list|show # first playbook (*) is the active system prompt -ncp history [query] # recent conversations +# Interactive TUI (Hermes/OpenClaw-style; needs pip install 'nexusos-ai[tui]'): +nexus # bare command opens the Textual chat TUI +nexus tui # same, explicit + +# Feature one-shots (dispatch to nexusos_cli/nexus_api.py — httpx): +nexus chat send "" # stream a reply (POST /chat/stream) +nexus memory list|add |rm +nexus playbook list|show # first playbook (*) is the active system prompt +nexus history [query] # recent conversations +nexus monitor # ASCII status dashboard (no prompt) ``` -The old curses TUIs (`nexus-chat.py`, `nexus-playbook.py`) were removed in favor of -these API-backed subcommands. The CLI covers chat, memory, playbooks, and history; -the web UI and control panel expose the remaining management features. -The CLI itself lives in `nexusos_cli/` (that is what the wheel ships and what -`nexus`/`ncp`/`nexusos` dispatch to); `management/` keeps the desktop-only -pieces — the shell wrappers, the Tk control panel, and the XFCE panel wiring. +The interactive TUI lives in `nexusos_cli/tui_app.py` (Textual, optional extra). +One-shot subcommands and `nexus monitor` remain for scripts. The CLI package is +`nexusos_cli/` (what the wheel ships); `management/` keeps desktop-only pieces — +shell wrappers, Tk control panel, XFCE panel wiring. `management/controlpanel.py` (tkinter GUI, wired into the XFCE panel via `bin/panel/nexus-popup.py`) stays. diff --git a/docs/CLI.md b/docs/CLI.md index 440a7a8..e8e2036 100644 --- a/docs/CLI.md +++ b/docs/CLI.md @@ -37,11 +37,14 @@ and seed playbooks. Extras keep platform-sensitive dependencies optional: - `desktop`: desktop process support and Windows pywebview - `search`: DuckDuckGo web search for chat - `mail`: IMAP mail reading +- `tui`: Textual interactive chat UI (`nexus` with no subcommand) - `all`: every optional capability at once ## Common commands ```text +nexus Interactive chat TUI (needs nexusos-ai[tui]) +nexus tui Same as bare nexus nexus init Create writable state and seed playbooks nexus doctor [--fix] [--json] Diagnose the install and provider nexus paths [--json] Show package, state, and asset locations diff --git a/nexusos_cli/cli.py b/nexusos_cli/cli.py index a5c0e93..43807c7 100644 --- a/nexusos_cli/cli.py +++ b/nexusos_cli/cli.py @@ -398,6 +398,24 @@ def cmd_monitor(args) -> int: ) +def cmd_tui(args) -> int: + """Interactive Hermes/OpenClaw-style chat TUI (requires nexusos-ai[tui]).""" + if not sys.stdin.isatty() or not sys.stdout.isatty(): + print( + "The TUI needs a terminal. Use: nexus chat send \"…\"\n" + "Or run `nexus` in an interactive shell.", + file=sys.stderr, + ) + return 2 + try: + from .tui_app import run_tui + except ImportError as exc: + print(str(exc), file=sys.stderr) + return 2 + api = getattr(args, "api_url", None) or settings.api_url + return run_tui(api_url=api) + + def _target_flag(target: str | None): return { "memory": "--memory", @@ -722,10 +740,20 @@ def _port(value: str) -> int: def build_parser() -> argparse.ArgumentParser: - parser = argparse.ArgumentParser(prog="nexus", description="NexusOS local AI runtime and API client") + parser = argparse.ArgumentParser( + prog="nexus", + description=( + "NexusOS local AI runtime and API client. " + "With no subcommand, opens the interactive TUI (needs nexusos-ai[tui])." + ), + ) parser.add_argument("--version", action="version", version=f"NexusOS {settings.version}") parser.add_argument("--api-url", help="override the NexusOS backend URL for this command") - sub = parser.add_subparsers(dest="command", required=True) + # Bare `nexus` → TUI. Subcommands remain for scripts and one-shots. + sub = parser.add_subparsers(dest="command", required=False) + + p = sub.add_parser("tui", help="interactive chat TUI (default when no subcommand)") + p.set_defaults(fn=cmd_tui) p = sub.add_parser("init", help="create user state and seed default playbooks"); _add_json(p); p.set_defaults(fn=cmd_init) p = sub.add_parser("paths", help="show resolved package and writable paths"); _add_json(p); p.set_defaults(fn=cmd_paths) @@ -827,7 +855,12 @@ def _normalize_legacy_argv(argv) -> list[str]: def main(argv=None) -> int: parser = build_parser() - args = parser.parse_args(_normalize_legacy_argv(argv)) + argv = _normalize_legacy_argv(argv) + args = parser.parse_args(argv) + if not getattr(args, "command", None): + # Bare `nexus` / `ncp` / `nexusos` → interactive TUI. + args.command = "tui" + args.fn = cmd_tui if args.command == "config": if args.action in ("get", "unset") and not args.key: parser.error(f"config {args.action} requires KEY") diff --git a/nexusos_cli/tui_app.py b/nexusos_cli/tui_app.py new file mode 100644 index 0000000..c3084ee --- /dev/null +++ b/nexusos_cli/tui_app.py @@ -0,0 +1,426 @@ +"""Hermes/OpenClaw-style interactive TUI for NexusOS. + +Optional: needs the ``tui`` extra (Textual). Launched by a bare ``nexus`` when +stdin/stdout are a TTY. Classic one-shots (``nexus chat send``, ``nexus monitor``, +``nexus status``, …) stay on the argparse tree. +""" +from __future__ import annotations + +import json +import threading +import uuid +from typing import Any + +import httpx + +from synapse.nexus_config import settings + +from . import nexus_api +from .monitor import collect_snapshot + +# Between SSE chunks a silent backend must not pin the UI forever. Connect stays +# short; the overall stream may run minutes. +_STREAM_TIMEOUT = httpx.Timeout(None, connect=5.0, read=120.0, write=30.0, pool=5.0) + + +def _require_textual(): + try: + from textual.app import App + from textual.widgets import Footer, Header, Input, RichLog, Static + except ImportError as e: # pragma: no cover - optional extra + raise ImportError( + "The interactive TUI needs the 'tui' extra — " + "pip install 'nexusos-ai[tui]' (or: pip install textual)." + ) from e + return App, Footer, Header, Input, RichLog, Static + + +def _escape(text: str) -> str: + """Make model/user text safe for Rich markup widgets. + + Rich's own escape is the only version that round-trips. Escaping every + backslash by hand looks equivalent but is not: Rich un-escapes ``\\[`` and + never collapses ``\\\\``, so doubling them puts the doubles on screen - + every Windows path and regex escape in a reply renders wrong. + + Imported inside the function so the module still loads without the ``tui`` + extra. Rich is not declared in pyproject: Textual depends on it, so it is + present whenever the TUI can run at all, and tests/test_packaging_deps.py + lists it in TRANSITIVE for that reason. + """ + from rich.markup import escape + + return escape(text) + + +def format_user_line(message: str) -> str: + return f"[bold green]you>[/] {_escape(message)}" + + +def format_assistant_line(text: str) -> str: + return f"[bold blue]nexus>[/] {_escape(text)}" + + +def _status_line(snap: dict | None = None) -> str: + """Format a snapshot. Pass ``snap`` — do not omit it on the UI thread.""" + if snap is None: + snap = collect_snapshot() + svcs = snap.get("services") or {} + api = snap.get("api") or {} + host = snap.get("host") or {} + parts = [f"NexusOS {snap.get('version', '')}"] + for key in ("backend", "memory", "provider"): + info = svcs.get(key) or {} + if key == "provider": + up = bool(info.get("reachable")) + else: + up = bool(info.get("running")) + parts.append(f"{key}={'UP' if up else 'DOWN'}") + if api.get("online"): + parts.append(f"tools={api.get('action_tool_policy') or '—'}") + cpu = host.get("cpu_pct") + if cpu is not None: + parts.append(f"cpu={cpu:.0f}%") + chains = snap.get("toolchains") or [] + ready = [c["lang"] for c in chains if c.get("ready")] + if ready: + parts.append("run=" + ",".join(ready)) + return " · ".join(parts) + + +def _compact_status(snap: dict | None = None) -> str: + """One-line strip for the bar under the chat log.""" + if snap is None: + snap = collect_snapshot() + host = snap.get("host") or {} + api = snap.get("api") or {} + recent = snap.get("recent_tools") or [] + cpu = host.get("cpu_pct") + mem = host.get("mem_pct") + bits = [] + if cpu is not None: + bits.append(f"cpu {cpu:.0f}%") + if mem is not None: + bits.append(f"mem {mem:.0f}%") + if api.get("online"): + bits.append( + f"memories={api.get('memories') if api.get('memories') is not None else '—'} " + f"chats={api.get('conversations') if api.get('conversations') is not None else '—'}" + ) + else: + bits.append("api DOWN — nexus start") + if recent: + bits.append("recent " + ", ".join(recent[:4])) + return " │ ".join(bits) + + +class NexusTUI: + """Factory so Textual imports stay lazy until run().""" + + @staticmethod + def build_app(*, api_url: str | None = None): + App, Footer, Header, Input, RichLog, Static = _require_textual() + base = (api_url or settings.api_url).rstrip("/") + + class AppImpl(App): + CSS = """ + Screen { layout: vertical; } + #status { + height: 1; + dock: top; + background: $boost; + color: $text; + padding: 0 1; + } + #strip { + height: 1; + background: $surface; + color: $text-muted; + padding: 0 1; + } + #log { + height: 1fr; + border: tall $accent; + padding: 0 1; + } + #live { + height: auto; + max-height: 8; + padding: 0 1; + color: $text; + } + #prompt { dock: bottom; } + """ + BINDINGS = [ + ("ctrl+c", "interrupt", "Interrupt"), + ("ctrl+d", "quit", "Quit"), + ] + + def __init__(self): + super().__init__() + self.api_url = base + self.conversation_id: str | None = None + self.history: list[dict] = [] + self._model: str | None = None + self._busy = False + self._stop_stream = threading.Event() + self._http_client: httpx.Client | None = None + self._status_lock = threading.Lock() + self._status_pending = False + + def compose(self): + # Placeholders only — never collect_snapshot() on the UI thread. + yield Header(show_clock=True) + yield Static("NexusOS …", id="status") + yield RichLog(id="log", highlight=True, markup=True, wrap=True) + yield Static("", id="live") + yield Static("collecting status…", id="strip") + yield Input( + placeholder="Message Nexus… (/help for commands)", + id="prompt", + ) + yield Footer() + + def on_mount(self) -> None: + self.title = "NexusOS" + self.sub_title = self.api_url + log = self.query_one("#log", RichLog) + log.write("[bold]NexusOS[/] interactive TUI") + log.write( + "Type a message and Enter. " + "Slash: /help /status /new /model /quit" + ) + log.write(f"API: {_escape(self.api_url)}") + log.write("") + self._schedule_status_refresh() + self.set_interval(2.0, self._schedule_status_refresh) + self.query_one("#prompt", Input).focus() + + def _schedule_status_refresh(self) -> None: + """Kick a worker; never call collect_snapshot on the event loop.""" + with self._status_lock: + if self._status_pending: + return + self._status_pending = True + + def worker(): + try: + snap = collect_snapshot() + self._call_ui(self._apply_status, snap) + except Exception: + pass + finally: + with self._status_lock: + self._status_pending = False + + threading.Thread(target=worker, daemon=True).start() + + def _apply_status(self, snap: dict) -> None: + self.query_one("#status", Static).update(_status_line(snap)) + self.query_one("#strip", Static).update(_compact_status(snap)) + + def _call_ui(self, callback, *args) -> None: + """call_from_thread, but never after quit (avoids CancelledError + traceback garbling the restored shell).""" + if not self.is_running: + return + try: + self.call_from_thread(callback, *args) + except BaseException: + # CancelledError is BaseException; also ignore post-exit races. + pass + + def action_quit(self) -> None: + self._stop_stream.set() + client = self._http_client + if client is not None: + try: + client.close() + except Exception: + pass + self.exit() + + def action_interrupt(self) -> None: + if self._busy: + self._stop_stream.set() + client = self._http_client + if client is not None: + try: + client.close() + except Exception: + pass + self.query_one("#log", RichLog).write( + "[yellow]▸ interrupt requested[/]" + ) + else: + self.exit() + + def on_input_submitted(self, event: Input.Submitted) -> None: + text = (event.value or "").strip() + event.input.value = "" + if not text: + return + if text.startswith("/"): + self._handle_slash(text) + return + if self._busy: + self.query_one("#log", RichLog).write( + "[yellow]Still streaming — wait or Ctrl+C to interrupt[/]" + ) + return + self._start_chat(text) + + def _handle_slash(self, text: str) -> None: + log = self.query_one("#log", RichLog) + cmd, _, rest = text[1:].partition(" ") + cmd = cmd.lower().strip() + rest = rest.strip() + if cmd in ("q", "quit", "exit"): + self.exit() + elif cmd in ("h", "help"): + log.write( + "[bold]/help[/] this list\n" + "[bold]/status[/] refresh service strip\n" + "[bold]/new[/] fresh conversation\n" + "[bold]/model[/] \\[name] pin model for next turns\n" + "[bold]/quit[/] leave the TUI\n" + "One-shot: [dim]nexus chat send \"…\"[/]" + ) + elif cmd == "status": + self._schedule_status_refresh() + log.write("[dim]refreshing status…[/]") + elif cmd == "new": + self.conversation_id = None + self.history = [] + log.write("[bold cyan]— new conversation —[/]") + elif cmd == "model": + if rest: + self._model = rest + log.write(f"[dim]model pinned:[/] {_escape(rest)}") + else: + log.write( + f"[dim]model:[/] {_escape(self._model or '(auto)')}" + ) + else: + log.write( + f"[red]unknown command[/] /{_escape(cmd)} — try /help" + ) + + def _start_chat(self, message: str) -> None: + log = self.query_one("#log", RichLog) + live = self.query_one("#live", Static) + log.write(format_user_line(message)) + live.update("[bold blue]nexus>[/] [dim]…[/]") + self._busy = True + self._stop_stream.clear() + if not self.conversation_id: + self.conversation_id = str(uuid.uuid4()) + body: dict[str, Any] = { + "message": message, + "conversation_id": self.conversation_id, + "history": list(self.history), + } + if self._model: + body["model"] = self._model + self.history.append({"role": "user", "content": message}) + + def worker(): + reply_parts: list[str] = [] + client = httpx.Client( + base_url=self.api_url, timeout=_STREAM_TIMEOUT + ) + self._http_client = client + try: + with client.stream( + "POST", "/chat/stream", json=body + ) as resp: + if resp.status_code >= 400: + detail = resp.read().decode( + "utf-8", errors="replace" + )[:300] + self._call_ui( + live.update, + f"[red]error HTTP {resp.status_code}[/] " + f"{_escape(detail)}", + ) + return + for kind, payload in nexus_api.iter_chunks( + resp.iter_lines() + ): + if self._stop_stream.is_set(): + break + if kind == "chunk": + reply_parts.append(payload) + preview = "".join(reply_parts) + if len(preview) > 4000: + preview = "…" + preview[-4000:] + self._call_ui( + live.update, + format_assistant_line(preview), + ) + elif kind == "tool_request": + self._call_ui( + log.write, + "[yellow]▸ tool approval needed — " + "Approve in the web UI, or set " + "action_tool_policy=allow[/]", + ) + elif kind == "error": + try: + detail = json.loads(payload).get( + "detail", payload + ) + except Exception: + detail = payload + self._call_ui( + live.update, + f"[red]error:[/] {_escape(str(detail))}", + ) + elif kind == "done": + break + except httpx.ConnectError: + self._call_ui( + live.update, + f"[red]Backend not reachable at " + f"{_escape(self.api_url)}. Start it: nexus start[/]", + ) + except Exception as exc: + self._call_ui( + live.update, + f"[red]{_escape(type(exc).__name__)}:[/] " + f"{_escape(str(exc))}", + ) + finally: + self._http_client = None + try: + client.close() + except Exception: + pass + text = "".join(reply_parts).strip() + self._call_ui(self._finish_stream, text) + + threading.Thread(target=worker, daemon=True).start() + + def _finish_stream(self, text: str) -> None: + log = self.query_one("#log", RichLog) + live = self.query_one("#live", Static) + try: + if text: + log.write(format_assistant_line(text)) + self.history.append( + {"role": "assistant", "content": text} + ) + finally: + # Always clear busy — a MarkupError must not wedge the TUI. + live.update("") + self._busy = False + self._schedule_status_refresh() + + return AppImpl() + + +def run_tui(*, api_url: str | None = None) -> int: + """Run the Textual app. Returns a process exit code.""" + app = NexusTUI.build_app(api_url=api_url) + app.run() + return 0 diff --git a/pyproject.toml b/pyproject.toml index f4a8cfc..c0eede4 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -42,6 +42,7 @@ mail = ["imap-tools>=1.7,<2"] # synapse/search.py imports this lazily behind a bare except, so without it # declared the chat web-search path silently returns nothing. search = ["duckduckgo-search>=6,<9"] +tui = ["textual>=1.0,<3"] desktop = [ "psutil>=5.9,<8", "pywebview>=5,<7; platform_system == 'Windows'", @@ -64,6 +65,7 @@ all = [ "faster-whisper>=1.1,<2", "imap-tools>=1.7,<2", "duckduckgo-search>=6,<9", + "textual>=1.0,<3", "pywebview>=5,<7; platform_system == 'Windows'", ] dev = [ diff --git a/tests/test_cli_packaging.py b/tests/test_cli_packaging.py index 873ea78..d0095be 100644 --- a/tests/test_cli_packaging.py +++ b/tests/test_cli_packaging.py @@ -31,8 +31,23 @@ def run_cli(tmp_path: Path, *args: str) -> subprocess.CompletedProcess[str]: def test_help_exposes_portable_command_tree(tmp_path): result = run_cli(tmp_path, "--help") assert result.returncode == 0, result.stderr - for command in ("init", "config", "provider", "doctor", "serve", "models", "chat", "monitor"): + for command in ("init", "config", "provider", "doctor", "serve", "models", "chat", "monitor", "tui"): assert command in result.stdout + assert "interactive TUI" in result.stdout or "TUI" in result.stdout + + +def test_bare_nexus_defaults_to_tui_command(): + """No subcommand → TUI entry (Hermes-style). Non-TTY exits 2 without launching.""" + from unittest import mock + + from nexusos_cli.cli import build_parser, cmd_tui + + parser = build_parser() + args = parser.parse_args([]) + assert args.command is None # filled in by main() + with mock.patch("sys.stdin.isatty", return_value=False), \ + mock.patch("sys.stdout.isatty", return_value=False): + assert cmd_tui(args) == 2 def test_legacy_cli_spellings_remain_compatible(): diff --git a/tests/test_packaging_deps.py b/tests/test_packaging_deps.py index ec3e4de..cde2a62 100644 --- a/tests/test_packaging_deps.py +++ b/tests/test_packaging_deps.py @@ -29,7 +29,8 @@ DISTRIBUTION_OF = { } # Provided by another declared distribution rather than named directly. -TRANSITIVE = {"starlette", "socketio", "engineio"} +# rich: Textual depends on it, so the tui extra already pulls it in. +TRANSITIVE = {"starlette", "socketio", "engineio", "rich"} # Modules that ship inside this repo. FIRST_PARTY = {"synapse", "nexusos_cli", "modules", "management", "bin", "tests"} diff --git a/tests/test_tui.py b/tests/test_tui.py new file mode 100644 index 0000000..4ca6f8a --- /dev/null +++ b/tests/test_tui.py @@ -0,0 +1,85 @@ +"""TUI helpers and headless App.run_test coverage.""" +from __future__ import annotations + +import asyncio + +import pytest +from rich.text import Text + +from nexusos_cli.tui_app import ( + _compact_status, + _escape, + _status_line, + format_assistant_line, + format_user_line, +) + + +def test_status_line_mentions_services(): + snap = { + "version": "1.0.0", + "services": { + "backend": {"running": True}, + "memory": {"running": False}, + "provider": {"reachable": True}, + }, + "api": {"online": True, "action_tool_policy": "ask"}, + "host": {"cpu_pct": 10.0}, + "toolchains": [{"lang": "python", "ready": True}], + } + line = _status_line(snap) + assert "backend=UP" in line + assert "memory=DOWN" in line + assert "provider=UP" in line + assert "tools=ask" in line + assert "run=python" in line + + +def test_compact_status_handles_api_down(): + snap = { + "host": {}, + "api": {"online": False}, + "recent_tools": [], + } + assert "api DOWN" in _compact_status(snap) + + +def test_escape_preserves_code_brackets_in_display(): + raw = "idx = arr[i] and rng = [a-z]+" + plain = Text.from_markup(format_assistant_line(raw)).plain + assert "arr[i]" in plain + assert "[a-z]+" in plain + # Unescaped markup would drop the bracket contents. + assert plain != "nexus> idx = arr and rng = +" + + +def test_closing_tag_in_model_output_does_not_raise(): + raw = "close with [/] please" + plain = Text.from_markup(format_assistant_line(raw)).plain + assert "[/]" in plain + + +def test_user_line_escapes_markup(): + plain = Text.from_markup(format_user_line("use [bold] please")).plain + assert "[bold]" in plain + + +def test_finish_stream_markup_does_not_wedge_busy(): + """A stray '[/]' used to raise before _busy=False and lock the TUI forever.""" + pytest.importorskip("textual") + from nexusos_cli.tui_app import NexusTUI + + app = NexusTUI.build_app(api_url="http://127.0.0.1:9") + + async def _run(): + async with app.run_test(): + app._busy = True + app._finish_stream("see [/] and arr[i]") + assert app._busy is False + assert app.history[-1]["content"] == "see [/] and arr[i]" + + asyncio.run(_run()) + + +def test_escape_round_trip_helper(): + assert "[" in _escape("x[y]") or "\\[" in _escape("x[y]") -- 2.39.5 From 6f5094b5fca1927a11d8a931641aee929192642f Mon Sep 17 00:00:00 2001 From: Athena Kaminsky Date: Fri, 21 Aug 2026 15:59:27 -0500 Subject: [PATCH 2/5] fix(tui): preserve errors and deny gated tools safely Keep stream failures in the persistent transcript instead of clearing them with the live preview. Use each tool request's capability token to deny actions immediately until the TUI has an interactive approval flow, and cover both behaviors with focused regressions. --- nexusos_cli/tui_app.py | 73 ++++++++++++++++++++++++++++++++++++------ tests/test_tui.py | 61 +++++++++++++++++++++++++++++++++++ 2 files changed, 124 insertions(+), 10 deletions(-) diff --git a/nexusos_cli/tui_app.py b/nexusos_cli/tui_app.py index c3084ee..54275f7 100644 --- a/nexusos_cli/tui_app.py +++ b/nexusos_cli/tui_app.py @@ -21,6 +21,7 @@ from .monitor import collect_snapshot # Between SSE chunks a silent backend must not pin the UI forever. Connect stays # short; the overall stream may run minutes. _STREAM_TIMEOUT = httpx.Timeout(None, connect=5.0, read=120.0, write=30.0, pool=5.0) +_APPROVAL_TIMEOUT = httpx.Timeout(10.0, connect=5.0) def _require_textual(): @@ -61,6 +62,40 @@ def format_assistant_line(text: str) -> str: return f"[bold blue]nexus>[/] {_escape(text)}" +def _deny_tool_request( + *, + api_url: str, + conversation_id: str, + payload: str, + client_factory=httpx.Client, +) -> list[str]: + """Immediately deny a TUI action request and let the stream resume. + + The web client presents an approval dialog, but the TUI does not yet have + that interaction. Denying with the stream's capability token preserves the + ``ask`` safety boundary without leaving the backend waiting for five minutes. + """ + request = json.loads(payload) + token = request.get("token") or "" + actions = request.get("actions") or [] + names = [ + action.get("name", "") + for action in actions + if isinstance(action, dict) and action.get("name") + ] + if not token or not names: + raise ValueError("invalid tool approval request") + body = { + "conversation_id": conversation_id, + "token": token, + "decisions": {name: False for name in names}, + } + with client_factory(base_url=api_url, timeout=_APPROVAL_TIMEOUT) as client: + response = client.post("/chat/approve", json=body) + response.raise_for_status() + return names + + def _status_line(snap: dict | None = None) -> str: """Format a snapshot. Pass ``snap`` — do not omit it on the UI thread.""" if snap is None: @@ -230,6 +265,10 @@ class NexusTUI: # CancelledError is BaseException; also ignore post-exit races. pass + def _show_error(self, message: str) -> None: + """Write a stream error to the persistent transcript.""" + self.query_one("#log", RichLog).write(message) + def action_quit(self) -> None: self._stop_stream.set() client = self._http_client @@ -339,7 +378,7 @@ class NexusTUI: "utf-8", errors="replace" )[:300] self._call_ui( - live.update, + self._show_error, f"[red]error HTTP {resp.status_code}[/] " f"{_escape(detail)}", ) @@ -359,12 +398,26 @@ class NexusTUI: format_assistant_line(preview), ) elif kind == "tool_request": - self._call_ui( - log.write, - "[yellow]▸ tool approval needed — " - "Approve in the web UI, or set " - "action_tool_policy=allow[/]", - ) + try: + names = _deny_tool_request( + api_url=self.api_url, + conversation_id=self.conversation_id, + payload=payload, + ) + shown = ", ".join(names) + self._call_ui( + log.write, + "[yellow]▸ denied action tool " + f"{_escape(shown)} — interactive approval " + "is not yet available in the TUI[/]", + ) + except Exception as exc: + self._call_ui( + self._show_error, + "[red]tool denial failed:[/] " + f"{_escape(str(exc))}", + ) + return elif kind == "error": try: detail = json.loads(payload).get( @@ -373,20 +426,20 @@ class NexusTUI: except Exception: detail = payload self._call_ui( - live.update, + self._show_error, f"[red]error:[/] {_escape(str(detail))}", ) elif kind == "done": break except httpx.ConnectError: self._call_ui( - live.update, + self._show_error, f"[red]Backend not reachable at " f"{_escape(self.api_url)}. Start it: nexus start[/]", ) except Exception as exc: self._call_ui( - live.update, + self._show_error, f"[red]{_escape(type(exc).__name__)}:[/] " f"{_escape(str(exc))}", ) diff --git a/tests/test_tui.py b/tests/test_tui.py index 4ca6f8a..4b80423 100644 --- a/tests/test_tui.py +++ b/tests/test_tui.py @@ -8,6 +8,7 @@ from rich.text import Text from nexusos_cli.tui_app import ( _compact_status, + _deny_tool_request, _escape, _status_line, format_assistant_line, @@ -15,6 +16,28 @@ from nexusos_cli.tui_app import ( ) +class _ApprovalResponse: + def raise_for_status(self): + return None + + +class _ApprovalClient: + calls = [] + + def __init__(self, **kwargs): + self.kwargs = kwargs + + def __enter__(self): + return self + + def __exit__(self, *args): + return None + + def post(self, path, *, json): + self.calls.append((path, json, self.kwargs)) + return _ApprovalResponse() + + def test_status_line_mentions_services(): snap = { "version": "1.0.0", @@ -81,5 +104,43 @@ def test_finish_stream_markup_does_not_wedge_busy(): asyncio.run(_run()) +def test_stream_error_remains_visible_after_finish(): + pytest.importorskip("textual") + from nexusos_cli.tui_app import NexusTUI + + app = NexusTUI.build_app(api_url="http://127.0.0.1:9") + + async def _run(): + async with app.run_test(): + app._busy = True + app._show_error("[red]Backend not reachable[/]") + app._finish_stream("") + log = app.query_one("#log") + assert any("Backend not reachable" in line.text for line in log.lines) + assert app._busy is False + + asyncio.run(_run()) + + +def test_tool_request_is_denied_with_stream_token(): + _ApprovalClient.calls.clear() + names = _deny_tool_request( + api_url="http://localhost:8000", + conversation_id="conversation-1", + payload='{"token":"secret","actions":[{"name":"run_snippet"}]}', + client_factory=_ApprovalClient, + ) + + assert names == ["run_snippet"] + path, body, client_kwargs = _ApprovalClient.calls[-1] + assert path == "/chat/approve" + assert body == { + "conversation_id": "conversation-1", + "token": "secret", + "decisions": {"run_snippet": False}, + } + assert client_kwargs["base_url"] == "http://localhost:8000" + + def test_escape_round_trip_helper(): assert "[" in _escape("x[y]") or "\\[" in _escape("x[y]") -- 2.39.5 From da3509eb04e1e0ddc512787c1d4dfc761a127420 Mon Sep 17 00:00:00 2001 From: Athena Kaminsky Date: Fri, 21 Aug 2026 16:19:59 -0500 Subject: [PATCH 3/5] fix(tui): retain stream conversation for tool denial Capture each stream's conversation ID before starting its worker so /new cannot redirect a later action denial. Add a headless regression that mutates the active conversation while a tool request is in flight. --- nexusos_cli/tui_app.py | 5 ++-- tests/test_tui.py | 63 ++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 66 insertions(+), 2 deletions(-) diff --git a/nexusos_cli/tui_app.py b/nexusos_cli/tui_app.py index 54275f7..1897bbe 100644 --- a/nexusos_cli/tui_app.py +++ b/nexusos_cli/tui_app.py @@ -354,9 +354,10 @@ class NexusTUI: self._stop_stream.clear() if not self.conversation_id: self.conversation_id = str(uuid.uuid4()) + conversation_id = self.conversation_id body: dict[str, Any] = { "message": message, - "conversation_id": self.conversation_id, + "conversation_id": conversation_id, "history": list(self.history), } if self._model: @@ -401,7 +402,7 @@ class NexusTUI: try: names = _deny_tool_request( api_url=self.api_url, - conversation_id=self.conversation_id, + conversation_id=conversation_id, payload=payload, ) shown = ", ".join(names) diff --git a/tests/test_tui.py b/tests/test_tui.py index 4b80423..7a0ac1e 100644 --- a/tests/test_tui.py +++ b/tests/test_tui.py @@ -2,6 +2,7 @@ from __future__ import annotations import asyncio +import threading import pytest from rich.text import Text @@ -142,5 +143,67 @@ def test_tool_request_is_denied_with_stream_token(): assert client_kwargs["base_url"] == "http://localhost:8000" +def test_inflight_tool_denial_uses_original_conversation_id(monkeypatch): + pytest.importorskip("textual") + import nexusos_cli.tui_app as tui_app + + stream_started = threading.Event() + release_stream = threading.Event() + denied_for = [] + + class _StreamResponse: + status_code = 200 + + def __enter__(self): + return self + + def __exit__(self, *args): + return None + + def iter_lines(self): + stream_started.set() + release_stream.wait(timeout=2) + yield "event: tool_request" + yield 'data: {"token":"secret","actions":[{"name":"run_snippet"}]}' + yield "" + yield "event: done" + yield "data: {}" + + class _StreamClient: + def __init__(self, **kwargs): + pass + + def stream(self, *args, **kwargs): + return _StreamResponse() + + def close(self): + return None + + def _capture_denial(*, conversation_id, **kwargs): + denied_for.append(conversation_id) + return ["run_snippet"] + + monkeypatch.setattr(tui_app.httpx, "Client", _StreamClient) + monkeypatch.setattr(tui_app, "_deny_tool_request", _capture_denial) + app = tui_app.NexusTUI.build_app(api_url="http://127.0.0.1:9") + + async def _run(): + async with app.run_test(): + app._start_chat("run it") + assert await asyncio.to_thread(stream_started.wait, 2) + original_id = app.conversation_id + app._handle_slash("/new") + assert app.conversation_id is None + release_stream.set() + for _ in range(200): + if not app._busy: + break + await asyncio.sleep(0.01) + assert app._busy is False + assert denied_for == [original_id] + + asyncio.run(_run()) + + def test_escape_round_trip_helper(): assert "[" in _escape("x[y]") or "\\[" in _escape("x[y]") -- 2.39.5 From 9ca37057eb66f490cc4a08d20a2f794bf2e4b25c Mon Sep 17 00:00:00 2001 From: Athena Kaminsky Date: Fri, 21 Aug 2026 17:07:29 -0500 Subject: [PATCH 4/5] fix(tui): cancel silent streams promptly Move chat streaming onto a cancellable async task so Ctrl+C interrupts a pending socket read on macOS instead of waiting for the 120-second read timeout. Add a headless silent-stream regression that verifies prompt recovery and a successful next message. --- nexusos_cli/tui_app.py | 194 +++++++++++++++++++++++------------------ tests/test_tui.py | 92 +++++++++++++++++-- 2 files changed, 193 insertions(+), 93 deletions(-) diff --git a/nexusos_cli/tui_app.py b/nexusos_cli/tui_app.py index 1897bbe..21945df 100644 --- a/nexusos_cli/tui_app.py +++ b/nexusos_cli/tui_app.py @@ -6,6 +6,7 @@ stdin/stdout are a TTY. Classic one-shots (``nexus chat send``, ``nexus monitor` """ from __future__ import annotations +import asyncio import json import threading import uuid @@ -15,7 +16,6 @@ import httpx from synapse.nexus_config import settings -from . import nexus_api from .monitor import collect_snapshot # Between SSE chunks a silent backend must not pin the UI forever. Connect stays @@ -199,7 +199,9 @@ class NexusTUI: self._model: str | None = None self._busy = False self._stop_stream = threading.Event() - self._http_client: httpx.Client | None = None + self._stream_cancel: ( + tuple[asyncio.AbstractEventLoop, asyncio.Task] | None + ) = None self._status_lock = threading.Lock() self._status_pending = False @@ -269,25 +271,26 @@ class NexusTUI: """Write a stream error to the persistent transcript.""" self.query_one("#log", RichLog).write(message) - def action_quit(self) -> None: + def _cancel_stream(self) -> None: + """Cancel the task that owns the socket read. + + Closing a synchronous httpx client from the UI thread does not + reliably unblock its worker-thread read on macOS. Async task + cancellation is delivered to the pending read itself. + """ self._stop_stream.set() - client = self._http_client - if client is not None: - try: - client.close() - except Exception: - pass + cancel = self._stream_cancel + if cancel is not None: + loop, task = cancel + loop.call_soon_threadsafe(task.cancel) + + def action_quit(self) -> None: + self._cancel_stream() self.exit() def action_interrupt(self) -> None: if self._busy: - self._stop_stream.set() - client = self._http_client - if client is not None: - try: - client.close() - except Exception: - pass + self._cancel_stream() self.query_one("#log", RichLog).write( "[yellow]▸ interrupt requested[/]" ) @@ -364,74 +367,96 @@ class NexusTUI: body["model"] = self._model self.history.append({"role": "user", "content": message}) - def worker(): + async def stream_worker(): reply_parts: list[str] = [] - client = httpx.Client( - base_url=self.api_url, timeout=_STREAM_TIMEOUT - ) - self._http_client = client + task = asyncio.current_task() + loop = asyncio.get_running_loop() + if task is None: # pragma: no cover - asyncio guarantees it + raise RuntimeError("stream worker has no task") + self._stream_cancel = (loop, task) try: - with client.stream( - "POST", "/chat/stream", json=body - ) as resp: - if resp.status_code >= 400: - detail = resp.read().decode( - "utf-8", errors="replace" - )[:300] - self._call_ui( - self._show_error, - f"[red]error HTTP {resp.status_code}[/] " - f"{_escape(detail)}", - ) - return - for kind, payload in nexus_api.iter_chunks( - resp.iter_lines() - ): - if self._stop_stream.is_set(): - break - if kind == "chunk": - reply_parts.append(payload) - preview = "".join(reply_parts) - if len(preview) > 4000: - preview = "…" + preview[-4000:] - self._call_ui( - live.update, - format_assistant_line(preview), - ) - elif kind == "tool_request": - try: - names = _deny_tool_request( - api_url=self.api_url, - conversation_id=conversation_id, - payload=payload, - ) - shown = ", ".join(names) - self._call_ui( - log.write, - "[yellow]▸ denied action tool " - f"{_escape(shown)} — interactive approval " - "is not yet available in the TUI[/]", - ) - except Exception as exc: - self._call_ui( - self._show_error, - "[red]tool denial failed:[/] " - f"{_escape(str(exc))}", - ) - return - elif kind == "error": - try: - detail = json.loads(payload).get( - "detail", payload - ) - except Exception: - detail = payload + if self._stop_stream.is_set(): + raise asyncio.CancelledError + async with httpx.AsyncClient( + base_url=self.api_url, timeout=_STREAM_TIMEOUT + ) as client: + async with client.stream( + "POST", "/chat/stream", json=body + ) as resp: + if resp.status_code >= 400: + detail = (await resp.aread()).decode( + "utf-8", errors="replace" + )[:300] self._call_ui( self._show_error, - f"[red]error:[/] {_escape(str(detail))}", + f"[red]error HTTP {resp.status_code}[/] " + f"{_escape(detail)}", ) - elif kind == "done": - break + return + event = "message" + async for line in resp.aiter_lines(): + if self._stop_stream.is_set(): + raise asyncio.CancelledError + if line == "": + event = "message" + continue + if line.startswith("event:"): + event = line[6:].strip() + continue + if not line.startswith("data:"): + continue + payload = line[5:].strip() + kind = event + if kind in ("message", ""): + kind = "chunk" + payload = json.loads(payload) + if kind == "chunk": + reply_parts.append(payload) + preview = "".join(reply_parts) + if len(preview) > 4000: + preview = "…" + preview[-4000:] + self._call_ui( + live.update, + format_assistant_line(preview), + ) + elif kind == "tool_request": + try: + names = _deny_tool_request( + api_url=self.api_url, + conversation_id=conversation_id, + payload=payload, + ) + shown = ", ".join(names) + self._call_ui( + log.write, + "[yellow]▸ denied action tool " + f"{_escape(shown)} — interactive " + "approval is not yet available in " + "the TUI[/]", + ) + except Exception as exc: + self._call_ui( + self._show_error, + "[red]tool denial failed:[/] " + f"{_escape(str(exc))}", + ) + return + elif kind == "error": + try: + detail = json.loads(payload).get( + "detail", payload + ) + except Exception: + detail = payload + self._call_ui( + self._show_error, + f"[red]error:[/] " + f"{_escape(str(detail))}", + ) + elif kind == "done": + break + except asyncio.CancelledError: + pass except httpx.ConnectError: self._call_ui( self._show_error, @@ -445,15 +470,14 @@ class NexusTUI: f"{_escape(str(exc))}", ) finally: - self._http_client = None - try: - client.close() - except Exception: - pass + if self._stream_cancel == (loop, task): + self._stream_cancel = None text = "".join(reply_parts).strip() self._call_ui(self._finish_stream, text) - threading.Thread(target=worker, daemon=True).start() + threading.Thread( + target=lambda: asyncio.run(stream_worker()), daemon=True + ).start() def _finish_stream(self, text: str) -> None: log = self.query_one("#log", RichLog) diff --git a/tests/test_tui.py b/tests/test_tui.py index 7a0ac1e..505c285 100644 --- a/tests/test_tui.py +++ b/tests/test_tui.py @@ -154,15 +154,15 @@ def test_inflight_tool_denial_uses_original_conversation_id(monkeypatch): class _StreamResponse: status_code = 200 - def __enter__(self): + async def __aenter__(self): return self - def __exit__(self, *args): + async def __aexit__(self, *args): return None - def iter_lines(self): + async def aiter_lines(self): stream_started.set() - release_stream.wait(timeout=2) + await asyncio.to_thread(release_stream.wait, 2) yield "event: tool_request" yield 'data: {"token":"secret","actions":[{"name":"run_snippet"}]}' yield "" @@ -173,17 +173,20 @@ def test_inflight_tool_denial_uses_original_conversation_id(monkeypatch): def __init__(self, **kwargs): pass + async def __aenter__(self): + return self + + async def __aexit__(self, *args): + return None + def stream(self, *args, **kwargs): return _StreamResponse() - def close(self): - return None - def _capture_denial(*, conversation_id, **kwargs): denied_for.append(conversation_id) return ["run_snippet"] - monkeypatch.setattr(tui_app.httpx, "Client", _StreamClient) + monkeypatch.setattr(tui_app.httpx, "AsyncClient", _StreamClient) monkeypatch.setattr(tui_app, "_deny_tool_request", _capture_denial) app = tui_app.NexusTUI.build_app(api_url="http://127.0.0.1:9") @@ -205,5 +208,78 @@ def test_inflight_tool_denial_uses_original_conversation_id(monkeypatch): asyncio.run(_run()) +def test_interrupt_cancels_silent_stream_and_accepts_next_message(monkeypatch): + pytest.importorskip("textual") + import nexusos_cli.tui_app as tui_app + + first_stream_started = threading.Event() + + class _StreamResponse: + status_code = 200 + + def __init__(self, call_number): + self.call_number = call_number + + async def __aenter__(self): + return self + + async def __aexit__(self, *args): + return None + + async def aiter_lines(self): + if self.call_number == 1: + first_stream_started.set() + await asyncio.Event().wait() + yield "data: \"READY\"" + yield "" + yield "event: done" + yield "data: {}" + + class _StreamClient: + calls = 0 + + def __init__(self, **kwargs): + pass + + async def __aenter__(self): + return self + + async def __aexit__(self, *args): + return None + + def stream(self, *args, **kwargs): + type(self).calls += 1 + return _StreamResponse(type(self).calls) + + monkeypatch.setattr(tui_app.httpx, "AsyncClient", _StreamClient) + app = tui_app.NexusTUI.build_app(api_url="http://127.0.0.1:9") + + async def _wait_until_idle(): + for _ in range(100): + if not app._busy: + return + await asyncio.sleep(0.01) + pytest.fail("stream did not become idle within one second") + + async def _run(): + async with app.run_test(): + app._start_chat("first") + assert await asyncio.to_thread(first_stream_started.wait, 2) + app.action_interrupt() + await _wait_until_idle() + + log = app.query_one("#log") + assert not any("ReadTimeout" in line.text for line in log.lines) + + app._start_chat("second") + await _wait_until_idle() + assert app.history[-1] == { + "role": "assistant", + "content": "READY", + } + + asyncio.run(_run()) + + def test_escape_round_trip_helper(): assert "[" in _escape("x[y]") or "\\[" in _escape("x[y]") -- 2.39.5 From 9ed29081708938e54ce0c5c07094216340d1f3ee Mon Sep 17 00:00:00 2001 From: Athena Kaminsky Date: Tue, 25 Aug 2026 15:30:03 -0500 Subject: [PATCH 5/5] fix(tui): prioritize interrupt and quit keys Declare Ctrl+C and Ctrl+D as priority Textual bindings so the focused prompt cannot consume them. Drive exit and silent-stream cancellation regressions through Pilot key events instead of calling action handlers directly. --- nexusos_cli/tui_app.py | 9 +++++---- tests/test_tui.py | 22 ++++++++++++++++++++-- 2 files changed, 25 insertions(+), 6 deletions(-) diff --git a/nexusos_cli/tui_app.py b/nexusos_cli/tui_app.py index 21945df..31c4135 100644 --- a/nexusos_cli/tui_app.py +++ b/nexusos_cli/tui_app.py @@ -27,13 +27,14 @@ _APPROVAL_TIMEOUT = httpx.Timeout(10.0, connect=5.0) def _require_textual(): try: from textual.app import App + from textual.binding import Binding from textual.widgets import Footer, Header, Input, RichLog, Static except ImportError as e: # pragma: no cover - optional extra raise ImportError( "The interactive TUI needs the 'tui' extra — " "pip install 'nexusos-ai[tui]' (or: pip install textual)." ) from e - return App, Footer, Header, Input, RichLog, Static + return App, Binding, Footer, Header, Input, RichLog, Static def _escape(text: str) -> str: @@ -154,7 +155,7 @@ class NexusTUI: @staticmethod def build_app(*, api_url: str | None = None): - App, Footer, Header, Input, RichLog, Static = _require_textual() + App, Binding, Footer, Header, Input, RichLog, Static = _require_textual() base = (api_url or settings.api_url).rstrip("/") class AppImpl(App): @@ -187,8 +188,8 @@ class NexusTUI: #prompt { dock: bottom; } """ BINDINGS = [ - ("ctrl+c", "interrupt", "Interrupt"), - ("ctrl+d", "quit", "Quit"), + Binding("ctrl+c", "interrupt", "Interrupt", priority=True), + Binding("ctrl+d", "quit", "Quit", priority=True), ] def __init__(self): diff --git a/tests/test_tui.py b/tests/test_tui.py index 505c285..9d11f32 100644 --- a/tests/test_tui.py +++ b/tests/test_tui.py @@ -123,6 +123,23 @@ def test_stream_error_remains_visible_after_finish(): asyncio.run(_run()) +@pytest.mark.parametrize("key", ["ctrl+c", "ctrl+d"]) +def test_priority_exit_bindings_reach_app_while_prompt_is_focused(key): + pytest.importorskip("textual") + from nexusos_cli.tui_app import NexusTUI + + app = NexusTUI.build_app(api_url="http://127.0.0.1:9") + + async def _run(): + async with app.run_test() as pilot: + assert app.is_running + await pilot.press(key) + await pilot.pause() + assert not app.is_running + + asyncio.run(_run()) + + def test_tool_request_is_denied_with_stream_token(): _ApprovalClient.calls.clear() names = _deny_tool_request( @@ -262,13 +279,14 @@ def test_interrupt_cancels_silent_stream_and_accepts_next_message(monkeypatch): pytest.fail("stream did not become idle within one second") async def _run(): - async with app.run_test(): + async with app.run_test() as pilot: app._start_chat("first") assert await asyncio.to_thread(first_stream_started.wait, 2) - app.action_interrupt() + await pilot.press("ctrl+c") await _wait_until_idle() log = app.query_one("#log") + assert any("interrupt requested" in line.text for line in log.lines) assert not any("ReadTimeout" in line.text for line in log.lines) app._start_chat("second") -- 2.39.5