Zion Boggan
repos/Darwin/term/main.py
zionboggan.com ↗
447 lines · python
History for this file →
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
    )