"""Process supervisor: one `node bot.js` per bot, kept alive. Responsibilities: * spawn/stop/restart bot processes * capture stdout+stderr into a bounded ring buffer with monotonic sequence numbers * parse the structured `__BB__{json}` lines the generated runtime emits (block steps, errors) and surface them alongside plain console output * restart crashed bots with exponential backoff * relaunch every enabled bot on boot Security: bot subprocesses get a *scrubbed* environment. HF_TOKEN and APP_KEY are removed so a generated bot can never reach the Hugging Face account or impersonate the control API. """ from __future__ import annotations import collections import json import os import shutil import signal import subprocess import threading import time from dataclasses import dataclass, field from pathlib import Path from typing import Any DATA_DIR = Path(os.environ.get("BB_DATA_DIR", "/home/user/data")) BOTS_DIR = DATA_DIR / "bots" SHARED_MODULES = DATA_DIR / "node_modules" LOG_CAPACITY = 500 # lines kept per bot BACKOFF_BASE = 2.0 # seconds BACKOFF_MAX = 120.0 CRASH_RESET_AFTER = 60.0 # run this long and the backoff resets # Variables never exposed to bot processes. SECRET_ENV_KEYS = ( "HF_TOKEN", "APP_KEY", "HUGGING_FACE_HUB_TOKEN", "HF_API_TOKEN", "BB_DATASET", "BB_SELF_URL", ) @dataclass class LogLine: seq: int ts: float stream: str # "stdout" | "stderr" | "system" text: str event: dict | None = None # decoded __BB__ payload, if this was a structured line @dataclass class Bot: id: str name: str = "" enabled: bool = False proc: subprocess.Popen | None = None logs: collections.deque = field(default_factory=lambda: collections.deque(maxlen=LOG_CAPACITY)) seq: int = 0 status: str = "stopped" # stopped | starting | running | crashed last_error: str | None = None current_block: str | None = None started_at: float | None = None restarts: int = 0 _backoff: float = BACKOFF_BASE _stop_requested: bool = False _lock: threading.Lock = field(default_factory=threading.Lock) _subscribers: list = field(default_factory=list) @property def dir(self) -> Path: return BOTS_DIR / self.id def public(self) -> dict[str, Any]: return { "id": self.id, "name": self.name, "enabled": self.enabled, "status": self.status, "pid": self.proc.pid if self.proc and self.proc.poll() is None else None, "uptime": (time.time() - self.started_at) if self.started_at and self.status == "running" else 0, "restarts": self.restarts, "last_error": self.last_error, "current_block": self.current_block, } class Supervisor: def __init__(self, node_bin: str = "node", on_change=None): self.node_bin = node_bin self.bots: dict[str, Bot] = {} self.on_change = on_change # called after state changes worth persisting self._lock = threading.Lock() BOTS_DIR.mkdir(parents=True, exist_ok=True) # ---------------------------------------------------------------- registry def load_from_disk(self) -> None: """Rebuild the bot registry from DATA_DIR (called after the dataset sync).""" if not BOTS_DIR.exists(): return for d in sorted(BOTS_DIR.iterdir()): if not d.is_dir(): continue meta_path = d / "meta.json" meta = {} if meta_path.exists(): try: meta = json.loads(meta_path.read_text()) except json.JSONDecodeError: pass bot = self.bots.get(d.name) or Bot(id=d.name) bot.name = meta.get("name", d.name) bot.enabled = bool(meta.get("enabled", False)) self.bots[bot.id] = bot def get(self, bot_id: str) -> Bot | None: return self.bots.get(bot_id) def ensure(self, bot_id: str, name: str = "") -> Bot: bot = self.bots.get(bot_id) if bot is None: bot = Bot(id=bot_id, name=name or bot_id) self.bots[bot_id] = bot elif name: bot.name = name bot.dir.mkdir(parents=True, exist_ok=True) return bot def write_meta(self, bot: Bot) -> None: (bot.dir / "meta.json").write_text(json.dumps({ "id": bot.id, "name": bot.name, "enabled": bot.enabled, "updated_at": time.time(), }, indent=2)) def delete(self, bot_id: str) -> bool: bot = self.bots.get(bot_id) if bot is None: return False # Kill the process *before* dropping the entry. Popping first means stop() can no # longer find the bot, so the node process is orphaned: still running, still talking # to Nerimity, and no longer reachable by any endpoint. bot._stop_requested = True bot.enabled = False self._terminate(bot) bot.status = "stopped" self.bots.pop(bot_id, None) shutil.rmtree(bot.dir, ignore_errors=True) return True # ---------------------------------------------------------------- logging def log(self, bot: Bot, text: str, stream: str = "system", event: dict | None = None) -> None: with bot._lock: bot.seq += 1 line = LogLine(seq=bot.seq, ts=time.time(), stream=stream, text=text, event=event) bot.logs.append(line) subs = list(bot._subscribers) payload = { "seq": line.seq, "ts": line.ts, "stream": line.stream, "text": line.text, "event": line.event, } for q in subs: try: q.put_nowait(payload) except Exception: pass def logs_since(self, bot: Bot, since: int = 0) -> list[dict]: with bot._lock: return [ {"seq": l.seq, "ts": l.ts, "stream": l.stream, "text": l.text, "event": l.event} for l in bot.logs if l.seq > since ] def subscribe(self, bot: Bot, q) -> None: with bot._lock: bot._subscribers.append(q) def unsubscribe(self, bot: Bot, q) -> None: with bot._lock: if q in bot._subscribers: bot._subscribers.remove(q) # ---------------------------------------------------------------- lifecycle def _child_env(self, bot: Bot) -> dict[str, str]: env = {k: v for k, v in os.environ.items() if k not in SECRET_ENV_KEYS} secret_path = bot.dir / "secret.json" token = "" if secret_path.exists(): try: token = json.loads(secret_path.read_text()).get("token", "") except json.JSONDecodeError: pass env["NERIMITY_TOKEN"] = token env["NODE_PATH"] = str(SHARED_MODULES) env["BB_BOT_ID"] = bot.id env["NODE_OPTIONS"] = "--max-old-space-size=256" return env def start(self, bot_id: str) -> tuple[bool, str]: bot = self.bots.get(bot_id) if bot is None: return False, "no such bot" if bot.proc and bot.proc.poll() is None: return True, "already running" entry = bot.dir / "bot.js" if not entry.exists(): return False, "no bot.js deployed" if not (bot.dir / "secret.json").exists(): return False, "no bot token configured" bot._stop_requested = False bot.enabled = True bot.status = "starting" bot.last_error = None self.write_meta(bot) try: bot.proc = subprocess.Popen( [self.node_bin, "bot.js"], cwd=str(bot.dir), env=self._child_env(bot), stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, bufsize=1, start_new_session=True, ) except Exception as exc: # noqa: BLE001 bot.status = "crashed" bot.last_error = str(exc) self.log(bot, f"failed to spawn: {exc}", "system") return False, str(exc) bot.started_at = time.time() bot.status = "running" self.log(bot, f"started (pid {bot.proc.pid})", "system") threading.Thread(target=self._pump, args=(bot, bot.proc.stdout, "stdout"), daemon=True).start() threading.Thread(target=self._pump, args=(bot, bot.proc.stderr, "stderr"), daemon=True).start() threading.Thread(target=self._watch, args=(bot,), daemon=True).start() if self.on_change: self.on_change(bot) return True, "started" def stop(self, bot_id: str, disable: bool = True) -> tuple[bool, str]: bot = self.bots.get(bot_id) if bot is None: return False, "no such bot" bot._stop_requested = True if disable: bot.enabled = False self.write_meta(bot) self._terminate(bot) bot.status = "stopped" bot.started_at = None bot.current_block = None self.log(bot, "stopped", "system") if self.on_change: self.on_change(bot) return True, "stopped" def _terminate(self, bot: Bot) -> None: """Kill a bot's process group, SIGTERM then SIGKILL. Safe to call repeatedly.""" proc = bot.proc if proc is None or proc.poll() is not None: bot.proc = None return try: os.killpg(os.getpgid(proc.pid), signal.SIGTERM) try: proc.wait(timeout=8) except subprocess.TimeoutExpired: os.killpg(os.getpgid(proc.pid), signal.SIGKILL) try: proc.wait(timeout=3) except subprocess.TimeoutExpired: pass except (ProcessLookupError, PermissionError, OSError): pass bot.proc = None def reap_orphans(self) -> list[int]: """Kill `node bot.js` processes the registry no longer knows about. A safety net: any bug that drops a bot from the registry without killing it leaves a process still talking to Nerimity that nothing can reach. Matching on the working directory keeps this from touching anything but our own bots. """ live = {str((BOTS_DIR / b.id).resolve()) for b in self.bots.values()} killed: list[int] = [] proc_root = Path("/proc") if not proc_root.exists(): return killed for entry in proc_root.iterdir(): if not entry.name.isdigit(): continue pid = int(entry.name) try: cmdline = (entry / "cmdline").read_bytes().split(b"\0") if len(cmdline) < 2 or not cmdline[0].endswith(b"node"): continue if cmdline[1] != b"bot.js": continue cwd = str((entry / "cwd").resolve()) except (OSError, PermissionError): continue if not cwd.startswith(str(BOTS_DIR.resolve())): continue if cwd in live: continue # a bot we still manage try: os.killpg(os.getpgid(pid), signal.SIGKILL) killed.append(pid) except (ProcessLookupError, PermissionError, OSError): try: os.kill(pid, signal.SIGKILL) killed.append(pid) except OSError: pass return killed def restart(self, bot_id: str) -> tuple[bool, str]: self.stop(bot_id, disable=False) time.sleep(0.4) return self.start(bot_id) def start_enabled(self) -> None: """Relaunch everything that was running before the Space restarted.""" for bot in list(self.bots.values()): if bot.enabled: ok, msg = self.start(bot.id) if not ok: self.log(bot, f"boot restore failed: {msg}", "system") def shutdown(self) -> None: for bot_id in list(self.bots): self.stop(bot_id, disable=False) # ---------------------------------------------------------------- internals def _pump(self, bot: Bot, pipe, stream: str) -> None: """Read a child pipe line by line, decoding structured runtime events.""" try: for raw in iter(pipe.readline, ""): line = raw.rstrip("\n") if not line: continue if line.startswith("__BB__"): try: event = json.loads(line[6:]) except json.JSONDecodeError: self.log(bot, line, stream) continue kind = event.get("t") if kind == "step": bot.current_block = event.get("block") continue # too chatty to log every step if kind == "error": bot.last_error = event.get("message") bot.current_block = event.get("block") self.log(bot, event.get("message", "error"), "stderr", event) continue self.log(bot, event.get("message", line), stream, event) else: self.log(bot, line, stream) except Exception: pass finally: try: pipe.close() except Exception: pass def _watch(self, bot: Bot) -> None: """Wait for exit, then decide whether to restart.""" proc = bot.proc if proc is None: return code = proc.wait() ran_for = time.time() - (bot.started_at or time.time()) if bot._stop_requested: bot.status = "stopped" return bot.status = "crashed" bot.restarts += 1 self.log(bot, f"exited with code {code} after {ran_for:.0f}s", "system") if ran_for > CRASH_RESET_AFTER: bot._backoff = BACKOFF_BASE delay = min(bot._backoff, BACKOFF_MAX) bot._backoff = min(bot._backoff * 2, BACKOFF_MAX) if not bot.enabled: return self.log(bot, f"restarting in {delay:.0f}s", "system") time.sleep(delay) if bot.enabled and not bot._stop_requested: self.start(bot.id)