Spaces:
Running on Zero
Running on Zero
Download supervisor.py from joddabod/bbuilder-host: direct link, hf CLI and curl.
- Browser
- Download file 14.7 kB
-
https://huggingface.co/spaces/joddabod/bbuilder-host/resolve/main/supervisor.py
- Command line
-
hf download hf://spaces/joddabod/bbuilder-host/supervisor.py
-
curl -L -o supervisor.py https://huggingface.co/spaces/joddabod/bbuilder-host/resolve/main/supervisor.py
14.7 kB
| """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", | |
| ) | |
| 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 | |
| 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) | |
| 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) | |