bbuilder-host / supervisor.py
joddabod's picture
deploy bbuilder host
a04249a verified
Raw History Blame Contribute Delete
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",
)
@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)