"""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