| 1 | """darwin-term - standalone WS→PTY→SSH terminal bridge for Darwin agent panes. |
| 2 | |
| 3 | Wire protocol (compatible with standard xterm.js client): |
| 4 | WS /ws?target=&agent=&mode=&session=&cwd=[&token=] |
| 5 | client→server: |
| 6 | binary frames - raw keystrokes/input bytes sent to PTY |
| 7 | text JSON - {"t":"resize","cols":N,"rows":N} → TIOCSWINSZ |
| 8 | server→client: |
| 9 | binary frames - raw PTY output |
| 10 | text (string) - banner lines + JSON control: |
| 11 | {"_cockpit":"session","name":<tmux-name>} (persistent mode only) |
| 12 | |
| 13 | Auth: DARWIN_TERM_TOKEN env var. Pass as ?token= query param or |
| 14 | Authorization: Bearer <token> header. Missing/wrong → close 4403. |
| 15 | (Future: replace with Authentik SSO.) |
| 16 | |
| 17 | Targets: ct215 (REDACTED-IP), ct241 (REDACTED-IP), ct247 (REDACTED-IP) |
| 18 | SSH key: /root/.ssh/termius_homelab (root auth, key-based, no password) |
| 19 | Port: 9894 |
| 20 | """ |
| 21 | |
| 22 | import asyncio |
| 23 | import fcntl |
| 24 | import json |
| 25 | import logging |
| 26 | import os |
| 27 | import pty |
| 28 | import re |
| 29 | import select |
| 30 | import shlex |
| 31 | import struct |
| 32 | import termios |
| 33 | import uuid |
| 34 | from contextlib import asynccontextmanager |
| 35 | |
| 36 | import uvicorn |
| 37 | from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Query |
| 38 | from fastapi.responses import JSONResponse |
| 39 | |
| 40 | logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s: %(message)s") |
| 41 | logger = logging.getLogger("darwin-term") |
| 42 | |
| 43 | SSH_KEY = os.environ.get("DARWIN_TERM_SSH_KEY", "/root/.ssh/termius_homelab") |
| 44 | KNOWN_HOSTS = "/root/.ssh/known_hosts" |
| 45 | DARWIN_TERM_TOKEN = os.environ.get("DARWIN_TERM_TOKEN", "") |
| 46 | |
| 47 | COMMON_DIRS = [ |
| 48 | {"label": "Default (home ~)", "path": ""}, |
| 49 | {"label": "Projects - /shared/projects", "path": "/shared/projects"}, |
| 50 | {"label": "Oversight - /shared/projects/Oversight", "path": "/shared/projects/Oversight"}, |
| 51 | ] |
| 52 | |
| 53 | TARGETS: dict[str, dict] = { |
| 54 | "ct215": { |
| 55 | "ip": "REDACTED-IP", |
| 56 | "label": "Claude Code - CT215", |
| 57 | "agents": ["claude", "shell"], |
| 58 | "default_agent": "claude", |
| 59 | "default_cwd": "/shared/projects", |
| 60 | }, |
| 61 | "ct241": { |
| 62 | "ip": "REDACTED-IP", |
| 63 | "label": "Codex - CT241", |
| 64 | "agents": ["codex", "shell"], |
| 65 | "default_agent": "codex", |
| 66 | "default_cwd": "/shared/projects", |
| 67 | }, |
| 68 | "ct247": { |
| 69 | "ip": "REDACTED-IP", |
| 70 | "label": "OpenCode / GLM - CT247", |
| 71 | "agents": ["opencode", "shell"], |
| 72 | "default_agent": "opencode", |
| 73 | "default_cwd": "", |
| 74 | }, |
| 75 | } |
| 76 | |
| 77 | AGENT_CMD: dict[str, str] = { |
| 78 | "claude": "CLAUDE_CODE_EXPERIMENTAL_AGENT_TEAMS=0 IS_SANDBOX=1 claude --dangerously-skip-permissions", |
| 79 | "codex": "codex", |
| 80 | "opencode": "opencode", |
| 81 | } |
| 82 | |
| 83 | SESSION_PREFIXES = tuple(f"{a}-" for a in list(AGENT_CMD) + ["shell"]) |
| 84 | |
| 85 | _NAME_RE = re.compile(r"^[A-Za-z0-9_.-]{1,80}$") |
| 86 | _CWD_RE = re.compile(r"^[A-Za-z0-9_./~ -]{1,256}$") |
| 87 | |
| 88 | |
| 89 | def _check_token(token_param: str | None, auth_header: str | None) -> bool: |
| 90 | """Return True if the request carries the correct bearer token.""" |
| 91 | if not DARWIN_TERM_TOKEN: |
| 92 | return True |
| 93 | candidate = token_param or "" |
| 94 | if not candidate and auth_header and auth_header.lower().startswith("bearer "): |
| 95 | candidate = auth_header[7:] |
| 96 | return candidate == DARWIN_TERM_TOKEN |
| 97 | |
| 98 | |
| 99 | def _ssh_argv(ip: str) -> list[str]: |
| 100 | return [ |
| 101 | "ssh", "-tt", |
| 102 | "-i", SSH_KEY, |
| 103 | "-o", "BatchMode=yes", |
| 104 | "-o", "ConnectTimeout=10", |
| 105 | "-o", "StrictHostKeyChecking=accept-new", |
| 106 | "-o", f"UserKnownHostsFile={KNOWN_HOSTS}", |
| 107 | "-o", "ServerAliveInterval=20", |
| 108 | f"root@{ip}", |
| 109 | ] |
| 110 | |
| 111 | |
| 112 | def _tmux_attach(session: str, payload: str) -> str: |
| 113 | """Create-or-attach a tmux session with deep history and mouse OFF. |
| 114 | `mouse off` keeps terminal output in xterm.js's own scrollback (the browser |
| 115 | viewport) instead of being captured into tmux copy-mode - which touch can't |
| 116 | drive - so the phone can finger-scroll terminal history. A large history-limit |
| 117 | still keeps the underlying tmux buffer deep across reattach.""" |
| 118 | q = shlex.quote(session) |
| 119 | create = f"tmux new-session -d -A -s {q} {shlex.quote(payload)}" |
| 120 | core = ( |
| 121 | "tmux set-option -g mouse off >/dev/null 2>&1; " |
| 122 | "tmux set-option -g history-limit 50000 >/dev/null 2>&1; " |
| 123 | "tmux set-option -g window-size smallest >/dev/null 2>&1; " |
| 124 | "tmux set-option -g aggressive-resize on >/dev/null 2>&1" |
| 125 | ) |
| 126 | skin = "tmux set-option -g status-style 'bg=colour234,fg=colour37' >/dev/null 2>&1 || true" |
| 127 | attach = f"tmux attach-session -t {q}" |
| 128 | return f"{create}; {core}; {skin}; {attach}" |
| 129 | |
| 130 | |
| 131 | def _remote_command(agent: str, mode: str, session: str, cwd: str | None) -> str: |
| 132 | if agent == "shell": |
| 133 | inner = "bash -l" |
| 134 | if cwd: |
| 135 | inner = f"bash -lc {shlex.quote(f'cd {shlex.quote(cwd)} && exec bash -l')}" |
| 136 | if mode == "persistent": |
| 137 | payload = inner if cwd else "bash -l" |
| 138 | return _tmux_attach(session, payload) |
| 139 | return inner |
| 140 | |
| 141 | base = AGENT_CMD[agent] |
| 142 | if cwd: |
| 143 | base = f"cd {shlex.quote(cwd)} && {base}" |
| 144 | run = f"bash -lc {shlex.quote(base)}" |
| 145 | if mode == "persistent": |
| 146 | return _tmux_attach(session, run) |
| 147 | return run |
| 148 | |
| 149 | |
| 150 | def _valid_session(name: str) -> bool: |
| 151 | return bool(name) and bool(_NAME_RE.match(name)) and name.startswith(SESSION_PREFIXES) |
| 152 | |
| 153 | |
| 154 | def _read_fd(fd: int) -> bytes | None: |
| 155 | """Blocking read via select; None = no data yet, b'' = EOF.""" |
| 156 | r, _, _ = select.select([fd], [], [], 1.0) |
| 157 | if r: |
| 158 | try: |
| 159 | data = os.read(fd, 65536) |
| 160 | return data if data else b"" |
| 161 | except OSError: |
| 162 | return b"" |
| 163 | return None |
| 164 | |
| 165 | |
| 166 | async def _run_ssh(ip: str, remote_cmd: str, timeout: float = 12.0) -> tuple[int, str, str]: |
| 167 | argv = [ |
| 168 | "ssh", "-i", SSH_KEY, |
| 169 | "-o", "BatchMode=yes", "-o", "ConnectTimeout=10", |
| 170 | "-o", "StrictHostKeyChecking=accept-new", |
| 171 | "-o", f"UserKnownHostsFile={KNOWN_HOSTS}", |
| 172 | f"root@{ip}", remote_cmd, |
| 173 | ] |
| 174 | proc = await asyncio.create_subprocess_exec( |
| 175 | *argv, |
| 176 | stdout=asyncio.subprocess.PIPE, |
| 177 | stderr=asyncio.subprocess.PIPE, |
| 178 | ) |
| 179 | try: |
| 180 | out, err = await asyncio.wait_for(proc.communicate(), timeout=timeout) |
| 181 | except asyncio.TimeoutError: |
| 182 | try: |
| 183 | proc.kill() |
| 184 | except ProcessLookupError: |
| 185 | pass |
| 186 | return -1, "", "timeout" |
| 187 | return proc.returncode, out.decode(errors="replace"), err.decode(errors="replace") |
| 188 | |
| 189 | |
| 190 | async def _list_sessions() -> list[dict]: |
| 191 | results: list[dict] = [] |
| 192 | |
| 193 | async def one(key: str, meta: dict): |
| 194 | _rc, out, _err = await _run_ssh( |
| 195 | meta["ip"], |
| 196 | "tmux list-sessions -F '#{session_name}|#{session_created}|#{session_attached}' 2>/dev/null || true", |
| 197 | ) |
| 198 | for line in out.splitlines(): |
| 199 | parts = line.strip().split("|") |
| 200 | if not parts or not parts[0]: |
| 201 | continue |
| 202 | name = parts[0] |
| 203 | if not name.startswith(SESSION_PREFIXES): |
| 204 | continue |
| 205 | results.append({ |
| 206 | "target": key, |
| 207 | "label": meta["label"], |
| 208 | "session": name, |
| 209 | "agent": name.split("-", 1)[0], |
| 210 | "created": parts[1] if len(parts) > 1 else "", |
| 211 | "attached": (parts[2] == "1") if len(parts) > 2 else False, |
| 212 | }) |
| 213 | |
| 214 | await asyncio.gather(*(one(k, m) for k, m in TARGETS.items())) |
| 215 | results.sort(key=lambda r: (r["target"], r["session"])) |
| 216 | return results |
| 217 | |
| 218 | |
| 219 | @asynccontextmanager |
| 220 | async def lifespan(_app: FastAPI): |
| 221 | if not DARWIN_TERM_TOKEN: |
| 222 | logger.warning("DARWIN_TERM_TOKEN is not set - service is OPEN, set it in /opt/darwin/.env") |
| 223 | else: |
| 224 | logger.info("darwin-term auth: token configured (length %d)", len(DARWIN_TERM_TOKEN)) |
| 225 | logger.info("SSH key: %s", SSH_KEY) |
| 226 | yield |
| 227 | |
| 228 | |
| 229 | app = FastAPI(title="darwin-term", version="1.0.0", lifespan=lifespan) |
| 230 | |
| 231 | |
| 232 | @app.get("/health") |
| 233 | async def health(): |
| 234 | return {"status": "ok", "service": "darwin-term"} |
| 235 | |
| 236 | |
| 237 | @app.get("/targets") |
| 238 | async def targets( |
| 239 | token: str | None = Query(default=None), |
| 240 | authorization: str | None = None, |
| 241 | ): |
| 242 | if not _check_token(token, authorization): |
| 243 | return JSONResponse({"error": "unauthorized"}, status_code=403) |
| 244 | return { |
| 245 | "targets": [ |
| 246 | { |
| 247 | "key": k, |
| 248 | "ip": m["ip"], |
| 249 | "label": m["label"], |
| 250 | "agents": m["agents"], |
| 251 | "default_agent": m["default_agent"], |
| 252 | "default_cwd": m.get("default_cwd", ""), |
| 253 | "dirs": COMMON_DIRS, |
| 254 | } |
| 255 | for k, m in TARGETS.items() |
| 256 | ], |
| 257 | "agent_cmds": AGENT_CMD, |
| 258 | "common_dirs": COMMON_DIRS, |
| 259 | } |
| 260 | |
| 261 | |
| 262 | @app.get("/sessions") |
| 263 | async def sessions( |
| 264 | token: str | None = Query(default=None), |
| 265 | authorization: str | None = None, |
| 266 | ): |
| 267 | if not _check_token(token, authorization): |
| 268 | return JSONResponse({"error": "unauthorized"}, status_code=403) |
| 269 | return {"sessions": await _list_sessions()} |
| 270 | |
| 271 | |
| 272 | @app.post("/kill") |
| 273 | async def kill( |
| 274 | target: str | None = Query(default=None), |
| 275 | session: str | None = Query(default=None), |
| 276 | token: str | None = Query(default=None), |
| 277 | authorization: str | None = None, |
| 278 | ): |
| 279 | if not _check_token(token, authorization): |
| 280 | return JSONResponse({"error": "unauthorized"}, status_code=403) |
| 281 | if target not in TARGETS: |
| 282 | return JSONResponse({"error": "unknown target"}, status_code=400) |
| 283 | if not _valid_session(session or ""): |
| 284 | return JSONResponse({"error": "invalid session"}, status_code=400) |
| 285 | meta = TARGETS[target] |
| 286 | rc, _out, err = await _run_ssh( |
| 287 | meta["ip"], |
| 288 | f"tmux kill-session -t {session} 2>&1 || true", |
| 289 | ) |
| 290 | return {"ok": rc == 0, "session": session, "detail": err.strip()} |
| 291 | |
| 292 | |
| 293 | @app.websocket("/ws") |
| 294 | async def terminal_ws(websocket: WebSocket): |
| 295 | q = websocket.query_params |
| 296 | token_param = q.get("token") |
| 297 | auth_header = websocket.headers.get("authorization") |
| 298 | |
| 299 | if not _check_token(token_param, auth_header): |
| 300 | await websocket.close(code=4403) |
| 301 | return |
| 302 | |
| 303 | target = q.get("target", "") |
| 304 | agent = q.get("agent", "shell") |
| 305 | mode = "persistent" if q.get("mode") == "persistent" else "ephemeral" |
| 306 | session = q.get("session", "") |
| 307 | cwd = q.get("cwd") or None |
| 308 | |
| 309 | if target not in TARGETS: |
| 310 | await websocket.accept() |
| 311 | await websocket.send_text("\r\n[darwin-term] unknown target\r\n") |
| 312 | await websocket.close() |
| 313 | return |
| 314 | |
| 315 | meta = TARGETS[target] |
| 316 | |
| 317 | if agent not in ("shell",) and agent not in meta["agents"]: |
| 318 | agent = meta["default_agent"] |
| 319 | |
| 320 | if cwd and not _CWD_RE.match(cwd): |
| 321 | await websocket.accept() |
| 322 | await websocket.send_text("\r\n[darwin-term] invalid cwd\r\n") |
| 323 | await websocket.close() |
| 324 | return |
| 325 | |
| 326 | if mode == "persistent": |
| 327 | if _valid_session(session): |
| 328 | pass |
| 329 | else: |
| 330 | session = f"{agent}-{uuid.uuid4().hex[:8]}" |
| 331 | else: |
| 332 | session = "" |
| 333 | |
| 334 | remote_cmd = _remote_command(agent, mode, session, cwd) |
| 335 | argv = _ssh_argv(meta["ip"]) + [remote_cmd] |
| 336 | |
| 337 | await websocket.accept() |
| 338 | banner = ( |
| 339 | f"[darwin-term] {meta['label']} · {agent} · {mode}" |
| 340 | + (f" · {session}" if session else "") |
| 341 | + "\r\n" |
| 342 | ) |
| 343 | await websocket.send_text(banner) |
| 344 | if mode == "persistent": |
| 345 | await websocket.send_text(json.dumps({"_cockpit": "session", "name": session}) + "\r\n") |
| 346 | |
| 347 | loop = asyncio.get_running_loop() |
| 348 | master_fd, slave_fd = pty.openpty() |
| 349 | |
| 350 | flags = fcntl.fcntl(master_fd, fcntl.F_GETFL) |
| 351 | fcntl.fcntl(master_fd, fcntl.F_SETFL, flags | os.O_NONBLOCK) |
| 352 | |
| 353 | try: |
| 354 | fcntl.ioctl(master_fd, termios.TIOCSWINSZ, struct.pack("HHHH", 24, 80, 0, 0)) |
| 355 | except OSError: |
| 356 | pass |
| 357 | |
| 358 | proc = await asyncio.create_subprocess_exec( |
| 359 | *argv, |
| 360 | stdin=slave_fd, |
| 361 | stdout=slave_fd, |
| 362 | stderr=slave_fd, |
| 363 | preexec_fn=os.setsid, |
| 364 | env={**os.environ, "TERM": "xterm-256color"}, |
| 365 | ) |
| 366 | os.close(slave_fd) |
| 367 | |
| 368 | def _set_winsize(rows: int, cols: int): |
| 369 | try: |
| 370 | fcntl.ioctl(master_fd, termios.TIOCSWINSZ, struct.pack("HHHH", rows, cols, 0, 0)) |
| 371 | except OSError: |
| 372 | pass |
| 373 | |
| 374 | stop = asyncio.Event() |
| 375 | |
| 376 | async def pty_to_ws(): |
| 377 | try: |
| 378 | while not stop.is_set(): |
| 379 | try: |
| 380 | data = await loop.run_in_executor(None, _read_fd, master_fd) |
| 381 | except OSError: |
| 382 | break |
| 383 | if data is None: |
| 384 | continue |
| 385 | if data == b"": |
| 386 | break |
| 387 | try: |
| 388 | await websocket.send_bytes(data) |
| 389 | except Exception: |
| 390 | break |
| 391 | finally: |
| 392 | stop.set() |
| 393 | |
| 394 | async def ws_to_pty(): |
| 395 | try: |
| 396 | while not stop.is_set(): |
| 397 | msg = await websocket.receive() |
| 398 | if msg.get("type") == "websocket.disconnect": |
| 399 | break |
| 400 | if msg.get("bytes") is not None: |
| 401 | os.write(master_fd, msg["bytes"]) |
| 402 | elif msg.get("text") is not None: |
| 403 | text = msg["text"] |
| 404 | if text.startswith('{"t":"resize"'): |
| 405 | try: |
| 406 | c = json.loads(text) |
| 407 | _set_winsize(int(c.get("rows", 24)), int(c.get("cols", 80))) |
| 408 | continue |
| 409 | except Exception: |
| 410 | pass |
| 411 | os.write(master_fd, text.encode()) |
| 412 | except WebSocketDisconnect: |
| 413 | pass |
| 414 | except Exception: |
| 415 | pass |
| 416 | finally: |
| 417 | stop.set() |
| 418 | |
| 419 | reader = asyncio.create_task(pty_to_ws()) |
| 420 | writer = asyncio.create_task(ws_to_pty()) |
| 421 | try: |
| 422 | await stop.wait() |
| 423 | finally: |
| 424 | reader.cancel() |
| 425 | writer.cancel() |
| 426 | try: |
| 427 | proc.terminate() |
| 428 | except ProcessLookupError: |
| 429 | pass |
| 430 | try: |
| 431 | os.close(master_fd) |
| 432 | except OSError: |
| 433 | pass |
| 434 | try: |
| 435 | await websocket.close() |
| 436 | except Exception: |
| 437 | pass |
| 438 | logger.info("session closed: target=%s agent=%s mode=%s session=%s", target, agent, mode, session or "-") |
| 439 | |
| 440 | |
| 441 | if __name__ == "__main__": |
| 442 | uvicorn.run( |
| 443 | "main:app", |
| 444 | host="0.0.0.0", |
| 445 | port=9894, |
| 446 | log_level="info", |
| 447 | ) |