365 lines
13 KiB
Python
365 lines
13 KiB
Python
|
|
#!/usr/bin/env python3
|
||
|
|
"""
|
||
|
|
chromebox-gateway.py — HTTPS fallback for chromebox/DM control when SSH is down.
|
||
|
|
|
||
|
|
The chromeboxes (headless Chromium per NetVM node) are normally driven over
|
||
|
|
SSH via muse-chat-api.py / dm.py, which speak CDP through the per-node
|
||
|
|
netns relay. When SSH to bl breaks, this gateway provides a constrained
|
||
|
|
HTTPS path to the same high-level operations.
|
||
|
|
|
||
|
|
CRITICAL: this does NOT expose raw CDP. Raw CDP (Runtime.evaluate,
|
||
|
|
Page.navigate, etc.) is arbitrary code execution inside the browser with
|
||
|
|
the agent's live session. This gateway exposes ONLY the allowlisted
|
||
|
|
high-level operations below, each mapped to an existing audited script.
|
||
|
|
|
||
|
|
Usage:
|
||
|
|
python3 chromebox-gateway.py --port 8444
|
||
|
|
|
||
|
|
Auth: Bearer token, reusing the shared exec per-agent token files
|
||
|
|
(~/.exec-tokens/<agent>). The master token (~/.exec-server-token) also
|
||
|
|
works. Identity = token filename; tokens are never logged.
|
||
|
|
|
||
|
|
Endpoints:
|
||
|
|
GET /health {"status":"ok"} — no auth (load-balancer friendly)
|
||
|
|
POST /api/v1/op {"op": "<name>", "params": {...}} — bearer auth
|
||
|
|
|
||
|
|
Audit: every call appended to ~/.chromebox-gateway-audit.jsonl (0600).
|
||
|
|
"""
|
||
|
|
import argparse
|
||
|
|
import hmac
|
||
|
|
import json
|
||
|
|
import os
|
||
|
|
import re
|
||
|
|
import ssl
|
||
|
|
import subprocess
|
||
|
|
import sys
|
||
|
|
import time
|
||
|
|
from http.server import HTTPServer, BaseHTTPRequestHandler
|
||
|
|
from urllib.parse import urlparse
|
||
|
|
|
||
|
|
# ---------------------------------------------------------------- config
|
||
|
|
|
||
|
|
BIN_DIR = os.path.expanduser("~/Projects/NetVM/bin")
|
||
|
|
TOKEN_FILE = "/home/super/.exec-server-token" # master token (shared exec token file)
|
||
|
|
TOKEN_DIR = "/home/super/.exec-tokens" # per-agent tokens (shared exec token files)
|
||
|
|
AUDIT_FILE = "/home/super/.chromebox-gateway-audit.jsonl"
|
||
|
|
CERT_FILE = "/home/super/.chromebox-gateway-cert.pem"
|
||
|
|
KEY_FILE = "/home/super/.chromebox-gateway-key.pem"
|
||
|
|
|
||
|
|
NODES = ("muse", "pip", "646", "opm")
|
||
|
|
MAX_MSG = 2000 # gateway-level cap; downstream scripts enforce their own
|
||
|
|
BACKEND_TIMEOUT = 120 # seconds per backend call
|
||
|
|
|
||
|
|
# Rate limit: token bucket per identity — 20 req/min sustained, burst 5.
|
||
|
|
RATE_PER_SEC = 20.0 / 60.0
|
||
|
|
RATE_BURST = 5
|
||
|
|
|
||
|
|
# ---------------------------------------------------------------- ops
|
||
|
|
# Each op maps to an argv builder for an existing script. No shell=True,
|
||
|
|
# ever. Params are validated before building argv.
|
||
|
|
|
||
|
|
def _node(params):
|
||
|
|
node = params.get("account") or params.get("agent")
|
||
|
|
if node not in NODES:
|
||
|
|
raise ValueError("account/agent must be one of %s" % (",".join(NODES)))
|
||
|
|
return node
|
||
|
|
|
||
|
|
def _msg(params):
|
||
|
|
m = params.get("message", "")
|
||
|
|
if not isinstance(m, str) or not m.strip():
|
||
|
|
raise ValueError("message must be a non-empty string")
|
||
|
|
if len(m) > MAX_MSG:
|
||
|
|
raise ValueError("message exceeds %d chars" % MAX_MSG)
|
||
|
|
return m
|
||
|
|
|
||
|
|
def _target(params):
|
||
|
|
t = params.get("target", "main")
|
||
|
|
if not isinstance(t, str) or not t or len(t) > 128:
|
||
|
|
raise ValueError("target must be a short string")
|
||
|
|
if not re.fullmatch(r"[A-Za-z0-9_./:-]+", t):
|
||
|
|
raise ValueError("target has invalid characters")
|
||
|
|
if ".." in t:
|
||
|
|
raise ValueError("target must not contain '..'")
|
||
|
|
return t
|
||
|
|
|
||
|
|
def _n(params, default=5, cap=50):
|
||
|
|
n = params.get("n", default)
|
||
|
|
try:
|
||
|
|
n = int(n)
|
||
|
|
except (TypeError, ValueError):
|
||
|
|
raise ValueError("n must be an integer")
|
||
|
|
if not 1 <= n <= cap:
|
||
|
|
raise ValueError("n must be 1..%d" % cap)
|
||
|
|
return n
|
||
|
|
|
||
|
|
def _tags(params):
|
||
|
|
tags = params.get("tags", [])
|
||
|
|
if not isinstance(tags, list):
|
||
|
|
raise ValueError("tags must be a list")
|
||
|
|
out = []
|
||
|
|
for t in tags:
|
||
|
|
if not isinstance(t, str) or len(t) > 128:
|
||
|
|
raise ValueError("bad tag")
|
||
|
|
if not re.fullmatch(r"[A-Za-z0-9_:=\-./]+", t):
|
||
|
|
raise ValueError("tag has invalid characters: %r" % t[:40])
|
||
|
|
out.append(t)
|
||
|
|
if len(out) > 10:
|
||
|
|
raise ValueError("too many tags (max 10)")
|
||
|
|
return out
|
||
|
|
|
||
|
|
CHAT = os.path.join(BIN_DIR, "muse-chat-api.py")
|
||
|
|
DM = os.path.join(BIN_DIR, "dm.py")
|
||
|
|
|
||
|
|
def op_chat_send(p):
|
||
|
|
return [CHAT, "--account", _node(p), "send", _msg(p)]
|
||
|
|
|
||
|
|
def op_chat_messages(p):
|
||
|
|
return [CHAT, "--account", _node(p), "messages", "--n", str(_n(p))]
|
||
|
|
|
||
|
|
def op_chat_sidechats(p):
|
||
|
|
return [CHAT, "--account", _node(p), "sidechat", "list"]
|
||
|
|
|
||
|
|
def op_chat_sidechat_create(p):
|
||
|
|
argv = [CHAT, "--account", _node(p), "sidechat", "create"]
|
||
|
|
name = p.get("name")
|
||
|
|
if name:
|
||
|
|
if not isinstance(name, str) or len(name) > 80 or not re.fullmatch(r"[A-Za-z0-9 _-]+", name):
|
||
|
|
raise ValueError("bad sidechat name")
|
||
|
|
argv += ["--name", name]
|
||
|
|
return argv
|
||
|
|
|
||
|
|
def op_chat_approvals(p):
|
||
|
|
return [CHAT, "--account", _node(p), "approvals"]
|
||
|
|
|
||
|
|
def op_chat_url(p):
|
||
|
|
return [CHAT, "--account", _node(p), "url"]
|
||
|
|
|
||
|
|
def op_dm_send(p):
|
||
|
|
node = _node(p)
|
||
|
|
to = p.get("to", node)
|
||
|
|
if to not in NODES:
|
||
|
|
raise ValueError("to must be one of %s" % (",".join(NODES)))
|
||
|
|
argv = [DM, "send", "--agent", node, "--to", to,
|
||
|
|
"--target", _target(p), _msg(p)]
|
||
|
|
for t in _tags(p):
|
||
|
|
argv += ["--tag", t]
|
||
|
|
return argv
|
||
|
|
|
||
|
|
def op_dm_read(p):
|
||
|
|
return [DM, "read", "--agent", _node(p),
|
||
|
|
"--target", _target(p), "--n", str(_n(p, default=5, cap=20))]
|
||
|
|
|
||
|
|
# The allowlist. Adding an op here is a security decision — review accordingly.
|
||
|
|
# Deliberately absent: wait (long-poll), upload (file ingress),
|
||
|
|
# sidechat use (raw navigation), anything raw-CDP.
|
||
|
|
ALLOWLIST = {
|
||
|
|
"chat.send": ("write", op_chat_send),
|
||
|
|
"chat.messages": ("read", op_chat_messages),
|
||
|
|
"chat.sidechats": ("read", op_chat_sidechats),
|
||
|
|
"chat.sidechat_create": ("write", op_chat_sidechat_create),
|
||
|
|
"chat.approvals": ("read", op_chat_approvals),
|
||
|
|
"chat.url": ("read", op_chat_url),
|
||
|
|
"dm.send": ("write", op_dm_send),
|
||
|
|
"dm.read": ("read", op_dm_read),
|
||
|
|
}
|
||
|
|
|
||
|
|
# ---------------------------------------------------------------- auth
|
||
|
|
|
||
|
|
def _read_token_file(path):
|
||
|
|
try:
|
||
|
|
with open(path) as f:
|
||
|
|
return f.read().strip()
|
||
|
|
except OSError:
|
||
|
|
return ""
|
||
|
|
|
||
|
|
def check_token(token):
|
||
|
|
"""Return identity label or None. Tokens never leave this function."""
|
||
|
|
if not token:
|
||
|
|
return None
|
||
|
|
if hmac.compare_digest(token, _read_token_file(TOKEN_FILE)):
|
||
|
|
return "master"
|
||
|
|
try:
|
||
|
|
names = os.listdir(TOKEN_DIR)
|
||
|
|
except OSError:
|
||
|
|
return None
|
||
|
|
for name in names:
|
||
|
|
if not re.fullmatch(r"[A-Za-z0-9_-]+", name):
|
||
|
|
continue
|
||
|
|
t = _read_token_file(os.path.join(TOKEN_DIR, name))
|
||
|
|
if t and hmac.compare_digest(token, t):
|
||
|
|
return name
|
||
|
|
return None
|
||
|
|
|
||
|
|
# ---------------------------------------------------------------- rate limit
|
||
|
|
|
||
|
|
_buckets = {} # identity -> [tokens, last_ts]
|
||
|
|
|
||
|
|
def rate_ok(identity):
|
||
|
|
now = time.monotonic()
|
||
|
|
tokens, last = _buckets.get(identity, (RATE_BURST, now))
|
||
|
|
tokens = min(RATE_BURST, tokens + (now - last) * RATE_PER_SEC)
|
||
|
|
if tokens < 1.0:
|
||
|
|
_buckets[identity] = (tokens, now)
|
||
|
|
return False
|
||
|
|
_buckets[identity] = (tokens - 1.0, now)
|
||
|
|
return True
|
||
|
|
|
||
|
|
# ---------------------------------------------------------------- audit
|
||
|
|
|
||
|
|
def audit(entry):
|
||
|
|
entry = dict(entry)
|
||
|
|
entry["ts"] = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime())
|
||
|
|
# Defense in depth: never let a full message body reach the audit file,
|
||
|
|
# even if a caller forgets to summarize first.
|
||
|
|
params = entry.get("params")
|
||
|
|
if isinstance(params, dict) and "message" in params:
|
||
|
|
params = dict(params)
|
||
|
|
params["message_len"] = len(params.pop("message"))
|
||
|
|
entry["params"] = params
|
||
|
|
try:
|
||
|
|
fd = os.open(AUDIT_FILE, os.O_WRONLY | os.O_CREAT | os.O_APPEND, 0o600)
|
||
|
|
with os.fdopen(fd, "a") as f:
|
||
|
|
f.write(json.dumps(entry) + "\n")
|
||
|
|
except OSError as e:
|
||
|
|
print("audit write failed: %s" % e, file=sys.stderr)
|
||
|
|
|
||
|
|
# ---------------------------------------------------------------- handler
|
||
|
|
|
||
|
|
class Handler(BaseHTTPRequestHandler):
|
||
|
|
server_version = "chromebox-gateway/1.0"
|
||
|
|
|
||
|
|
def log_message(self, fmt, *args): # quiet; audit log is the record
|
||
|
|
pass
|
||
|
|
|
||
|
|
def _json(self, code, obj):
|
||
|
|
body = json.dumps(obj).encode()
|
||
|
|
self.send_response(code)
|
||
|
|
self.send_header("Content-Type", "application/json")
|
||
|
|
self.send_header("Content-Length", str(len(body)))
|
||
|
|
self.end_headers()
|
||
|
|
self.wfile.write(body)
|
||
|
|
|
||
|
|
def do_GET(self):
|
||
|
|
if urlparse(self.path).path == "/health":
|
||
|
|
self._json(200, {"status": "ok", "ops": sorted(ALLOWLIST)})
|
||
|
|
return
|
||
|
|
self._json(404, {"error": "not_found"})
|
||
|
|
|
||
|
|
def do_POST(self):
|
||
|
|
if urlparse(self.path).path != "/api/v1/op":
|
||
|
|
self._json(404, {"error": "not_found"})
|
||
|
|
return
|
||
|
|
|
||
|
|
auth = self.headers.get("Authorization", "")
|
||
|
|
token = auth[7:] if auth.startswith("Bearer ") else ""
|
||
|
|
identity = check_token(token)
|
||
|
|
if not identity:
|
||
|
|
self._json(401, {"error": "unauthorized"})
|
||
|
|
return
|
||
|
|
if not rate_ok(identity):
|
||
|
|
audit({"identity": identity, "op": None, "result": "rate_limited"})
|
||
|
|
self._json(429, {"error": "rate_limited"})
|
||
|
|
return
|
||
|
|
|
||
|
|
try:
|
||
|
|
length = int(self.headers.get("Content-Length", 0))
|
||
|
|
except ValueError:
|
||
|
|
length = 0
|
||
|
|
if length > 65536:
|
||
|
|
self._json(413, {"error": "body_too_large"})
|
||
|
|
return
|
||
|
|
try:
|
||
|
|
req = json.loads(self.rfile.read(length) or b"{}")
|
||
|
|
except (ValueError, OSError):
|
||
|
|
self._json(400, {"error": "bad_json"})
|
||
|
|
return
|
||
|
|
|
||
|
|
op = req.get("op")
|
||
|
|
params = req.get("params") or {}
|
||
|
|
if not isinstance(params, dict):
|
||
|
|
self._json(400, {"error": "params_must_be_object"})
|
||
|
|
return
|
||
|
|
entry = ALLOWLIST.get(op)
|
||
|
|
if not entry:
|
||
|
|
audit({"identity": identity, "op": op, "result": "unknown_op"})
|
||
|
|
self._json(400, {"error": "unknown_op", "allowed": sorted(ALLOWLIST)})
|
||
|
|
return
|
||
|
|
cls, builder = entry
|
||
|
|
|
||
|
|
t0 = time.monotonic()
|
||
|
|
try:
|
||
|
|
argv = builder(params)
|
||
|
|
except ValueError as e:
|
||
|
|
audit({"identity": identity, "op": op, "class": cls, "result": "bad_params",
|
||
|
|
"detail": str(e)[:120]})
|
||
|
|
self._json(400, {"error": "bad_params", "detail": str(e)})
|
||
|
|
return
|
||
|
|
|
||
|
|
# Redacted summary for the audit log — never the full message.
|
||
|
|
summary = {k: (v[:80] + "…" if isinstance(v, str) and len(v) > 80 else v)
|
||
|
|
for k, v in params.items() if k != "message"}
|
||
|
|
if "message" in params:
|
||
|
|
summary["message_len"] = len(params["message"])
|
||
|
|
|
||
|
|
try:
|
||
|
|
proc = subprocess.run(argv, capture_output=True, text=True,
|
||
|
|
timeout=BACKEND_TIMEOUT)
|
||
|
|
ok = proc.returncode == 0
|
||
|
|
result = "ok" if ok else "backend_error"
|
||
|
|
self._json(200 if ok else 502, {
|
||
|
|
"ok": ok,
|
||
|
|
"op": op,
|
||
|
|
"returncode": proc.returncode,
|
||
|
|
"stdout": proc.stdout[-8000:],
|
||
|
|
"stderr": proc.stderr[-2000:],
|
||
|
|
})
|
||
|
|
except subprocess.TimeoutExpired:
|
||
|
|
result = "timeout"
|
||
|
|
self._json(504, {"ok": False, "op": op, "error": "backend_timeout"})
|
||
|
|
except OSError as e:
|
||
|
|
result = "exec_failed"
|
||
|
|
self._json(500, {"ok": False, "op": op, "error": "exec_failed"})
|
||
|
|
finally:
|
||
|
|
audit({"identity": identity, "op": op, "class": cls,
|
||
|
|
"node": params.get("account") or params.get("agent"),
|
||
|
|
"params": summary, "result": result,
|
||
|
|
"latency_ms": int((time.monotonic() - t0) * 1000)})
|
||
|
|
|
||
|
|
# ---------------------------------------------------------------- main
|
||
|
|
|
||
|
|
def ensure_cert():
|
||
|
|
if os.path.exists(CERT_FILE) and os.path.exists(KEY_FILE):
|
||
|
|
return
|
||
|
|
print("generating self-signed cert...", file=sys.stderr)
|
||
|
|
subprocess.run([
|
||
|
|
"openssl", "req", "-x509", "-newkey", "rsa:2048",
|
||
|
|
"-keyout", KEY_FILE, "-out", CERT_FILE,
|
||
|
|
"-days", "825", "-nodes", "-subj", "/CN=chromebox-gateway",
|
||
|
|
], check=True)
|
||
|
|
os.chmod(KEY_FILE, 0o600)
|
||
|
|
os.chmod(CERT_FILE, 0o600)
|
||
|
|
|
||
|
|
def main():
|
||
|
|
ap = argparse.ArgumentParser()
|
||
|
|
ap.add_argument("--port", type=int, default=8444)
|
||
|
|
ap.add_argument("--bind", default="100.123.153.75",
|
||
|
|
help="tailnet IP; use 127.0.0.1 for local-only")
|
||
|
|
args = ap.parse_args()
|
||
|
|
|
||
|
|
ensure_cert()
|
||
|
|
context = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER)
|
||
|
|
context.load_cert_chain(CERT_FILE, KEY_FILE)
|
||
|
|
|
||
|
|
srv = HTTPServer((args.bind, args.port), Handler)
|
||
|
|
srv.socket = context.wrap_socket(srv.socket, server_side=True)
|
||
|
|
print("chromebox-gateway listening on https://%s:%d" % (args.bind, args.port),
|
||
|
|
file=sys.stderr)
|
||
|
|
print("ops: %s" % ", ".join(sorted(ALLOWLIST)), file=sys.stderr)
|
||
|
|
try:
|
||
|
|
srv.serve_forever()
|
||
|
|
except KeyboardInterrupt:
|
||
|
|
pass
|
||
|
|
|
||
|
|
if __name__ == "__main__":
|
||
|
|
main()
|