import "server-only"; import { spawn as nodeSpawn, type ChildProcess } from "node:child_process"; /** * Runs MediaMTX as a child process and keeps it running (vrek dec-xkn4z0e): restarts it * after a crash with exponential backoff, and stops it with the server. The spawn and * timer functions are injectable so tests never start a real process (pri-e14bahk). */ export type BridgeState = "stopped" | "missing" | "starting" | "running" | "crashed"; export interface BridgeStatus { state: BridgeState; /** Short, fixed text for the UI and logs; never includes secrets. */ detail?: string; restarts: number; } export interface SupervisorOptions { binary: string; configPath: string; cwd: string; /** Resolves once MediaMTX answers (e.g. its API); rejects if it never does. */ waitReady: () => Promise; spawn?: typeof nodeSpawn; setTimer?: (fn: () => void, ms: number) => unknown; clearTimer?: (handle: unknown) => void; now?: () => number; log?: (line: string) => void; } const MIN_BACKOFF_MS = 1_000; const MAX_BACKOFF_MS = 30_000; /** A run this long counts as healthy, so the next crash starts backoff from the minimum. */ const STABLE_MS = 60_000; export class MediamtxSupervisor { private child: ChildProcess | null = null; private timer: unknown = null; private backoff = MIN_BACKOFF_MS; private startedAt = 0; private stopping = false; private current: BridgeStatus = { state: "stopped", restarts: 0 }; private readonly spawnFn: typeof nodeSpawn; private readonly setTimer: (fn: () => void, ms: number) => unknown; private readonly clearTimer: (handle: unknown) => void; private readonly now: () => number; private readonly log: (line: string) => void; constructor(private readonly opts: SupervisorOptions) { this.spawnFn = opts.spawn ?? nodeSpawn; this.setTimer = opts.setTimer ?? ((fn, ms) => setTimeout(fn, ms).unref()); this.clearTimer = opts.clearTimer ?? ((h) => clearTimeout(h as NodeJS.Timeout)); this.now = opts.now ?? Date.now; this.log = opts.log ?? ((line) => console.log(line)); } status(): BridgeStatus { return { ...this.current }; } start(): void { if (this.child) return; this.stopping = false; this.launch(); } stop(): void { this.stopping = true; if (this.timer !== null) this.clearTimer(this.timer); this.timer = null; this.child?.kill("SIGTERM"); this.child = null; this.current = { ...this.current, state: "stopped", detail: undefined }; } private set(state: BridgeState, detail?: string) { this.current = { ...this.current, state, detail }; } private launch() { this.set("starting"); this.startedAt = this.now(); const child = this.spawnFn(this.opts.binary, [this.opts.configPath], { cwd: this.opts.cwd, stdio: ["ignore", "pipe", "pipe"], }); this.child = child; const forward = (chunk: Buffer) => { for (const line of chunk.toString("utf8").split("\n")) if (line.trim()) this.log(`[mediamtx] ${line.trimEnd()}`); }; child.stdout?.on("data", forward); child.stderr?.on("data", forward); child.once("error", (err: NodeJS.ErrnoException) => { if (this.child !== child) return; this.child = null; if (err.code === "ENOENT") { // Not installed: retrying won't help until someone runs the installer. this.set("missing", "MediaMTX is not installed. Run npm run video:install."); this.log(`[video] ${this.current.detail}`); return; } this.crashed(`MediaMTX failed to start (${err.code ?? "error"})`); }); child.once("exit", (code, signal) => { if (this.child !== child) return; this.child = null; if (this.stopping) return; this.crashed(`MediaMTX exited (${signal ?? `code ${code}`})`); }); this.opts.waitReady().then( () => { if (this.child === child) this.set("running"); }, () => { if (this.child !== child) return; this.log("[video] MediaMTX did not become ready; restarting it"); child.kill("SIGTERM"); }, ); } private crashed(detail: string) { if (this.now() - this.startedAt >= STABLE_MS) this.backoff = MIN_BACKOFF_MS; const wait = this.backoff; this.backoff = Math.min(this.backoff * 2, MAX_BACKOFF_MS); this.current = { state: "crashed", detail, restarts: this.current.restarts + 1 }; this.log(`[video] ${detail}; restarting in ${wait / 1000} s`); this.timer = this.setTimer(() => { this.timer = null; if (!this.stopping) this.launch(); }, wait); } }